diff --git a/Case_02_Comparison_Research/YAML_Prompts/1. Stage_1/v.7/Part_3_to_Part_4_modification_summary.md b/Case_02_Comparison_Research/YAML_Prompts/1. Stage_1/v.7/Part_3_to_Part_4_modification_summary.md new file mode 100644 index 00000000..4f21289f --- /dev/null +++ b/Case_02_Comparison_Research/YAML_Prompts/1. Stage_1/v.7/Part_3_to_Part_4_modification_summary.md @@ -0,0 +1,35 @@ +# Stage 1 Part 4 v3에서 v4로의 개정 작업 요약 + +## 1. 개정 목적 + +`Stage_1_Part_4_Codex_v3.yml`은 `Fact_Ledger_base.json`의 27개 key 존재 여부만 검사하고 field type과 nested structure는 강제하지 않았다. 이에 따라 FL1이 upstream `amount` Object를 JSON String으로 바꾸고, `object_spec`도 자산 ID Array를 문자열화할 수 있었으며 FL3는 이를 차단하지 못했다. v4는 실행 코드의 의미 구조를 유지하면서 machine-readable JSON Schema를 단일 기준으로 도입하여 생성, 보정, 최종 저장 전 과정에 같은 자료형 계약을 적용한다. + +## 2. 핵심 개정 + +| 항목 | v3 | v4 | +|---|---|---| +| Schema | `FINAL_FIELDS` key 목록 | Draft 2020-12 `FACT_LEDGER_ROW_SCHEMA` 및 `fact_ledger_row_schema.json` | +| `amount` | Object를 String으로 변환 | upstream Object를 `deepcopy`, Object 또는 null 유지 | +| `object_spec` | JSON String 가능 | `source_object`, `asset_cluster_ids`를 가진 Object 또는 null | +| 파생 정보 | `derived_fact` String | `derived_fact` String/null과 `is_derived` Boolean 분리 | +| 도출 근거 | 세미콜론 String | BO, evidence, structure 참조를 담은 Array of Object | +| FL3 gate | key 집합만 확인 | type, enum, nested required key, ID, reference universe 검증 | + +FL1에는 `_structured_amount()`, `_object_spec()`, `_derivation_basis()`, `_validate_row_schema()`를 추가하였다. 최종 row는 `is_derived`가 추가된 28개 필드이다. `fact_id`는 `^F-[0-9]{3}$`, `credibility`는 `high|medium|low`, `must_consider`와 `is_derived`는 Boolean으로 고정하였다. `claim_chain_ref`, `linked_structures`, `linked_actio_structures`는 named Object로 유지하고 `legal_calculation_object`의 nested required key도 명시하였다. + +FL1은 schema의 `$id`와 SHA-256 digest를 candidate bundle의 `row_schema_ref` 및 `compile_gate`에 연결한다. FL3는 같은 artifact의 경로, `$id`, digest와 closed-schema 조건을 확인한 후 다음 네 시점에 row를 재검증한다. + +1. FL1 candidate 수신 시 +2. 각 FL2 `PATCH` 또는 `BLOCK_REVIEW` 적용 직후 +3. 최종 정렬 및 Fact ID 재부여 후 쓰기 직전 +4. `Fact_Ledger_base.json` 쓰기 후 재읽기 시 + +`amount`나 `object_spec`이 JSON처럼 보이는 String이면 hard failure 처리한다. `BLOCK_REVIEW`의 근거도 String 연결 대신 `derivation_basis`의 review provenance Object로 기록한다. candidate bundle, FL1 manifest, writer report와 final writer의 관련 계약은 v3으로 승격하였다. + +## 3. 보존 범위 및 검증 + +`FL0 -> FL1 -> 조건부 FL2 -> FL3 -> OUT` DAG, task 이름, Agent/Stage 메타데이터, FL0와 FL2는 변경하지 않았다. FL2 common cache prefix도 기존과 byte 단위로 동일하며 SHA-256은 `6e01f72f1803767b6bb16b5aef6cf34760e6dc1b293c316c75e990d09ad5fe63`이다. + +v4는 YAML 파싱과 FL0, FL1, FL3 embedded Python 컴파일을 통과하였다. 실제 v.7 source pack의 BO 48개를 fixture로 투입한 결과 48개 row가 새 schema를 통과했고 `amount` Object 48개가 값 손실 없이 deep copy되었다. `amount`와 `object_spec`의 JSON String 변환은 0건이며, 잘못된 ID, 필드 누락, String형 도출 근거, nested key 누락, BO universe 외 참조 등 7개 음성 사례는 모두 차단되었다. + +이 개정으로 문서상 format, FL1 생성값, FL3 통과 조건이 하나의 authoritative schema로 통합된다. Stage 2는 재파싱이나 의미 추정 없이 Object, Boolean, provenance 구조를 직접 소비할 수 있다. 이번 결과는 로컬 fixture 및 정적 검증 기준이며 원격 Liti-agent/MCP 통합 실행은 별도로 수행해야 한다. diff --git a/Case_02_Comparison_Research/YAML_Prompts/1. Stage_1/v.7/Stage_1_Part_4_Codex_v4.yml b/Case_02_Comparison_Research/YAML_Prompts/1. Stage_1/v.7/Stage_1_Part_4_Codex_v4.yml new file mode 100644 index 00000000..6ac59575 --- /dev/null +++ b/Case_02_Comparison_Research/YAML_Prompts/1. Stage_1/v.7/Stage_1_Part_4_Codex_v4.yml @@ -0,0 +1,2541 @@ +--- +Agent: + name: Liti-agent_Civil_Suit_Plaintiff_Stage_1_Part_4 + description: 민사소송 원고 송무 초지능 AI변호사 - Stage 1 Fact Ledger Generation + version: v2.0 + Stages: + - name: stage1_fact_ledger_generation + description: BO, Legal Effect 사용해서 Fact_Ledger_base.json 생성 + tools: + mcpServers: + localdocs: + type: streamable-http + url: http://mcp-localdocs:8012/mcp + description: Get the content of local documents + code-executor: + type: streamable-http + url: https://code-executor.mcp.eroomai.com/mcp + description: Run scripts of programming languages + headers: + Authorization: Bearer rR8OXqWrVZA1gFEo8oWfBkw2XgpWoGrrspw5ObsxTCM= + + task_procedure: + IN: + nexts: + - Task_FL0_fact_source_pack_compiler + wait_until: [] + Task_FL0_fact_source_pack_compiler: + nexts: + - Task_FL1_deterministic_fact_ledger_candidate_builder + wait_until: + - IN + Task_FL1_deterministic_fact_ledger_candidate_builder: + nexts: + - Task_FL2_fact_exception_adjudicator_* + - Task_FL3_final_fact_ledger_gate_and_writer + wait_until: + - Task_FL0_fact_source_pack_compiler + Task_FL2_fact_exception_adjudicator_*: + nexts: + - Task_FL3_final_fact_ledger_gate_and_writer + wait_until: + - Task_FL1_deterministic_fact_ledger_candidate_builder + Task_FL3_final_fact_ledger_gate_and_writer: + nexts: + - OUT + wait_until: + - Task_FL1_deterministic_fact_ledger_candidate_builder + - all Task_FL2_fact_exception_adjudicator_* + OUT: + nexts: [] + wait_until: + - Task_FL3_final_fact_ledger_gate_and_writer + + tasks: + - task_name: Task_FL0_fact_source_pack_compiler + mcp: code-executor + tool_name: run_code + parameters: + language: python + requirements: httpx + network: agent-network + timeout: 240 + code: | + #!/usr/bin/env python3 + from __future__ import annotations + + import copy + import hashlib + import itertools + import json + import re + from typing import Any + + import httpx + + LOCALDOCS_URL = "http://mcp-localdocs:8012/mcp" + MCP_HEADERS = { + "Content-Type": "application/json", + "Accept": "application/json, text/event-stream", + } + CLIENT = httpx.Client(timeout=60) + MSG_ID_COUNTER = itertools.count(10) + JSON_DECODER = json.JSONDecoder() + + + def _next_msg_id() -> int: + return next(MSG_ID_COUNTER) + + + def _safe_error(exc: BaseException) -> str: + return re.sub(r"\s+", " ", str(exc)).strip()[:300] + + + def _parse_json_value(raw: Any, label: str, allow_trailing: bool = False) -> Any: + if isinstance(raw, (dict, list)): + return raw + if not isinstance(raw, str): + raise ValueError(f"{label} is not JSON text") + text = raw.lstrip("\ufeff").strip() + if not text: + raise ValueError(f"{label} is empty") + try: + return json.loads(text) + except json.JSONDecodeError: + try: + value, end = JSON_DECODER.raw_decode(text) + except json.JSONDecodeError as exc: + raise ValueError(f"{label} invalid JSON: {text[:300]}") from exc + if not allow_trailing and text[end:].strip(): + raise ValueError(f"{label} has trailing content: {text[end:end + 300]}") + return value + + + def _parse_mcp_response(text: str) -> dict[str, Any]: + parsed: dict[str, Any] | None = None + for line in text.strip().splitlines(): + if not line.startswith("data: "): + continue + try: + item = _parse_json_value(line[6:], "MCP SSE data", True) + except ValueError: + continue + if isinstance(item, dict): + parsed = item + if parsed is None: + item = _parse_json_value(text, "MCP response", True) + if not isinstance(item, dict): + raise RuntimeError("MCP response root is not an object") + parsed = item + return parsed + + + def _init_mcp_session() -> None: + response = CLIENT.post( + LOCALDOCS_URL, + json={ + "jsonrpc": "2.0", + "id": 1, + "method": "initialize", + "params": { + "protocolVersion": "2025-03-26", + "capabilities": {}, + "clientInfo": { + "name": "stage-1-fact-ledger-fl0-source-compiler", + "version": "2.0", + "user_id": "{{__user_hash__}}", + "workspace_id": "{{__workspace_hash__}}", + }, + }, + }, + headers=MCP_HEADERS, + ) + response.raise_for_status() + payload = _parse_mcp_response(response.text) + if payload.get("error"): + raise RuntimeError(f"MCP initialize failed: {payload['error']}") + session_id = response.headers.get("mcp-session-id") + if session_id: + MCP_HEADERS["mcp-session-id"] = session_id + initialized = CLIENT.post( + LOCALDOCS_URL, + json={"jsonrpc": "2.0", "method": "notifications/initialized"}, + headers=MCP_HEADERS, + ) + initialized.raise_for_status() + + + def _call_tool(name: str, arguments: dict[str, Any]) -> dict[str, Any]: + response = CLIENT.post( + LOCALDOCS_URL, + json={ + "jsonrpc": "2.0", + "id": _next_msg_id(), + "method": "tools/call", + "params": {"name": name, "arguments": arguments}, + }, + headers=MCP_HEADERS, + ) + response.raise_for_status() + payload = _parse_mcp_response(response.text) + if payload.get("error"): + raise RuntimeError(f"MCP {name} failed: {payload['error']}") + result = payload.get("result") + if not isinstance(result, dict) or result.get("isError") is True: + raise RuntimeError(f"MCP {name} returned an invalid result") + return result + + + def _extract_doc(result: dict[str, Any], doc_name: str) -> Any: + blocks = result.get("content") or [] + if not blocks or not isinstance(blocks[0], dict): + raise RuntimeError(f"Empty response: {doc_name}") + text = blocks[0].get("text") + if not isinstance(text, str) or not text.strip(): + raise RuntimeError(f"Empty text response: {doc_name}") + outer = _parse_json_value(text, f"{doc_name} outer envelope", True) + inner: Any = outer + if isinstance(outer, dict) and "results" in outer: + results = outer.get("results") or [] + if not results or not isinstance(results[0], dict): + raise RuntimeError(f"Empty results envelope: {doc_name}") + inner = results[0].get("content") + if inner in (None, ""): + inner = results[0].get("text") + if isinstance(inner, (dict, list)): + return inner + return _parse_json_value(inner, f"{doc_name} inner document", False) + + + def _read_json(path: str) -> Any: + return _extract_doc(_call_tool("read_docs", {"doc_names": [path]}), path) + + + def _write_text(path: str, content: str) -> None: + _call_tool("write_file", {"path": path, "content": content, "overwrite": True}) + + + def _canonical(value: Any) -> str: + return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + + + def _digest(value: Any) -> str: + return "sha256:" + hashlib.sha256(_canonical(value).encode("utf-8")).hexdigest() + + + def _write_json(path: str, value: Any) -> None: + _write_text(path, json.dumps(value, ensure_ascii=False, sort_keys=True, indent=2)) + + + def _verify_reread(path: str, expected: Any) -> None: + actual = _read_json(path) + if _digest(actual) != _digest(expected): + raise RuntimeError(f"Post-write verification failed: {path}") + + + def _as_dict(value: Any) -> dict[str, Any]: + return value if isinstance(value, dict) else {} + + + def _as_list(value: Any) -> list[Any]: + return value if isinstance(value, list) else [] + + + def _strings(value: Any) -> list[str]: + values = value if isinstance(value, list) else [value] + out: list[str] = [] + for item in values: + if item is None: + continue + if isinstance(item, str): + text = item.strip() + elif isinstance(item, (int, float)) and not isinstance(item, bool): + text = str(item) + else: + continue + if text and text not in out: + out.append(text) + return out + + + def _sorted_strings(value: Any) -> list[str]: + return sorted(set(_strings(value))) + + + def _text(value: Any, limit: int = 500) -> str: + if value in (None, "", [], {}): + return "" + raw = value if isinstance(value, str) else json.dumps(value, ensure_ascii=False, sort_keys=True) + return re.sub(r"\s+", " ", raw).strip()[:limit] + + INPUT_FILES = { + "bo": "BO.json", + "structures": "legal_effect_structures.json", + "actio": "actio_case_signals.json", + "liability": "case_liability_signals.json", + "legal_effect": "legal_effect_signals.json", + "evidence": "evidence_indexed.json", + } + OUT_PATH = "stage1_tmp/fact_ledger/fact_source_pack.json" + STRUCTURE_INDEX_KEYS = ( + "by_bo_id", + "by_issue_cluster_id", + "by_asset_cluster_id", + "by_evidence_index", + "by_candidate_structure_type", + ) + DOMAIN_KEYS = ( + "money_claim_details", + "possession_state", + "lien_state", + "notice_lifecycle", + "asset_alias", + "fractional_effect_events", + ) + + + def _bo_id(bo: dict[str, Any]) -> str: + return _text(bo.get("BO_ID") or bo.get("id"), 120) + + + def _evidence_id(row: dict[str, Any]) -> str: + return _text(row.get("evidence_index") or row.get("evidence_id") or row.get("id"), 120) + + + def _sanitize(value: Any, depth: int = 0) -> Any: + if depth > 8: + return None + if isinstance(value, dict): + return {str(k): _sanitize(v, depth + 1) for k, v in sorted(value.items())} + if isinstance(value, list): + out: list[Any] = [] + seen: set[str] = set() + for item in value: + compact = _sanitize(item, depth + 1) + key = _canonical(compact) + if key not in seen: + seen.add(key) + out.append(compact) + return out + if isinstance(value, str): + return _text(value, 500) + if value is None or isinstance(value, (bool, int, float)): + return value + return _text(value, 500) + + + def _bo_evidence_refs(bo: dict[str, Any]) -> list[str]: + refs = _strings(bo.get("source_evidence_indexes")) + for row in _as_list(bo.get("Evidence")): + if isinstance(row, dict): + refs.extend(_strings(row.get("evidence_index"))) + return sorted(set(refs)) + + + def _compact_bo(bo: dict[str, Any]) -> dict[str, Any]: + core = _as_dict(bo.get("core_field_base")) + payload = _as_dict(_as_dict(bo.get("extensions")).get("domain_payload")) + domain = {key: _sanitize(payload.get(key)) for key in DOMAIN_KEYS if payload.get(key) not in (None, {}, [])} + juristic = bo.get("JuristicAct") + return { + "bo_id": _bo_id(bo), + "bo_type": bo.get("BOType") or bo.get("bo_type"), + "action": bo.get("Action") or bo.get("action"), + "action_type": bo.get("ActionType") or bo.get("action_type"), + "juristic_act": _sanitize(juristic), + "reason": _text(bo.get("Reason") or bo.get("reason"), 500) or None, + "prior_act": _sanitize(bo.get("PriorAct") or bo.get("prior_act")), + "legal_keywords": _sorted_strings(bo.get("Legal_Keywords") or bo.get("legal_keywords")), + "amount": _sanitize(bo.get("amount") if bo.get("amount") is not None else core.get("Amount")), + "date": core.get("BehaviorTime") or bo.get("date"), + "core_field_base": { + key: _sanitize(core.get(key)) + for key in ("BehaviorTime", "Performer", "Subject", "Object", "Amount", "Action") + }, + "source_evidence_indexes": _bo_evidence_refs(bo), + "domain_payload_compact": domain, + } + + + def _compact_structure(row: dict[str, Any]) -> dict[str, Any]: + audit = _as_dict(row.get("routing_audit")) + return { + "structure_id": _text(row.get("structure_id"), 120), + "candidate_structure_type": _text(row.get("candidate_structure_type"), 120) or None, + "candidate_structure_types": _sorted_strings(audit.get("candidate_structure_types")), + "structure_label": _text(row.get("structure_label"), 300), + "source_bo_ids": _sorted_strings(row.get("source_bo_ids")), + "issue_cluster_ids": _sorted_strings(row.get("issue_cluster_ids")), + "asset_cluster_ids": _sorted_strings(row.get("asset_cluster_ids")), + "evidence_indexes": _sorted_strings(row.get("evidence_indexes")), + "linked_claim_group_ids": _sorted_strings(row.get("linked_claim_group_ids")), + "linked_liability_group_ids": _sorted_strings(row.get("linked_liability_group_ids")), + "linked_actio_signal_ids": _sorted_strings(row.get("linked_actio_signal_ids")), + "legal_effect_roles": _sorted_strings(row.get("legal_effect_roles")), + "review_required": row.get("review_required") is True, + "routing_audit": { + "blocked_exception_ids": _sorted_strings(audit.get("blocked_exception_ids")), + "deterministic_review_ids": _sorted_strings(audit.get("deterministic_review_ids")), + }, + } + + + def _compact_evidence(row: dict[str, Any]) -> dict[str, Any]: + def compact_values(value: Any) -> list[str]: + out: list[str] = [] + for item in _as_list(value): + text = _text(item, 500) + if text and text not in out: + out.append(text) + return out + + return { + "evidence_index": _evidence_id(row), + "title": _text(row.get("title") or row.get("source_title"), 500), + "doc_type": _text(row.get("doc_type"), 160) or None, + "key_facts": compact_values(row.get("key_facts")), + "key_dates": compact_values(row.get("key_dates")), + "key_amounts": compact_values(row.get("key_amounts")), + "source_pointer": _sanitize(row.get("source_pointer")) if isinstance(row.get("source_pointer"), dict) else {}, + "authentication": _sanitize(row.get("authentication") or row.get("authentication_status")), + } + + + def _signal_rows(doc: dict[str, Any], root: str, schema: str) -> list[dict[str, Any]]: + if doc.get("schema_version") != schema: + raise ValueError(f"{root} schema_version mismatch") + if doc.get("status") not in {"READY", "READY_WITH_REVIEW"}: + raise ValueError(f"{root} status invalid") + rows = doc.get(root) + if not isinstance(rows, list): + raise ValueError(f"{root} must be an array") + return [row for row in rows if isinstance(row, dict)] + + + def _compact_signal(row: dict[str, Any], signal_type: str) -> dict[str, Any]: + return { + "signal_type": signal_type, + "signal_id": _text(row.get("signal_id") or row.get("id"), 120) or None, + "candidate_structure_types": _sorted_strings(row.get("candidate_structure_types") or row.get("candidate_structure_type")), + "linked_structure_ids": _sorted_strings(row.get("linked_structure_ids")), + "related_evidence_indexes": _sorted_strings(row.get("related_evidence_indexes") or row.get("evidence_indexes")), + "party_roles": _sanitize(_as_list(row.get("party_roles"))), + "actio_role_hints": _sanitize(_as_list(row.get("actio_role_hints"))), + "claim_group_candidates": _sanitize(_as_list(row.get("claim_group_candidates")) + _strings(row.get("claim_group_id"))), + "liability_group_candidates": _sanitize(_as_list(row.get("liability_group_candidates")) + _strings(row.get("liability_group_id"))), + } + + + def _signal_bo_ids(row: dict[str, Any], signal_type: str) -> list[str]: + if signal_type == "legal_effect": + return _strings(row.get("bo_id") or row.get("related_bo_ids")) + return _strings(row.get("related_bo_ids") or row.get("source_bo_ids") or row.get("bo_id")) + + + def _review_for_bo(queue_item: dict[str, Any], bo_id: str) -> dict[str, Any]: + raw_type = _text(queue_item.get("review_type"), 120) + if raw_type == "blocked_exception": + review_type = "blocked_exception" + elif raw_type in {"deterministic_review", "legal_structure_deterministic_review"}: + review_type = "deterministic_review" + elif raw_type in {"review_sensitive_split", "split_review"}: + review_type = "review_sensitive_split" + else: + review_type = "structure_review" + structure_ids = _strings(queue_item.get("structure_ids") or queue_item.get("structure_id")) + return { + "review_type": review_type, + "review_id": _text(queue_item.get("review_id"), 120) or None, + "exception_id": _text(queue_item.get("exception_id"), 120) or None, + "review_code": _text(queue_item.get("review_code"), 160) or None, + "reason_code": _text(queue_item.get("reason_code") or queue_item.get("block_reason"), 160) or None, + "downstream_owner": _text(queue_item.get("downstream_owner"), 80) or ("HUMAN_REVIEW" if review_type == "blocked_exception" else "STAGE_2"), + "structure_ids": sorted(set(structure_ids)), + "blocking": review_type == "blocked_exception" or queue_item.get("blocking") is True, + "source_bo_id": bo_id, + } + + + def main() -> None: + _init_mcp_session() + docs = {key: _read_json(path) for key, path in INPUT_FILES.items()} + input_digests = {path: _digest(docs[key]) for key, path in INPUT_FILES.items()} + run_fingerprint = _digest({"contract": "stage1_fact_ledger_run.v3", "input_digests": input_digests}) + + bo_rows = docs["bo"] + if not isinstance(bo_rows, list) or not bo_rows: + raise ValueError("BO.json must be a non-empty array") + bo_by_id: dict[str, Any] = {} + for row in bo_rows: + if not isinstance(row, dict): + raise ValueError("BO.json contains a non-object item") + bo_id = _bo_id(row) + if not bo_id or bo_id in bo_by_id: + raise ValueError(f"BO ID missing or duplicated: {bo_id}") + bo_by_id[bo_id] = _compact_bo(row) + bo_ids = set(bo_by_id) + + structures_doc = _as_dict(docs["structures"]) + if structures_doc.get("schema_version") != "stage1_legal_effect_structures.v1": + raise ValueError("legal_effect_structures schema_version mismatch") + if structures_doc.get("status") not in {"READY", "READY_WITH_REVIEW"}: + raise ValueError("legal_effect_structures status invalid") + structures_raw = structures_doc.get("legal_effect_structures") + indexes = structures_doc.get("structure_index") + quality_gate = structures_doc.get("quality_gate") + if not isinstance(structures_raw, list) or not isinstance(indexes, dict) or not isinstance(quality_gate, dict): + raise ValueError("legal_effect_structures required roots missing") + for key in STRUCTURE_INDEX_KEYS: + if not isinstance(indexes.get(key), dict): + raise ValueError(f"structure_index.{key} missing") + for gate_name in ("bo_coverage", "evidence_coverage"): + if _as_dict(quality_gate.get(gate_name)).get("coverage_pass") is not True: + raise ValueError(f"quality_gate.{gate_name} failed") + + structure_by_id: dict[str, Any] = {} + expected_by_bo: dict[str, list[str]] = {bo_id: [] for bo_id in bo_ids} + for raw in structures_raw: + if not isinstance(raw, dict): + raise ValueError("legal_effect_structures contains a non-object item") + compact = _compact_structure(raw) + sid = compact["structure_id"] + if not sid or sid in structure_by_id: + raise ValueError(f"structure ID missing or duplicated: {sid}") + unknown_bos = sorted(set(compact["source_bo_ids"]) - bo_ids) + if unknown_bos: + raise ValueError(f"structure {sid} has unknown BO refs: {unknown_bos}") + structure_by_id[sid] = compact + for bo_id in compact["source_bo_ids"]: + expected_by_bo[bo_id].append(sid) + + by_bo = _as_dict(indexes.get("by_bo_id")) + for bo_id in sorted(bo_ids): + indexed = _sorted_strings(_as_dict(by_bo.get(bo_id)).get("structure_ids")) + expected = sorted(set(expected_by_bo.get(bo_id, []))) + if indexed != expected: + raise ValueError(f"structure_index.by_bo_id mismatch: {bo_id}") + + evidence_doc = _as_dict(docs["evidence"]) + evidence_rows = evidence_doc.get("items") or evidence_doc.get("evidence_indexed") + if not isinstance(evidence_rows, list) or not evidence_rows: + raise ValueError("evidence_indexed conservation array missing") + evidence_by_id: dict[str, Any] = {} + for raw in evidence_rows: + if not isinstance(raw, dict): + raise ValueError("evidence_indexed contains a non-object item") + compact = _compact_evidence(raw) + eid = compact["evidence_index"] + if not eid or eid in evidence_by_id: + raise ValueError(f"evidence ID missing or duplicated: {eid}") + evidence_by_id[eid] = compact + evidence_ids = set(evidence_by_id) + for sid, structure in structure_by_id.items(): + unknown = sorted(set(structure["evidence_indexes"]) - evidence_ids) + if unknown: + raise ValueError(f"structure {sid} has unknown evidence refs: {unknown}") + + signal_specs = ( + ("actio", _as_dict(docs["actio"]), "actio_case_signals", "actio_case_signals.v1"), + ("case_liability", _as_dict(docs["liability"]), "case_liability_signals", "case_liability_signals.v1"), + ("legal_effect", _as_dict(docs["legal_effect"]), "bo_legal_effect_routes", "legal_effect_signals.v1"), + ) + signal_hints: dict[str, dict[str, list[dict[str, Any]]]] = { + bo_id: {"actio": [], "case_liability": [], "legal_effect": []} for bo_id in bo_ids + } + source_schema_versions = { + "BO.json": "array", + "legal_effect_structures.json": structures_doc.get("schema_version"), + "evidence_indexed.json": evidence_doc.get("schema_contract_version") or evidence_doc.get("schema_version") or "conservation-array", + } + for signal_type, doc, root, schema in signal_specs: + rows = _signal_rows(doc, root, schema) + source_schema_versions[INPUT_FILES[{"actio": "actio", "case_liability": "liability", "legal_effect": "legal_effect"}[signal_type]]] = schema + for raw in rows: + compact = _compact_signal(raw, signal_type) + for bo_id in _signal_bo_ids(raw, signal_type): + if bo_id in signal_hints: + signal_hints[bo_id][signal_type].append(compact) + for bo_id in signal_hints: + for signal_type in signal_hints[bo_id]: + signal_hints[bo_id][signal_type].sort(key=_canonical) + + structure_refs_by_bo: dict[str, Any] = {} + for bo_id in sorted(bo_ids): + row = _as_dict(by_bo.get(bo_id)) + structure_refs_by_bo[bo_id] = { + "structure_ids": _sorted_strings(row.get("structure_ids")), + "issue_cluster_ids": _sorted_strings(row.get("issue_cluster_ids")), + "asset_cluster_ids": _sorted_strings(row.get("asset_cluster_ids")), + "evidence_index_ids": _sorted_strings(row.get("evidence_indexes")), + "candidate_structure_types": _sorted_strings(row.get("candidate_structure_types")), + "claim_group_ids": _sorted_strings(row.get("claim_group_ids")), + "liability_group_ids": _sorted_strings(row.get("liability_group_ids")), + "downstream_fact_roles": _sorted_strings(row.get("legal_effect_roles") or row.get("downstream_fact_roles")), + } + + upstream_reviews: dict[str, list[dict[str, Any]]] = {bo_id: [] for bo_id in bo_ids} + queue = _as_list(quality_gate.get("review_queue")) + for raw in queue: + if not isinstance(raw, dict): + continue + for bo_id in _strings(raw.get("source_bo_ids") or raw.get("source_bo_id")): + if bo_id in upstream_reviews: + upstream_reviews[bo_id].append(_review_for_bo(raw, bo_id)) + for sid, structure in structure_by_id.items(): + if structure.get("review_required") is not True: + continue + for bo_id in structure["source_bo_ids"]: + upstream_reviews[bo_id].append(_review_for_bo({"review_type": "structure_review", "structure_id": sid}, bo_id)) + for bo_id, rows in upstream_reviews.items(): + unique = {_canonical(row): row for row in rows} + upstream_reviews[bo_id] = sorted(unique.values(), key=lambda row: (_text(row.get("review_type")), _text(row.get("review_id") or row.get("exception_id")), _canonical(row))) + + warnings: list[dict[str, Any]] = [] + for bo_id, bo in bo_by_id.items(): + for eid in bo["source_evidence_indexes"]: + if eid not in evidence_ids: + warnings.append({"code": "UNKNOWN_BO_EVIDENCE_REF", "source_bo_id": bo_id, "evidence_index": eid}) + review_count = sum(len(rows) for rows in upstream_reviews.values()) + blocking_count = sum(1 for rows in upstream_reviews.values() for row in rows if row.get("blocking") is True) + status = "READY_WITH_REVIEW" if review_count or warnings else "READY" + source_pack = { + "schema_version": "stage1_fact_source_pack.v3", + "status": status, + "run_fingerprint": run_fingerprint, + "input_manifest": {"files": input_digests, "source_schema_versions": source_schema_versions}, + "bo_by_id": {bo_id: bo_by_id[bo_id] for bo_id in sorted(bo_by_id)}, + "structure_by_id": {sid: structure_by_id[sid] for sid in sorted(structure_by_id)}, + "structure_refs_by_bo_id": structure_refs_by_bo, + "signal_hints_by_bo_id": {bo_id: signal_hints[bo_id] for bo_id in sorted(signal_hints)}, + "evidence_authority_by_index": {eid: evidence_by_id[eid] for eid in sorted(evidence_by_id)}, + "upstream_review_by_bo_id": {bo_id: upstream_reviews[bo_id] for bo_id in sorted(upstream_reviews)}, + "reference_universe": { + "bo_ids": sorted(bo_ids), + "structure_ids": sorted(structure_by_id), + "evidence_indexes": sorted(evidence_ids), + }, + "source_integrity": { + "warnings": sorted(warnings, key=_canonical), + "warning_count": len(warnings), + "reference_validation_pass": True, + }, + "source_counts": { + "bo_count": len(bo_by_id), + "structure_count": len(structure_by_id), + "evidence_count": len(evidence_by_id), + "review_count": review_count, + "blocking_review_count": blocking_count, + }, + } + payload = {"fact_source_pack": source_pack} + _write_json(OUT_PATH, payload) + _verify_reread(OUT_PATH, payload) + print(_canonical({"fact_source_pack_manifest": { + "schema_version": "stage1_fact_source_pack_manifest.v2", + "status": status, + "path": OUT_PATH, + "sha256": _digest(payload), + "run_fingerprint": run_fingerprint, + "bo_count": len(bo_by_id), + "structure_count": len(structure_by_id), + "evidence_count": len(evidence_by_id), + "review_count": review_count, + }})) + + + try: + main() + except Exception as exc: + print(_canonical({"fact_source_pack_manifest": { + "schema_version": "stage1_fact_source_pack_manifest.v2", + "status": "FAILED", + "error": _safe_error(exc), + }})) + raise + + - task_name: Task_FL1_deterministic_fact_ledger_candidate_builder + mcp: code-executor + tool_name: run_code + parameters: + language: python + requirements: |- + httpx + jsonschema>=4.23.0 + network: agent-network + timeout: 240 + code: | + #!/usr/bin/env python3 + from __future__ import annotations + + import copy + import hashlib + import itertools + import json + import re + from typing import Any + + import httpx + from jsonschema import Draft202012Validator + from jsonschema.exceptions import ValidationError + + LOCALDOCS_URL = "http://mcp-localdocs:8012/mcp" + MCP_HEADERS = { + "Content-Type": "application/json", + "Accept": "application/json, text/event-stream", + } + CLIENT = httpx.Client(timeout=60) + MSG_ID_COUNTER = itertools.count(10) + JSON_DECODER = json.JSONDecoder() + + + def _next_msg_id() -> int: + return next(MSG_ID_COUNTER) + + + def _safe_error(exc: BaseException) -> str: + return re.sub(r"\s+", " ", str(exc)).strip()[:300] + + + def _parse_json_value(raw: Any, label: str, allow_trailing: bool = False) -> Any: + if isinstance(raw, (dict, list)): + return raw + if not isinstance(raw, str): + raise ValueError(f"{label} is not JSON text") + text = raw.lstrip("\ufeff").strip() + if not text: + raise ValueError(f"{label} is empty") + try: + return json.loads(text) + except json.JSONDecodeError: + try: + value, end = JSON_DECODER.raw_decode(text) + except json.JSONDecodeError as exc: + raise ValueError(f"{label} invalid JSON: {text[:300]}") from exc + if not allow_trailing and text[end:].strip(): + raise ValueError(f"{label} has trailing content: {text[end:end + 300]}") + return value + + + def _parse_mcp_response(text: str) -> dict[str, Any]: + parsed: dict[str, Any] | None = None + for line in text.strip().splitlines(): + if not line.startswith("data: "): + continue + try: + item = _parse_json_value(line[6:], "MCP SSE data", True) + except ValueError: + continue + if isinstance(item, dict): + parsed = item + if parsed is None: + item = _parse_json_value(text, "MCP response", True) + if not isinstance(item, dict): + raise RuntimeError("MCP response root is not an object") + parsed = item + return parsed + + + def _init_mcp_session() -> None: + response = CLIENT.post( + LOCALDOCS_URL, + json={ + "jsonrpc": "2.0", + "id": 1, + "method": "initialize", + "params": { + "protocolVersion": "2025-03-26", + "capabilities": {}, + "clientInfo": { + "name": "stage-1-fact-ledger-fl1-candidate-builder", + "version": "2.0", + "user_id": "{{__user_hash__}}", + "workspace_id": "{{__workspace_hash__}}", + }, + }, + }, + headers=MCP_HEADERS, + ) + response.raise_for_status() + payload = _parse_mcp_response(response.text) + if payload.get("error"): + raise RuntimeError(f"MCP initialize failed: {payload['error']}") + session_id = response.headers.get("mcp-session-id") + if session_id: + MCP_HEADERS["mcp-session-id"] = session_id + initialized = CLIENT.post( + LOCALDOCS_URL, + json={"jsonrpc": "2.0", "method": "notifications/initialized"}, + headers=MCP_HEADERS, + ) + initialized.raise_for_status() + + + def _call_tool(name: str, arguments: dict[str, Any]) -> dict[str, Any]: + response = CLIENT.post( + LOCALDOCS_URL, + json={ + "jsonrpc": "2.0", + "id": _next_msg_id(), + "method": "tools/call", + "params": {"name": name, "arguments": arguments}, + }, + headers=MCP_HEADERS, + ) + response.raise_for_status() + payload = _parse_mcp_response(response.text) + if payload.get("error"): + raise RuntimeError(f"MCP {name} failed: {payload['error']}") + result = payload.get("result") + if not isinstance(result, dict) or result.get("isError") is True: + raise RuntimeError(f"MCP {name} returned an invalid result") + return result + + + def _extract_doc(result: dict[str, Any], doc_name: str) -> Any: + blocks = result.get("content") or [] + if not blocks or not isinstance(blocks[0], dict): + raise RuntimeError(f"Empty response: {doc_name}") + text = blocks[0].get("text") + if not isinstance(text, str) or not text.strip(): + raise RuntimeError(f"Empty text response: {doc_name}") + outer = _parse_json_value(text, f"{doc_name} outer envelope", True) + inner: Any = outer + if isinstance(outer, dict) and "results" in outer: + results = outer.get("results") or [] + if not results or not isinstance(results[0], dict): + raise RuntimeError(f"Empty results envelope: {doc_name}") + inner = results[0].get("content") + if inner in (None, ""): + inner = results[0].get("text") + if isinstance(inner, (dict, list)): + return inner + return _parse_json_value(inner, f"{doc_name} inner document", False) + + + def _read_json(path: str) -> Any: + return _extract_doc(_call_tool("read_docs", {"doc_names": [path]}), path) + + + def _write_text(path: str, content: str) -> None: + _call_tool("write_file", {"path": path, "content": content, "overwrite": True}) + + + def _canonical(value: Any) -> str: + return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + + + def _digest(value: Any) -> str: + return "sha256:" + hashlib.sha256(_canonical(value).encode("utf-8")).hexdigest() + + + def _write_json(path: str, value: Any) -> None: + _write_text(path, json.dumps(value, ensure_ascii=False, sort_keys=True, indent=2)) + + + def _verify_reread(path: str, expected: Any) -> None: + actual = _read_json(path) + if _digest(actual) != _digest(expected): + raise RuntimeError(f"Post-write verification failed: {path}") + + + def _as_dict(value: Any) -> dict[str, Any]: + return value if isinstance(value, dict) else {} + + + def _as_list(value: Any) -> list[Any]: + return value if isinstance(value, list) else [] + + + def _strings(value: Any) -> list[str]: + values = value if isinstance(value, list) else [value] + out: list[str] = [] + for item in values: + if item is None: + continue + if isinstance(item, str): + text = item.strip() + elif isinstance(item, (int, float)) and not isinstance(item, bool): + text = str(item) + else: + continue + if text and text not in out: + out.append(text) + return out + + + def _sorted_strings(value: Any) -> list[str]: + return sorted(set(_strings(value))) + + + def _text(value: Any, limit: int = 500) -> str: + if value in (None, "", [], {}): + return "" + raw = value if isinstance(value, str) else json.dumps(value, ensure_ascii=False, sort_keys=True) + return re.sub(r"\s+", " ", raw).strip()[:limit] + + SOURCE_PACK_PATH = "stage1_tmp/fact_ledger/fact_source_pack.json" + CANDIDATE_PATH = "stage1_tmp/fact_ledger/Fact_Ledger_base_candidate.json" + ROW_SCHEMA_PATH = "stage1_tmp/fact_ledger/fact_ledger_row_schema.json" + EXCEPTION_MANIFEST_PATH = "stage1_tmp/fact_ledger/fact_exception_manifest.json" + INPUT_PART_DIR = "stage1_tmp/fact_ledger/fact_exception_input_parts" + OUTPUT_PART_DIR = "stage1_tmp/fact_ledger/fact_exception_adjudication_parts" + MAX_EXCEPTION_ITEMS_PER_PACK = 4 + TARGET_EXCEPTION_ITEMS_PER_PACK = 3 + MAX_EXCEPTION_ITEM_BYTES = 12000 + MAX_EXCEPTION_PACK_BYTES = 25000 + MAX_EXCEPTION_PACKS = 8 + MAX_CONCURRENCY = 2 + FINAL_FIELDS = [ + "fact_id", "source_bo_id", "type", "date", "parties", "object_spec", "amount", "action", + "evidence_refs", "credibility", "legal_centrality", "proof_strength", "legal_effect_roles", + "claim_chain_ref", "linked_structures", "must_not_drop_in_claim_types", "must_consider", + "state_context", "money_claim_effect", "commercial_successor_effect", "actio_roles", + "linked_actio_structures", "downstream_module_candidates", "skeleton_validation_codes", + "derived_fact", "is_derived", "derivation_basis", "legal_calculation_object", + ] + FINAL_FIELD_SET = set(FINAL_FIELDS) + FACT_LEDGER_ROW_SCHEMA_ID = "https://eroomai.com/schemas/stage1/fact-ledger-row.v4.json" + FACT_LEDGER_ROW_SCHEMA = { + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": FACT_LEDGER_ROW_SCHEMA_ID, + "title": "Stage 1 Fact Ledger Row", + "type": "object", + "required": FINAL_FIELDS, + "additionalProperties": False, + "$defs": { + "string_array": { + "type": "array", + "items": {"type": "string", "minLength": 1}, + "uniqueItems": True, + }, + "effect_flag": { + "type": "object", + "required": ["enabled", "source_bo_id"], + "additionalProperties": False, + "properties": { + "enabled": {"type": "boolean"}, + "source_bo_id": {"type": ["string", "null"]}, + }, + }, + }, + "properties": { + "fact_id": {"type": "string", "pattern": "^F-[0-9]{3}$"}, + "source_bo_id": {"type": "string", "minLength": 1}, + "type": {"type": ["string", "null"]}, + "date": {"type": ["string", "null"]}, + "parties": { + "type": "array", + "items": { + "type": "object", + "required": ["party", "role", "source"], + "additionalProperties": False, + "properties": { + "party": {"type": "string", "minLength": 1}, + "role": {"type": "string"}, + "source": {"type": "string", "minLength": 1}, + }, + }, + }, + "object_spec": { + "oneOf": [ + {"type": "null"}, + { + "type": "object", + "required": ["source_object", "asset_cluster_ids"], + "additionalProperties": False, + "properties": { + "source_object": {}, + "asset_cluster_ids": {"$ref": "#/$defs/string_array"}, + }, + }, + ], + }, + "amount": {"type": ["object", "null"]}, + "action": {"type": ["string", "null"]}, + "evidence_refs": {"$ref": "#/$defs/string_array"}, + "credibility": {"type": "string", "enum": ["high", "medium", "low"]}, + "legal_centrality": {"type": "string", "enum": ["critical", "medium"]}, + "proof_strength": { + "type": "string", + "enum": ["high", "medium", "low", "meeting_only_core", "documentary"], + }, + "legal_effect_roles": {"$ref": "#/$defs/string_array"}, + "claim_chain_ref": { + "type": "object", + "required": ["claim_group_ids", "liability_group_ids", "signal_count"], + "additionalProperties": False, + "properties": { + "claim_group_ids": {"$ref": "#/$defs/string_array"}, + "liability_group_ids": {"$ref": "#/$defs/string_array"}, + "signal_count": {"type": "integer", "minimum": 0}, + }, + }, + "linked_structures": { + "type": "object", + "required": [ + "structure_ids", "issue_cluster_ids", "asset_cluster_ids", "evidence_index_ids", + "candidate_structure_types", "claim_group_ids", "liability_group_ids", + ], + "additionalProperties": False, + "properties": { + "structure_ids": {"$ref": "#/$defs/string_array"}, + "issue_cluster_ids": {"$ref": "#/$defs/string_array"}, + "asset_cluster_ids": {"$ref": "#/$defs/string_array"}, + "evidence_index_ids": {"$ref": "#/$defs/string_array"}, + "candidate_structure_types": {"$ref": "#/$defs/string_array"}, + "claim_group_ids": {"$ref": "#/$defs/string_array"}, + "liability_group_ids": {"$ref": "#/$defs/string_array"}, + }, + }, + "must_not_drop_in_claim_types": {"$ref": "#/$defs/string_array"}, + "must_consider": {"type": "boolean"}, + "state_context": { + "type": "object", + "required": ["bo_type", "issue_cluster_ids", "asset_cluster_ids"], + "additionalProperties": False, + "properties": { + "bo_type": {"type": ["string", "null"]}, + "issue_cluster_ids": {"$ref": "#/$defs/string_array"}, + "asset_cluster_ids": {"$ref": "#/$defs/string_array"}, + }, + }, + "money_claim_effect": {"$ref": "#/$defs/effect_flag"}, + "commercial_successor_effect": {"$ref": "#/$defs/effect_flag"}, + "actio_roles": {"type": "array", "items": {"type": "object"}}, + "linked_actio_structures": { + "type": "object", + "required": ["structure_ids", "actio_signal_count"], + "additionalProperties": False, + "properties": { + "structure_ids": {"$ref": "#/$defs/string_array"}, + "actio_signal_count": {"type": "integer", "minimum": 0}, + }, + }, + "downstream_module_candidates": {"$ref": "#/$defs/string_array"}, + "skeleton_validation_codes": {"$ref": "#/$defs/string_array"}, + "derived_fact": {"type": ["string", "null"]}, + "is_derived": {"type": "boolean"}, + "derivation_basis": { + "type": "array", + "items": { + "type": "object", + "required": ["source_type", "source_id", "source_field", "role"], + "additionalProperties": False, + "properties": { + "source_type": { + "type": "string", + "enum": ["bo", "evidence", "legal_effect_structure", "review"], + }, + "source_id": {"type": "string", "minLength": 1}, + "source_field": {"type": "string", "minLength": 1}, + "role": { + "type": "string", + "enum": ["primary", "support", "structure_link", "review"], + }, + }, + }, + }, + "legal_calculation_object": { + "oneOf": [ + {"type": "null"}, + { + "type": "object", + "required": [ + "calculation_required", "amount_source", "amount_text", "object_basis", + "candidate_structure_types", + ], + "additionalProperties": False, + "properties": { + "calculation_required": {"type": "boolean"}, + "amount_source": {"type": ["string", "null"]}, + "amount_text": {"type": ["string", "null"]}, + "object_basis": {"type": ["object", "null"]}, + "candidate_structure_types": {"$ref": "#/$defs/string_array"}, + }, + }, + ], + }, + }, + } + Draft202012Validator.check_schema(FACT_LEDGER_ROW_SCHEMA) + FACT_LEDGER_ROW_VALIDATOR = Draft202012Validator(FACT_LEDGER_ROW_SCHEMA) + STRUCTURE_TYPE_TO_MODULE = { + "money_claim": "money_claim", + "commercial_successor": "commercial_successor", + "secured_debt": "secured_registry", + "registry_invalidity": "secured_registry", + "land_use_gain": "land_valuation", + "valuation": "land_valuation", + "fraudulent_transfer": "actio", + "preserved_claim": "actio", + "succession_notice_lien_asset_defense": "succession_notice_lien_asset_defense", + "general_legal_effect": "legal_effect", + } + STRUCTURE_TYPE_TO_ROLE = { + "money_claim": "claim_chain_ref", + "commercial_successor": "claim_chain_ref", + "secured_debt": "linked_structures", + "registry_invalidity": "linked_structures", + "land_use_gain": "legal_calculation_object", + "valuation": "legal_calculation_object", + "fraudulent_transfer": "actio_roles", + "preserved_claim": "actio_roles", + "succession_notice_lien_asset_defense": "linked_structures", + "general_legal_effect": "legal_effect_roles", + } + CLOSED_EXCEPTION_TYPES = { + "conflicting_fact_projection", + "conflicting_calculation_projection", + "conflicting_asset_projection", + "conflicting_structure_role_projection", + } + CLOSED_TARGET_FIELDS = { + "date", "amount", "action", "object_spec", "derived_fact", "claim_chain_ref", + "linked_structures", "state_context", "legal_calculation_object", + } + + + def _first(*values: Any) -> Any: + for value in values: + if value not in (None, "", [], {}): + return value + return None + + + def _structured_amount(value: Any) -> dict[str, Any] | None: + if value in (None, "", [], {}): + return None + if isinstance(value, dict): + return copy.deepcopy(value) + return {"source_value": copy.deepcopy(value)} + + + def _object_spec(core: dict[str, Any], linked: dict[str, Any]) -> dict[str, Any] | None: + source_object = copy.deepcopy(core.get("Object")) + if source_object in ("", [], {}): + source_object = None + asset_cluster_ids = _sorted_strings(linked.get("asset_cluster_ids")) + if source_object is None and not asset_cluster_ids: + return None + return { + "source_object": source_object, + "asset_cluster_ids": asset_cluster_ids, + } + + + def _derivation_basis(bo_id: str, evidence_refs: list[str], structure_ids: list[str], is_derived: bool) -> list[dict[str, Any]]: + basis = [{ + "source_type": "bo", + "source_id": bo_id, + "source_field": "core_field_base" if is_derived else "action", + "role": "primary", + }] + basis.extend({ + "source_type": "evidence", + "source_id": evidence_id, + "source_field": "source_evidence_indexes", + "role": "support", + } for evidence_id in evidence_refs) + basis.extend({ + "source_type": "legal_effect_structure", + "source_id": structure_id, + "source_field": "linked_structures.structure_ids", + "role": "structure_link", + } for structure_id in structure_ids) + return basis + + + def _looks_like_json_container_string(value: Any) -> bool: + if not isinstance(value, str) or not value.strip().startswith(("{", "[")): + return False + try: + return isinstance(json.loads(value), (dict, list)) + except json.JSONDecodeError: + return False + + + def _validate_row_schema(row: dict[str, Any], label: str) -> None: + for field in ("amount", "object_spec"): + if _looks_like_json_container_string(row.get(field)): + raise ValueError(f"{label}.{field} contains a JSON-encoded String") + errors: list[ValidationError] = sorted( + FACT_LEDGER_ROW_VALIDATOR.iter_errors(row), + key=lambda item: ([str(part) for part in item.absolute_path], item.message), + ) + if errors: + error = errors[0] + path = ".".join(str(part) for part in error.absolute_path) or "$" + raise ValueError(f"{label} JSON Schema violation at {path}: {error.message}") + + + def _add_code(row: dict[str, Any], code: str) -> None: + codes = _strings(row.get("skeleton_validation_codes")) + if code not in codes: + codes.append(code) + row["skeleton_validation_codes"] = sorted(codes) + + + def _score_evidence(rows: list[dict[str, Any]]) -> tuple[str, str, str]: + if not rows: + return "low", "low", "NO_EVIDENCE" + texts = " ".join( + " ".join([_text(row.get("title"), 200), _text(row.get("doc_type"), 100)] + _strings(row.get("key_facts"))) + for row in rows + ).lower() + meeting = any(marker in texts for marker in ("meeting", "회의", "면담")) + documentary = any(marker in texts for marker in ("등기", "계약", "판결", "결정", "공문", "금융", "송금", "계좌", "registry", "contract", "judgment")) + has_structured_fact = any(_as_list(row.get("key_facts")) for row in rows) + has_pointer = any(isinstance(row.get("source_pointer"), dict) and row.get("source_pointer") for row in rows) + if meeting and not documentary and all("meeting" in _text(row.get("doc_type"), 100).lower() or "회의" in _text(row.get("doc_type"), 100) or "면담" in _text(row.get("doc_type"), 100) for row in rows): + return "medium", "meeting_only_core", "MEETING_ONLY" + if documentary and has_structured_fact and has_pointer: + return "high", "documentary", "DOCUMENTARY_STRUCTURED" + if len(rows) >= 2 and has_structured_fact: + return "high", "high", "MULTI_SOURCE_STRUCTURED" + return "medium", "medium", "SINGLE_STRUCTURED_SOURCE" + + + def _linked_view(linked: dict[str, Any]) -> dict[str, Any]: + return { + "structure_ids": _sorted_strings(linked.get("structure_ids")), + "issue_cluster_ids": _sorted_strings(linked.get("issue_cluster_ids")), + "asset_cluster_ids": _sorted_strings(linked.get("asset_cluster_ids")), + "evidence_index_ids": _sorted_strings(linked.get("evidence_index_ids")), + "candidate_structure_types": _sorted_strings(linked.get("candidate_structure_types")), + "claim_group_ids": _sorted_strings(linked.get("claim_group_ids")), + "liability_group_ids": _sorted_strings(linked.get("liability_group_ids")), + } + + + def _legal_roles(types: list[str], structures: list[dict[str, Any]], signals: dict[str, Any]) -> list[str]: + roles: list[str] = [] + for structure in structures: + roles.extend(_strings(structure.get("legal_effect_roles"))) + roles.extend(STRUCTURE_TYPE_TO_ROLE.get(item, "legal_effect_roles") for item in types) + if _as_list(signals.get("actio")): + roles.append("actio_roles") + if _as_list(signals.get("case_liability")): + roles.append("claim_chain_ref") + if _as_list(signals.get("legal_effect")): + roles.append("legal_effect_roles") + return sorted(set(roles)) + + + def _modules(types: list[str], signals: dict[str, Any]) -> list[str]: + modules = [STRUCTURE_TYPE_TO_MODULE.get(item, "legal_effect") for item in types] + if _as_list(signals.get("actio")): + modules.append("actio") + if _as_list(signals.get("case_liability")): + modules.append("claim_liability") + return sorted(set(modules)) + + + def _parties(bo: dict[str, Any], signals: dict[str, Any]) -> list[dict[str, Any]]: + core = _as_dict(bo.get("core_field_base")) + out: list[dict[str, Any]] = [] + for role, field in (("actor", "Performer"), ("counterparty", "Subject")): + value = _text(core.get(field), 160) + if value: + out.append({"party": value, "role": role, "source": "BO.core_field_base"}) + for signal in _as_list(signals.get("case_liability")): + for item in _as_list(_as_dict(signal).get("party_roles")): + if not isinstance(item, dict): + continue + party = _text(item.get("party") or item.get("name"), 160) + if party: + out.append({"party": party, "role": _text(item.get("role"), 80), "source": "case_liability_signals"}) + unique = {(row["party"], row["role"]): row for row in out} + return [unique[key] for key in sorted(unique)] + + + def _claim_groups(linked: dict[str, Any], signals: dict[str, Any]) -> dict[str, Any]: + claims = _strings(linked.get("claim_group_ids")) + liabilities = _strings(linked.get("liability_group_ids")) + for signal in _as_list(signals.get("case_liability")): + claims.extend(_strings(_as_dict(signal).get("claim_group_candidates"))) + liabilities.extend(_strings(_as_dict(signal).get("liability_group_candidates"))) + return { + "claim_group_ids": sorted(set(claims)), + "liability_group_ids": sorted(set(liabilities)), + "signal_count": len(_as_list(signals.get("case_liability"))), + } + + + def _fallback_derived(bo: dict[str, Any], row: dict[str, Any]) -> str: + if _text(row.get("action"), 900): + return _text(row.get("action"), 900) + core = _as_dict(bo.get("core_field_base")) + fields = [core.get(key) for key in ("BehaviorTime", "Performer", "Subject", "Object", "Action", "Amount")] + parts = [_text(value, 240) for value in fields] + return " ".join(value for value in parts if value).strip() + + + def _review_template(code: str, row: dict[str, Any], reason: str, severity: str = "WARNING", owner: str = "STAGE_2") -> dict[str, Any]: + return { + "review_id": None, + "review_code": code, + "severity": severity, + "fact_id": row["fact_id"], + "source_bo_id": row["source_bo_id"], + "structure_ids": _strings(_as_dict(row.get("linked_structures")).get("structure_ids")), + "evidence_refs": _strings(row.get("evidence_refs")), + "reason_code": reason, + "downstream_owner": owner, + } + + + def _build_row(index: int, bo_id: str, source: dict[str, Any]) -> tuple[dict[str, Any], list[dict[str, Any]], list[dict[str, Any]], str]: + bo = _as_dict(_as_dict(source.get("bo_by_id")).get(bo_id)) + linked = _linked_view(_as_dict(_as_dict(source.get("structure_refs_by_bo_id")).get(bo_id))) + signals = _as_dict(_as_dict(source.get("signal_hints_by_bo_id")).get(bo_id)) + structure_map = _as_dict(source.get("structure_by_id")) + structures = [_as_dict(structure_map.get(sid)) for sid in linked["structure_ids"] if isinstance(structure_map.get(sid), dict)] + types = sorted(set(linked["candidate_structure_types"] + [str(row.get("candidate_structure_type")) for row in structures if row.get("candidate_structure_type")])) + evidence_refs = sorted(set(_strings(bo.get("source_evidence_indexes")) + linked["evidence_index_ids"])) + evidence_map = _as_dict(source.get("evidence_authority_by_index")) + evidence_rows = [_as_dict(evidence_map.get(eid)) for eid in evidence_refs if isinstance(evidence_map.get(eid), dict)] + credibility, proof_strength, score_basis = _score_evidence(evidence_rows) + core = _as_dict(bo.get("core_field_base")) + modules = _modules(types, signals) + amount_value = _first(bo.get("amount"), core.get("Amount")) + action_text = _text(_first(bo.get("action"), core.get("Action")), 900) or None + row = { + "fact_id": f"F-{index:03d}", + "source_bo_id": bo_id, + "type": _first(bo.get("bo_type"), types[0] if types else None), + "date": _first(core.get("BehaviorTime"), bo.get("date")), + "parties": _parties(bo, signals), + "object_spec": _object_spec(core, linked), + "amount": _structured_amount(amount_value), + "action": _first(bo.get("action"), core.get("Action")), + "evidence_refs": evidence_refs, + "credibility": credibility, + "legal_centrality": "critical" if linked["structure_ids"] or _as_list(signals.get("actio")) else "medium", + "proof_strength": proof_strength, + "legal_effect_roles": _legal_roles(types, structures, signals), + "claim_chain_ref": _claim_groups(linked, signals), + "linked_structures": linked, + "must_not_drop_in_claim_types": modules if linked["structure_ids"] or _as_list(signals.get("actio")) else [], + "must_consider": bool(linked["structure_ids"] or _as_list(signals.get("actio"))), + "state_context": {"bo_type": bo.get("bo_type"), "issue_cluster_ids": linked["issue_cluster_ids"], "asset_cluster_ids": linked["asset_cluster_ids"]}, + "money_claim_effect": {"enabled": "money_claim" in types, "source_bo_id": bo_id if "money_claim" in types else None}, + "commercial_successor_effect": {"enabled": "commercial_successor" in types, "source_bo_id": bo_id if "commercial_successor" in types else None}, + "actio_roles": copy.deepcopy(_as_list(signals.get("actio"))), + "linked_actio_structures": { + "structure_ids": sorted(s["structure_id"] for s in structures if s.get("candidate_structure_type") in {"fraudulent_transfer", "preserved_claim"}), + "actio_signal_count": len(_as_list(signals.get("actio"))), + }, + "downstream_module_candidates": modules, + "skeleton_validation_codes": [], + "derived_fact": None, + "is_derived": False, + "derivation_basis": [], + "legal_calculation_object": None, + } + needs_calc = any(item in {"money_claim", "secured_debt", "land_use_gain", "valuation", "fraudulent_transfer"} for item in types) + row["legal_calculation_object"] = { + "calculation_required": needs_calc, + "amount_source": "BO.amount" if row.get("amount") else None, + "amount_text": _text(row.get("amount"), 500) or None, + "object_basis": copy.deepcopy(row.get("object_spec")), + "candidate_structure_types": types, + } + derived_fact = _fallback_derived(bo, row) or None + row["derived_fact"] = derived_fact + row["is_derived"] = bool(derived_fact and not action_text) + row["derivation_basis"] = _derivation_basis( + bo_id, + evidence_refs, + linked["structure_ids"], + row["is_derived"], + ) + reviews: list[dict[str, Any]] = [] + if not row["derived_fact"]: + _add_code(row, "EMPTY_DERIVED_FACT") + _add_code(row, "EMPTY_DERIVED_FACT_AFTER_FALLBACK") + row["must_consider"] = True + reviews.append(_review_template("EMPTY_DERIVED_FACT_AFTER_FALLBACK", row, "NO_SOURCE_BACKED_FALLBACK")) + if not evidence_refs: + _add_code(row, "NO_EVIDENCE_REFS") + if not linked["structure_ids"]: + _add_code(row, "NO_LINKED_LEGAL_EFFECT_STRUCTURE") + if needs_calc and not row.get("amount") and not row.get("object_spec"): + _add_code(row, "CALCULATION_BASIS_MISSING_STAGE2") + reviews.append(_review_template("CALCULATION_BASIS_MISSING_STAGE2", row, "CALCULATION_SOURCE_DEFERRED")) + + upstream = _as_list(_as_dict(source.get("upstream_review_by_bo_id")).get(bo_id)) + for item in upstream: + review = _as_dict(item) + review_type = review.get("review_type") + if review_type == "blocked_exception": + code, severity, owner = "LES_BLOCKED_EXCEPTION_REVIEW", "BLOCK_FINAL_DRAFTING", "HUMAN_REVIEW" + elif review_type == "deterministic_review": + code, severity, owner = "LES_DETERMINISTIC_REVIEW", "WARNING", "STAGE_2" + elif review_type == "review_sensitive_split": + code, severity, owner = "LES_REVIEW_SENSITIVE_SPLIT", "WARNING", "STAGE_2" + else: + code, severity, owner = "LES_REVIEW_REQUIRED", "WARNING", "STAGE_2" + _add_code(row, code) + row["must_consider"] = True + reviews.append(_review_template(code, row, _text(review.get("reason_code") or review_type, 160), severity, owner)) + return row, reviews, evidence_rows, score_basis + + + def _clean_date(value: Any) -> str: + return _text(value, 40) + + + def _periods_overlap(a_start: str, a_end: str, b_start: str, b_end: str) -> bool: + a1, a2 = a_start or "0000-00-00", a_end or "9999-99-99" + b1, b2 = b_start or "0000-00-00", b_end or "9999-99-99" + return max(a1, b1) <= min(a2, b2) + + + def _domain_scans(rows: list[dict[str, Any]], source: dict[str, Any]) -> list[dict[str, Any]]: + reviews: list[dict[str, Any]] = [] + bo_map = _as_dict(source.get("bo_by_id")) + possession_entries: list[tuple[dict[str, Any], set[str], str, str, list[dict[str, Any]]]] = [] + receipt_by_notice: dict[str, str] = {} + for row in rows: + payload = _as_dict(_as_dict(bo_map.get(row["source_bo_id"])).get("domain_payload_compact")) + possession = _as_dict(payload.get("possession_state")) + events = [item for item in _as_list(possession.get("possession_events")) if isinstance(item, dict)] + causes = [item for item in _as_list(possession.get("legal_cause_assertions")) if isinstance(item, dict)] + assets = set(_strings(possession.get("asset_ref")) + _strings(_as_dict(row.get("state_context")).get("asset_cluster_ids"))) + if events and not causes: + _add_code(row, "LEGAL_CAUSE_STATE_UNKNOWN_REVIEW") + reviews.append(_review_template("LEGAL_CAUSE_STATE_UNKNOWN_REVIEW", row, "POSSESSION_WITHOUT_LEGAL_CAUSE")) + cause_names = {_text(item.get("cause_asserted"), 120) for item in causes if _text(item.get("cause_asserted"), 120)} + conflict = any("무권원" in left and any(term in right for term in ("유치권", "임대차", "사용대차")) for left in cause_names for right in cause_names if left != right) + if conflict: + _add_code(row, "POSSESSION_CAUSE_CONFLICT") + reviews.append(_review_template("POSSESSION_CAUSE_CONFLICT", row, "INCOMPATIBLE_LEGAL_CAUSE_ASSERTIONS", "HARD_WARNING", "HUMAN_REVIEW")) + for event in events: + possession_entries.append((row, assets, _clean_date(event.get("start_date")), _clean_date(event.get("end_date")), causes)) + + notice = _as_dict(payload.get("notice_lifecycle")) + notice_id = _text(notice.get("notice_id"), 120) + received = _clean_date(notice.get("received_date")) or _clean_date(notice.get("arrival_date")) + if notice_id and received: + receipt_by_notice[notice_id] = min(received, receipt_by_notice.get(notice_id, "9999-99-99")) + if notice.get("requires_receipt_for_effect") is True and not received: + _add_code(row, "NOTICE_RECEIPT_UNVERIFIED") + reviews.append(_review_template("NOTICE_RECEIPT_UNVERIFIED", row, "RECEIPT_DATE_REQUIRED", "BLOCK_FINAL_DRAFTING", "HUMAN_REVIEW")) + if notice.get("requires_receipt_for_effect") is True and received and not _strings(notice.get("receipt_evidence_refs")) and not _text(notice.get("recipient_admission_text"), 300): + _add_code(row, "NOTICE_RECEIPT_EVIDENCE_WEAK_REVIEW") + reviews.append(_review_template("NOTICE_RECEIPT_EVIDENCE_WEAK_REVIEW", row, "RECEIPT_SUPPORT_WEAK")) + + for row in rows: + payload = _as_dict(_as_dict(bo_map.get(row["source_bo_id"])).get("domain_payload_compact")) + lien = _as_dict(payload.get("lien_state")) + if not lien: + continue + dates = [receipt_by_notice[ref] for ref in _strings(lien.get("extinction_notice_refs")) if ref in receipt_by_notice] + lien_assets = set(_strings(lien.get("asset_ref"))) + if not dates or not lien_assets: + continue + receipt = min(dates) + for possession_row, assets, start, _, _ in possession_entries: + if assets & lien_assets and start and start > receipt: + _add_code(possession_row, "LIEN_ASSERTED_AFTER_EXTINCTION_NOTICE_REVIEW") + reviews.append(_review_template("LIEN_ASSERTED_AFTER_EXTINCTION_NOTICE_REVIEW", possession_row, "POSSESSION_AFTER_NOTICE_RECEIPT", "HARD_WARNING", "HUMAN_REVIEW")) + return reviews + + + def _walk_explicit_candidates(value: Any, path: str = "") -> list[dict[str, Any]]: + out: list[dict[str, Any]] = [] + if isinstance(value, dict): + for key, child in value.items(): + child_path = f"{path}.{key}" if path else str(key) + if key in {"candidate_values", "projection_candidates", "conflict_candidates"} and isinstance(child, list): + for item in child: + if isinstance(item, dict): + candidate = copy.deepcopy(item) + candidate["_path"] = child_path + out.append(candidate) + else: + out.extend(_walk_explicit_candidates(child, child_path)) + elif isinstance(value, list): + for index, child in enumerate(value): + out.extend(_walk_explicit_candidates(child, f"{path}[{index}]")) + return out + + + def _true_exceptions(row: dict[str, Any], source: dict[str, Any]) -> list[dict[str, Any]]: + bo = _as_dict(_as_dict(source.get("bo_by_id")).get(row["source_bo_id"])) + payload = _as_dict(bo.get("domain_payload_compact")) + explicit = _walk_explicit_candidates(payload) + grouped: dict[str, list[dict[str, Any]]] = {} + for item in explicit: + field = _text(item.get("field") or item.get("target_field"), 120) + refs = sorted(set(_strings(item.get("source_refs")) + _strings(item.get("evidence_refs")) + _strings(item.get("structure_refs")))) + if field not in CLOSED_TARGET_FIELDS or not refs or "value" not in item: + continue + grouped.setdefault(field, []).append({ + "field": field, + "value": copy.deepcopy(item.get("value")), + "basis_codes": _sorted_strings(item.get("basis_codes")) or ["EXPLICIT_SOURCE_BACKED_CANDIDATE"], + "source_refs": refs, + }) + out: list[dict[str, Any]] = [] + for field, candidates in sorted(grouped.items()): + unique_values = {_canonical(item["value"]) for item in candidates} + source_sets = {tuple(item["source_refs"]) for item in candidates} + if len(unique_values) < 2 or len(source_sets) < 2: + continue + if field in {"amount", "legal_calculation_object"}: + exception_type = "conflicting_calculation_projection" + elif field in {"object_spec", "state_context"}: + exception_type = "conflicting_asset_projection" + elif field in {"claim_chain_ref", "linked_structures"}: + exception_type = "conflicting_structure_role_projection" + else: + exception_type = "conflicting_fact_projection" + out.append({ + "exception_id": None, + "exception_type": exception_type, + "fact_id": row["fact_id"], + "source_bo_id": row["source_bo_id"], + "target_fields": [field], + "candidate_values": sorted(candidates, key=lambda item: (_canonical(item["value"]), _canonical(item["source_refs"]))), + "candidate_fact_subset": {field: copy.deepcopy(row.get(field))}, + "structure_refs": _strings(_as_dict(row.get("linked_structures")).get("structure_ids")), + "evidence_facts": [], + "allowed_output_fields": [field], + "allowed_decisions": ["KEEP_AS_IS", "PATCH", "BLOCK_REVIEW"], + "conflict_codes": ["EXPLICIT_SOURCE_BACKED_CONFLICT"], + }) + return out + + + def _assign_reviews(items: list[dict[str, Any]]) -> list[dict[str, Any]]: + unique: dict[str, dict[str, Any]] = {} + for item in items: + key = _canonical({key: value for key, value in item.items() if key != "review_id"}) + unique[key] = item + ordered = sorted(unique.values(), key=lambda item: (item["source_bo_id"], item["review_code"], item["reason_code"], item["fact_id"])) + for index, item in enumerate(ordered, start=1): + item["review_id"] = f"FRV-{index:03d}" + return ordered + + + def _pack_exceptions(exceptions: list[dict[str, Any]], reviews: list[dict[str, Any]], rows_by_fact: dict[str, dict[str, Any]], run_fingerprint: str) -> tuple[list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]]]: + eligible: list[dict[str, Any]] = [] + for item in sorted(exceptions, key=lambda value: (value["source_bo_id"], value["exception_type"], value["target_fields"][0])): + if len(_canonical(item).encode("utf-8")) > MAX_EXCEPTION_ITEM_BYTES: + row = rows_by_fact[item["fact_id"]] + _add_code(row, "FL2_REVIEW_REQUIRED") + reviews.append(_review_template("PACK_CONTEXT_TOO_LARGE", row, "EXCEPTION_ITEM_OVER_12000_BYTES", "BLOCK_FINAL_DRAFTING", "HUMAN_REVIEW")) + else: + eligible.append(item) + for index, item in enumerate(eligible, start=1): + item["exception_id"] = f"FX-{index:03d}" + + groups: list[list[dict[str, Any]]] = [] + for exception_type in sorted(CLOSED_EXCEPTION_TYPES): + typed = [item for item in eligible if item["exception_type"] == exception_type] + current: list[dict[str, Any]] = [] + for item in typed: + trial = current + [item] + trial_bytes = len(_canonical({"items": trial}).encode("utf-8")) + if current and (len(trial) > MAX_EXCEPTION_ITEMS_PER_PACK or trial_bytes > MAX_EXCEPTION_PACK_BYTES or len(current) >= TARGET_EXCEPTION_ITEMS_PER_PACK): + groups.append(current) + current = [item] + else: + current = trial + if current: + groups.append(current) + if len(groups) > MAX_EXCEPTION_PACKS: + raise ValueError(f"exception pack count exceeds {MAX_EXCEPTION_PACKS}") + + pack_rows: list[dict[str, Any]] = [] + fanout: list[dict[str, Any]] = [] + for ordinal, items in enumerate(groups, start=1): + file_id = f"PACK-{ordinal:03d}" + pack_id = f"FL-PACK-{ordinal:03d}" + input_path = f"{INPUT_PART_DIR}/{file_id}.json" + output_path = f"{OUTPUT_PART_DIR}/{file_id}.json" + exception_ids = [item["exception_id"] for item in items] + exception_type = items[0]["exception_type"] + part = { + "schema_version": "stage1_fact_exception_input_part.v2", + "status": "READY", + "pack_id": pack_id, + "run_fingerprint": run_fingerprint, + "exception_type": exception_type, + "input_exception_ids": exception_ids, + "items": items, + "output_path": output_path, + } + serialized_bytes = len(_canonical(part).encode("utf-8")) + if serialized_bytes > MAX_EXCEPTION_PACK_BYTES: + raise ValueError(f"serialized exception pack exceeds budget: {pack_id}") + _write_json(input_path, part) + _verify_reread(input_path, part) + pack_rows.append({ + "pack_id": pack_id, + "exception_type": exception_type, + "exception_ids": exception_ids, + "input_path": input_path, + "output_path": output_path, + "serialized_bytes": serialized_bytes, + }) + fanout.append({ + "pack_id": pack_id, + "exception_type": exception_type, + "exception_ids": exception_ids, + "items": items, + "output_path": output_path, + "run_fingerprint": run_fingerprint, + }) + return eligible, pack_rows, fanout + + + def main() -> None: + _init_mcp_session() + source_outer = _read_json(SOURCE_PACK_PATH) + source = _as_dict(_as_dict(source_outer).get("fact_source_pack")) + if source.get("schema_version") != "stage1_fact_source_pack.v3": + raise ValueError("fact_source_pack schema_version mismatch") + if source.get("status") not in {"READY", "READY_WITH_REVIEW"}: + raise ValueError("fact_source_pack status invalid") + run_fingerprint = _text(source.get("run_fingerprint"), 100) + if not run_fingerprint.startswith("sha256:"): + raise ValueError("fact_source_pack run_fingerprint invalid") + bo_ids = _sorted_strings(_as_dict(source.get("reference_universe")).get("bo_ids")) + if not bo_ids or bo_ids != sorted(_as_dict(source.get("bo_by_id")).keys()): + raise ValueError("source BO universe mismatch") + _write_json(ROW_SCHEMA_PATH, FACT_LEDGER_ROW_SCHEMA) + _verify_reread(ROW_SCHEMA_PATH, FACT_LEDGER_ROW_SCHEMA) + row_schema_sha256 = _digest(FACT_LEDGER_ROW_SCHEMA) + + rows: list[dict[str, Any]] = [] + reviews: list[dict[str, Any]] = [] + evidence_score_basis: dict[str, str] = {} + exceptions: list[dict[str, Any]] = [] + for index, bo_id in enumerate(bo_ids, start=1): + row, row_reviews, _, score_basis = _build_row(index, bo_id, source) + _validate_row_schema(row, f"candidate_items[{index - 1}]") + rows.append(row) + reviews.extend(row_reviews) + evidence_score_basis[row["fact_id"]] = score_basis + exceptions.extend(_true_exceptions(row, source)) + reviews.extend(_domain_scans(rows, source)) + rows_by_fact = {row["fact_id"]: row for row in rows} + exceptions, pack_rows, fanout = _pack_exceptions(exceptions, reviews, rows_by_fact, run_fingerprint) + reviews = _assign_reviews(reviews) + for index, row in enumerate(rows): + _validate_row_schema(row, f"candidate_items[{index}]") + + exception_ids = [item["exception_id"] for item in exceptions] + if sorted(item for pack in pack_rows for item in pack["exception_ids"]) != sorted(exception_ids): + raise ValueError("exception pack conservation failed") + manifest_status = "READY" if exceptions else ("READY_WITH_REVIEW" if reviews else "READY_NO_EXCEPTIONS") + manifest = { + "fact_exception_manifest": { + "schema_version": "stage1_fact_exception_manifest.v2", + "status": manifest_status, + "run_fingerprint": run_fingerprint, + "expected_exception_ids": exception_ids, + "exception_count": len(exceptions), + "deterministic_review_items": reviews, + "pack_count": len(pack_rows), + "packs": pack_rows, + } + } + _write_json(EXCEPTION_MANIFEST_PATH, manifest) + _verify_reread(EXCEPTION_MANIFEST_PATH, manifest) + + candidate_bo_ids = [row["source_bo_id"] for row in rows] + missing = sorted(set(bo_ids) - set(candidate_bo_ids)) + unknown = sorted(set(candidate_bo_ids) - set(bo_ids)) + duplicates = sorted({item for item in candidate_bo_ids if candidate_bo_ids.count(item) > 1}) + if missing or unknown or duplicates or len(rows) != len(bo_ids): + raise ValueError("candidate BO conservation failed") + candidate_status = "READY_WITH_EXCEPTIONS" if exceptions else ("READY_WITH_REVIEW" if reviews else "READY") + bundle = { + "fact_ledger_candidate_bundle": { + "schema_version": "stage1_fact_ledger_candidate_bundle.v3", + "status": candidate_status, + "run_fingerprint": run_fingerprint, + "source_pack_ref": { + "path": SOURCE_PACK_PATH, + "schema_version": "stage1_fact_source_pack.v3", + "sha256": _digest(source_outer), + }, + "row_schema_ref": { + "path": ROW_SCHEMA_PATH, + "$id": FACT_LEDGER_ROW_SCHEMA_ID, + "sha256": row_schema_sha256, + }, + "candidate_items": rows, + "exception_manifest_ref": { + "path": EXCEPTION_MANIFEST_PATH, + "schema_version": "stage1_fact_exception_manifest.v2", + "sha256": _digest(manifest), + "expected_exception_ids": exception_ids, + }, + "compile_gate": { + "source_bo_count": len(bo_ids), + "candidate_count": len(rows), + "source_bo_ids": bo_ids, + "candidate_bo_ids": candidate_bo_ids, + "missing_bo_ids": missing, + "unknown_bo_ids": unknown, + "duplicate_bo_ids": duplicates, + "bo_conservation_pass": True, + "candidate_schema_closed": True, + "candidate_schema_validation_pass": True, + "row_schema_id": FACT_LEDGER_ROW_SCHEMA_ID, + "row_schema_sha256": row_schema_sha256, + "reference_validation_pass": True, + "evidence_score_basis": evidence_score_basis, + }, + } + } + _write_json(CANDIDATE_PATH, bundle) + _verify_reread(CANDIDATE_PATH, bundle) + stdout_status = "READY" if exceptions else ("READY_WITH_REVIEW" if reviews else "READY_NO_EXCEPTIONS") + print(_canonical({ + "fact_ledger_fl1_manifest": { + "schema_version": "stage1_fact_ledger_fl1_manifest.v3", + "status": stdout_status, + "candidate_path": CANDIDATE_PATH, + "row_schema_path": ROW_SCHEMA_PATH, + "row_schema_sha256": row_schema_sha256, + "exception_manifest_path": EXCEPTION_MANIFEST_PATH, + "run_fingerprint": run_fingerprint, + "candidate_count": len(rows), + "exception_count": len(exceptions), + "review_count": len(reviews), + "pack_count": len(pack_rows), + }, + "dynamic_fanout": fanout, + })) + + + try: + main() + except Exception as exc: + print(_canonical({"fact_ledger_fl1_manifest": { + "schema_version": "stage1_fact_ledger_fl1_manifest.v3", + "status": "FAILED", + "error": _safe_error(exc), + }, "dynamic_fanout": []})) + raise + + - task_name: Task_FL2_fact_exception_adjudicator_* + llm_provider: google + llm_model: gemini-3.1-flash-lite + llm_reasoning: medium + llm_verbosity: low + max_iterations: 2 + max_concurrency: 2 + use_tools: + - localdocs + cache_control: + mode: auto + ttl: 1h + preflight: false + prompts: + - role: user + content: | + + + - You are executing one LLM sub-task inside Stage 1 of a Korean civil-litigation complaint-generation pipeline. + - The current task's static block, role overlay, assigned inputs, output schema, and writer boundary control. + - This common prefix cannot expand the current task's input set, output set, legal domain, validation authority, writer authority, or reasoning depth. + - If any common rule appears broader than the current task, apply only the narrower current-task version. + - Stage 1 prepares verified structured artifacts. Do not draft complaint prose or final counsel-level conclusions unless the current task explicitly authorizes a validation or gate conclusion. + + + + - Use only assigned files, provided context inputs, prior outputs, and allowed tools. + - Do not import facts, law, procedural history, parties, dates, amounts, IDs, document contents, or source meanings from memory, outside knowledge, or unassigned files. + - Treat prior outputs as authority only to the extent the current task names them or provides them as context. + - If a value is unsupported, missing, conflicting, stale, or out of scope, use only the current schema's allowed null, empty, unknown, warning, blocked, or needs_review path. + + + + - Preserve exact source identifiers required by the current schema. + - Maintain separation among raw fact, inferred fact, legal signal, evidence support, fact support, validation issue, and final gate decision when the current schema distinguishes them. + - Do not upgrade meeting-only or indirect material into direct proof. + - Do not silently resolve material conflicts. If the current schema has a conflict or uncertainty field, use it; otherwise stay within the task's allowed warning or review path. + + + + - Follow required JSON shape, key names, enum values, ordering, file names, and status strings exactly. + - Do not add arbitrary keys, prose, markdown fences, alternative files, unauthorized repair, or explanatory material outside allowed fields. + - Create, mutate, normalize, merge, or finalize IDs only when the current task explicitly authorizes it. + - Write final files only when the current task is the authorized writer. Validators and guards report issues in their own authorized schema and do not silently repair unless instructed. + + + + - Prefer the current prompt and schema, assigned structured upstream artifacts, compact indexes, ledgers, manifests, bundles, and gates. + - Read raw evidence or meeting text only when the current task requires direct provenance, ambiguity resolution, or a schema-required value missing from structured artifacts. + - For map or projection tasks, process only the assigned item, domain, or batch. Reducers aggregate only the inputs assigned to them. + - Do not restate, summarize, cite, or copy this common prefix in any output. + + + + - Return only the requested structured artifact, concise allowed rationale fields, validation notes, or status object. + - Keep chain-of-thought private. + - Stop when the current schema is complete and safe. + + + + + TASK_NAME: Task_FL2_fact_exception_adjudicator_* + STAGE: Stage 1 Fact Ledger generation + ROLE: compact source-backed Fact projection exception adjudicator + + + + pack_id: {{item.pack_id}} + exception_type: {{item.exception_type}} + exception_ids: {{item.exception_ids}} + items: {{item.items}} + output_path: {{item.output_path}} + run_fingerprint: {{item.run_fingerprint}} + + + + - Act as a senior Korean civil-litigation lawyer and AI systems architect, but decide only the assigned compact projection conflicts. + - Use only ASSIGNED_PACK. Do not read or request any file, any other pack, BO, evidence, signal, structure, source pack, candidate bundle, or final ledger. + - Process every assigned exception ID exactly once, without omission, duplication, merge, split, or reorder. + - Select only KEEP_AS_IS, PATCH, or BLOCK_REVIEW when that decision appears in the item's allowed_decisions. + - PATCH keys must be a subset of allowed_output_fields. Every PATCH value must exactly equal one candidate_values[].value for the same field. + - source_refs_used must be a subset of source_refs embedded in that exception item. + - Never create a fact, party, date, amount, asset, BO ID, Fact ID, exception ID, structure ID, evidence ID, candidate value, source ref, or legal theory. + - If compact context is insufficient, use BLOCK_REVIEW. Never guess. + - Do not write Fact_Ledger_base.json or any input, source, candidate, manifest, or other adjudication file. + - The only permitted tool action is localdocs write_file with overwrite=true to ASSIGNED_PACK.output_path. Do not call read_docs, list_docs, or any other tool. + - Write exactly one raw JSON adjudication part, then return only the compact manifest JSON. Do not return markdown, code fences, explanation, or the file body. + + + + - KEEP_AS_IS: preserve the deterministic candidate when its source priority remains supportable; patch must be {}, review_required=false. + - PATCH: choose exactly one source-backed candidate for each patched field; review_required=false. + - BLOCK_REVIEW: use when sources remain materially balanced, compact context is insufficient, or a safe selection requires legal theory; patch must be {}, review_required=true. + - No split decision is allowed. + + + + Each decisions[] item must contain exactly: + { + "exception_id": "FX-001", + "exception_type": "conflicting_fact_projection|conflicting_calculation_projection|conflicting_asset_projection|conflicting_structure_role_projection", + "fact_id": "F-001", + "source_bo_id": "bh1", + "decision": "KEEP_AS_IS|PATCH|BLOCK_REVIEW", + "patch": {}, + "review_required": false, + "reason_code": "KEEP_EXISTING_SOURCE_PRIORITY|SELECT_SOURCE_BACKED_CANDIDATE|INSUFFICIENT_COMPACT_CONTEXT|UNSUPPORTED_PATCH", + "basis_codes_used": [], + "source_refs_used": [], + "conflict_codes_acknowledged": [] + } + + + + Write this raw JSON object to ASSIGNED_PACK.output_path: + { + "schema_version": "stage1_fact_exception_adjudication_part.v2", + "status": "READY|READY_WITH_BLOCKS|FAILED", + "pack_id": "the assigned pack_id", + "run_fingerprint": "the assigned run_fingerprint", + "exception_type": "the assigned exception_type", + "input_exception_ids": ["all assigned exception IDs exactly once"], + "decisions": ["one decision per assigned exception"], + "coverage": { + "expected": 0, + "actual": 0, + "missing_exception_ids": [], + "unknown_exception_ids": [], + "duplicate_exception_ids": [], + "coverage_pass": true + } + } + Use READY_WITH_BLOCKS when any decision is BLOCK_REVIEW; otherwise use READY. expected and actual must equal the assigned exception count. + + + + After localdocs confirms the write, return exactly: + { + "pack_id": "the assigned pack_id", + "status": "READY|READY_WITH_BLOCKS|FAILED", + "output_path": "the assigned output_path", + "decision_count": 0 + } + + + + - task_name: Task_FL3_final_fact_ledger_gate_and_writer + mcp: code-executor + tool_name: run_code + parameters: + language: python + requirements: |- + httpx + jsonschema>=4.23.0 + network: agent-network + timeout: 240 + code: | + #!/usr/bin/env python3 + from __future__ import annotations + + import copy + import hashlib + import itertools + import json + import re + from typing import Any + + import httpx + from jsonschema import Draft202012Validator + from jsonschema.exceptions import ValidationError + + LOCALDOCS_URL = "http://mcp-localdocs:8012/mcp" + MCP_HEADERS = { + "Content-Type": "application/json", + "Accept": "application/json, text/event-stream", + } + CLIENT = httpx.Client(timeout=60) + MSG_ID_COUNTER = itertools.count(10) + JSON_DECODER = json.JSONDecoder() + + + def _next_msg_id() -> int: + return next(MSG_ID_COUNTER) + + + def _safe_error(exc: BaseException) -> str: + return re.sub(r"\s+", " ", str(exc)).strip()[:300] + + + def _parse_json_value(raw: Any, label: str, allow_trailing: bool = False) -> Any: + if isinstance(raw, (dict, list)): + return raw + if not isinstance(raw, str): + raise ValueError(f"{label} is not JSON text") + text = raw.lstrip("\ufeff").strip() + if not text: + raise ValueError(f"{label} is empty") + try: + return json.loads(text) + except json.JSONDecodeError: + try: + value, end = JSON_DECODER.raw_decode(text) + except json.JSONDecodeError as exc: + raise ValueError(f"{label} invalid JSON: {text[:300]}") from exc + if not allow_trailing and text[end:].strip(): + raise ValueError(f"{label} has trailing content: {text[end:end + 300]}") + return value + + + def _parse_mcp_response(text: str) -> dict[str, Any]: + parsed: dict[str, Any] | None = None + for line in text.strip().splitlines(): + if not line.startswith("data: "): + continue + try: + item = _parse_json_value(line[6:], "MCP SSE data", True) + except ValueError: + continue + if isinstance(item, dict): + parsed = item + if parsed is None: + item = _parse_json_value(text, "MCP response", True) + if not isinstance(item, dict): + raise RuntimeError("MCP response root is not an object") + parsed = item + return parsed + + + def _init_mcp_session() -> None: + response = CLIENT.post( + LOCALDOCS_URL, + json={ + "jsonrpc": "2.0", + "id": 1, + "method": "initialize", + "params": { + "protocolVersion": "2025-03-26", + "capabilities": {}, + "clientInfo": { + "name": "stage-1-fact-ledger-fl3-final-writer", + "version": "2.0", + "user_id": "{{__user_hash__}}", + "workspace_id": "{{__workspace_hash__}}", + }, + }, + }, + headers=MCP_HEADERS, + ) + response.raise_for_status() + payload = _parse_mcp_response(response.text) + if payload.get("error"): + raise RuntimeError(f"MCP initialize failed: {payload['error']}") + session_id = response.headers.get("mcp-session-id") + if session_id: + MCP_HEADERS["mcp-session-id"] = session_id + initialized = CLIENT.post( + LOCALDOCS_URL, + json={"jsonrpc": "2.0", "method": "notifications/initialized"}, + headers=MCP_HEADERS, + ) + initialized.raise_for_status() + + + def _call_tool(name: str, arguments: dict[str, Any]) -> dict[str, Any]: + response = CLIENT.post( + LOCALDOCS_URL, + json={ + "jsonrpc": "2.0", + "id": _next_msg_id(), + "method": "tools/call", + "params": {"name": name, "arguments": arguments}, + }, + headers=MCP_HEADERS, + ) + response.raise_for_status() + payload = _parse_mcp_response(response.text) + if payload.get("error"): + raise RuntimeError(f"MCP {name} failed: {payload['error']}") + result = payload.get("result") + if not isinstance(result, dict) or result.get("isError") is True: + raise RuntimeError(f"MCP {name} returned an invalid result") + return result + + + def _extract_doc(result: dict[str, Any], doc_name: str) -> Any: + blocks = result.get("content") or [] + if not blocks or not isinstance(blocks[0], dict): + raise RuntimeError(f"Empty response: {doc_name}") + text = blocks[0].get("text") + if not isinstance(text, str) or not text.strip(): + raise RuntimeError(f"Empty text response: {doc_name}") + outer = _parse_json_value(text, f"{doc_name} outer envelope", True) + inner: Any = outer + if isinstance(outer, dict) and "results" in outer: + results = outer.get("results") or [] + if not results or not isinstance(results[0], dict): + raise RuntimeError(f"Empty results envelope: {doc_name}") + inner = results[0].get("content") + if inner in (None, ""): + inner = results[0].get("text") + if isinstance(inner, (dict, list)): + return inner + return _parse_json_value(inner, f"{doc_name} inner document", False) + + + def _read_json(path: str) -> Any: + return _extract_doc(_call_tool("read_docs", {"doc_names": [path]}), path) + + + def _write_text(path: str, content: str) -> None: + _call_tool("write_file", {"path": path, "content": content, "overwrite": True}) + + + def _canonical(value: Any) -> str: + return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")) + + + def _digest(value: Any) -> str: + return "sha256:" + hashlib.sha256(_canonical(value).encode("utf-8")).hexdigest() + + + def _write_json(path: str, value: Any) -> None: + _write_text(path, json.dumps(value, ensure_ascii=False, sort_keys=True, indent=2)) + + + def _verify_reread(path: str, expected: Any) -> None: + actual = _read_json(path) + if _digest(actual) != _digest(expected): + raise RuntimeError(f"Post-write verification failed: {path}") + + + def _as_dict(value: Any) -> dict[str, Any]: + return value if isinstance(value, dict) else {} + + + def _as_list(value: Any) -> list[Any]: + return value if isinstance(value, list) else [] + + + def _strings(value: Any) -> list[str]: + values = value if isinstance(value, list) else [value] + out: list[str] = [] + for item in values: + if item is None: + continue + if isinstance(item, str): + text = item.strip() + elif isinstance(item, (int, float)) and not isinstance(item, bool): + text = str(item) + else: + continue + if text and text not in out: + out.append(text) + return out + + + def _sorted_strings(value: Any) -> list[str]: + return sorted(set(_strings(value))) + + + def _text(value: Any, limit: int = 500) -> str: + if value in (None, "", [], {}): + return "" + raw = value if isinstance(value, str) else json.dumps(value, ensure_ascii=False, sort_keys=True) + return re.sub(r"\s+", " ", raw).strip()[:limit] + + SOURCE_PACK_PATH = "stage1_tmp/fact_ledger/fact_source_pack.json" + CANDIDATE_PATH = "stage1_tmp/fact_ledger/Fact_Ledger_base_candidate.json" + ROW_SCHEMA_PATH = "stage1_tmp/fact_ledger/fact_ledger_row_schema.json" + EXCEPTION_MANIFEST_PATH = "stage1_tmp/fact_ledger/fact_exception_manifest.json" + TARGET_PATH = "Fact_Ledger_base.json" + REPORT_PATH = "stage1_tmp/fact_ledger/fact_ledger_writer_report.json" + FINAL_FIELDS = [ + "fact_id", "source_bo_id", "type", "date", "parties", "object_spec", "amount", "action", + "evidence_refs", "credibility", "legal_centrality", "proof_strength", "legal_effect_roles", + "claim_chain_ref", "linked_structures", "must_not_drop_in_claim_types", "must_consider", + "state_context", "money_claim_effect", "commercial_successor_effect", "actio_roles", + "linked_actio_structures", "downstream_module_candidates", "skeleton_validation_codes", + "derived_fact", "is_derived", "derivation_basis", "legal_calculation_object", + ] + FINAL_FIELD_SET = set(FINAL_FIELDS) + FACT_LEDGER_ROW_SCHEMA_ID = "https://eroomai.com/schemas/stage1/fact-ledger-row.v4.json" + CLOSED_EXCEPTION_TYPES = { + "conflicting_fact_projection", + "conflicting_calculation_projection", + "conflicting_asset_projection", + "conflicting_structure_role_projection", + } + CLOSED_DECISIONS = {"KEEP_AS_IS", "PATCH", "BLOCK_REVIEW"} + + + def _add_code(row: dict[str, Any], code: str) -> None: + row["skeleton_validation_codes"] = sorted(set(_strings(row.get("skeleton_validation_codes")) + [code])) + + + def _looks_like_json_container_string(value: Any) -> bool: + if not isinstance(value, str) or not value.strip().startswith(("{", "[")): + return False + try: + return isinstance(json.loads(value), (dict, list)) + except json.JSONDecodeError: + return False + + + def _validate_row_schema(row: dict[str, Any], validator: Draft202012Validator, label: str) -> None: + for field in ("amount", "object_spec"): + if _looks_like_json_container_string(row.get(field)): + raise ValueError(f"{label}.{field} contains a JSON-encoded String") + errors: list[ValidationError] = sorted( + validator.iter_errors(row), + key=lambda item: ([str(part) for part in item.absolute_path], item.message), + ) + if errors: + error = errors[0] + path = ".".join(str(part) for part in error.absolute_path) or "$" + raise ValueError(f"{label} JSON Schema violation at {path}: {error.message}") + + + def _load_row_schema(bundle: dict[str, Any]) -> tuple[dict[str, Any], Draft202012Validator, str]: + schema_ref = _as_dict(bundle.get("row_schema_ref")) + if schema_ref.get("path") != ROW_SCHEMA_PATH or schema_ref.get("$id") != FACT_LEDGER_ROW_SCHEMA_ID: + raise ValueError("row schema ref path/$id mismatch") + schema = _as_dict(_read_json(ROW_SCHEMA_PATH)) + if schema.get("$schema") != "https://json-schema.org/draft/2020-12/schema" or schema.get("$id") != FACT_LEDGER_ROW_SCHEMA_ID: + raise ValueError("row schema identity mismatch") + if schema.get("required") != FINAL_FIELDS or schema.get("additionalProperties") is not False: + raise ValueError("row schema closed-field contract mismatch") + schema_sha256 = _digest(schema) + if schema_ref.get("sha256") != schema_sha256: + raise ValueError("row schema hash mismatch") + Draft202012Validator.check_schema(schema) + return schema, Draft202012Validator(schema), schema_sha256 + + + def _validate_rows(bundle: dict[str, Any], validator: Draft202012Validator, source_bo_ids: set[str], label: str) -> list[dict[str, Any]]: + if bundle.get("schema_version") != "stage1_fact_ledger_candidate_bundle.v3": + raise ValueError("candidate bundle schema_version mismatch") + if bundle.get("status") not in {"READY", "READY_WITH_REVIEW", "READY_WITH_EXCEPTIONS"}: + raise ValueError("candidate bundle status invalid") + rows = bundle.get("candidate_items") + if not isinstance(rows, list) or not rows: + raise ValueError("candidate_items missing") + for index, row in enumerate(rows): + if not isinstance(row, dict): + raise ValueError(f"{label}[{index}] is not an object") + _validate_row_schema(row, validator, f"{label}[{index}]") + if row.get("source_bo_id") not in source_bo_ids: + raise ValueError(f"{label}[{index}].source_bo_id is outside upstream BO universe") + return copy.deepcopy(rows) + + + def _load_exception_inputs(manifest: dict[str, Any], run_fingerprint: str) -> tuple[dict[str, dict[str, Any]], list[str]]: + items: dict[str, dict[str, Any]] = {} + input_ids: list[str] = [] + packs = _as_list(manifest.get("packs")) + if len(packs) != manifest.get("pack_count"): + raise ValueError("exception manifest pack_count mismatch") + for pack in packs: + row = _as_dict(pack) + part = _as_dict(_read_json(_text(row.get("input_path"), 300))) + if part.get("schema_version") != "stage1_fact_exception_input_part.v2" or part.get("status") != "READY": + raise ValueError("exception input part schema/status invalid") + if part.get("pack_id") != row.get("pack_id") or part.get("run_fingerprint") != run_fingerprint: + raise ValueError("exception input part identity mismatch") + if part.get("output_path") != row.get("output_path") or part.get("exception_type") != row.get("exception_type"): + raise ValueError("exception input part routing mismatch") + part_ids = _strings(part.get("input_exception_ids")) + if part_ids != _strings(row.get("exception_ids")): + raise ValueError("exception input ID mismatch") + for item in _as_list(part.get("items")): + if not isinstance(item, dict): + raise ValueError("exception input item invalid") + exception_id = _text(item.get("exception_id"), 120) + if not exception_id or exception_id in items: + raise ValueError(f"duplicate exception input: {exception_id}") + if item.get("exception_type") not in CLOSED_EXCEPTION_TYPES: + raise ValueError("exception input type invalid") + items[exception_id] = item + input_ids.append(exception_id) + return items, input_ids + + + def _load_decisions(manifest: dict[str, Any], run_fingerprint: str) -> tuple[dict[str, dict[str, Any]], list[str]]: + decisions: dict[str, dict[str, Any]] = {} + decision_ids: list[str] = [] + for pack in _as_list(manifest.get("packs")): + expected_pack = _as_dict(pack) + part = _as_dict(_read_json(_text(expected_pack.get("output_path"), 300))) + if part.get("schema_version") != "stage1_fact_exception_adjudication_part.v2": + raise ValueError("adjudication part schema_version mismatch") + if part.get("status") not in {"READY", "READY_WITH_BLOCKS"}: + raise ValueError("adjudication part failed or has invalid status") + if part.get("pack_id") != expected_pack.get("pack_id") or part.get("run_fingerprint") != run_fingerprint: + raise ValueError("adjudication part identity mismatch") + if part.get("exception_type") != expected_pack.get("exception_type"): + raise ValueError("adjudication part type mismatch") + if _strings(part.get("input_exception_ids")) != _strings(expected_pack.get("exception_ids")): + raise ValueError("adjudication input IDs mismatch") + coverage = _as_dict(part.get("coverage")) + if coverage.get("coverage_pass") is not True or _as_list(coverage.get("missing_exception_ids")) or _as_list(coverage.get("unknown_exception_ids")) or _as_list(coverage.get("duplicate_exception_ids")): + raise ValueError("adjudication part coverage failed") + for decision in _as_list(part.get("decisions")): + if not isinstance(decision, dict): + raise ValueError("adjudication decision invalid") + exception_id = _text(decision.get("exception_id"), 120) + if not exception_id or exception_id in decisions: + raise ValueError(f"duplicate adjudication decision: {exception_id}") + decisions[exception_id] = decision + decision_ids.append(exception_id) + return decisions, decision_ids + + + def _allowed_values(item: dict[str, Any], field: str) -> list[Any]: + return [candidate.get("value") for candidate in _as_list(item.get("candidate_values")) if isinstance(candidate, dict) and candidate.get("field") == field] + + + def _all_source_refs(item: dict[str, Any]) -> set[str]: + refs: set[str] = set() + for candidate in _as_list(item.get("candidate_values")): + if isinstance(candidate, dict): + refs.update(_strings(candidate.get("source_refs"))) + refs.update(_strings(item.get("structure_refs"))) + return refs + + + def _apply_decision(row: dict[str, Any], item: dict[str, Any], decision: dict[str, Any]) -> bool: + for key in ("exception_id", "exception_type", "fact_id", "source_bo_id"): + if decision.get(key) != item.get(key): + raise ValueError(f"decision identity mismatch: {key}") + decision_type = decision.get("decision") + if decision_type not in CLOSED_DECISIONS or decision_type not in _as_list(item.get("allowed_decisions")): + raise ValueError("decision enum invalid") + patch = decision.get("patch") + if not isinstance(patch, dict): + raise ValueError("decision patch must be an object") + used_refs = set(_strings(decision.get("source_refs_used"))) + if not used_refs.issubset(_all_source_refs(item)): + raise ValueError("decision used unknown source refs") + if decision_type == "KEEP_AS_IS": + if patch or decision.get("review_required") is True: + raise ValueError("KEEP_AS_IS contract violated") + return False + if decision_type == "BLOCK_REVIEW": + if patch or decision.get("review_required") is not True: + raise ValueError("BLOCK_REVIEW contract violated") + _add_code(row, "BLOCK_REVIEW_PRESENT") + _add_code(row, "FL2_REVIEW_REQUIRED") + row["must_consider"] = True + basis = row.get("derivation_basis") + if not isinstance(basis, list): + raise ValueError("derivation_basis must remain an Array of Object") + basis.append({ + "source_type": "review", + "source_id": str(item.get("exception_id")), + "source_field": "exception_adjudication", + "role": "review", + }) + return True + if decision.get("review_required") is True or not patch: + raise ValueError("PATCH contract violated") + allowed_fields = set(_strings(item.get("allowed_output_fields"))) + if not set(patch).issubset(allowed_fields) or not set(patch).issubset(FINAL_FIELD_SET): + raise ValueError("PATCH field outside hard whitelist") + for field, value in patch.items(): + if not any(_canonical(value) == _canonical(candidate) for candidate in _allowed_values(item, field)): + raise ValueError("PATCH value is not source-backed") + row[field] = copy.deepcopy(value) + return False + + + def _reference_universes(source: dict[str, Any]) -> dict[str, set[str]]: + structures = _as_dict(source.get("structure_by_id")) + claims: set[str] = set() + liabilities: set[str] = set() + actio_signals: set[str] = set() + for structure in structures.values(): + row = _as_dict(structure) + claims.update(_strings(row.get("linked_claim_group_ids"))) + liabilities.update(_strings(row.get("linked_liability_group_ids"))) + actio_signals.update(_strings(row.get("linked_actio_signal_ids"))) + for by_type in _as_dict(source.get("signal_hints_by_bo_id")).values(): + for signal_type, rows in _as_dict(by_type).items(): + for signal in _as_list(rows): + item = _as_dict(signal) + claims.update(_strings(item.get("claim_group_candidates"))) + liabilities.update(_strings(item.get("liability_group_candidates"))) + if signal_type == "actio": + actio_signals.update(_strings(item.get("signal_id"))) + return {"claim": claims, "liability": liabilities, "actio": actio_signals} + + + def _validate_row_refs(row: dict[str, Any], source: dict[str, Any], group_universes: dict[str, set[str]], report: dict[str, Any]) -> None: + universe = _as_dict(source.get("reference_universe")) + bo_ids = set(_strings(universe.get("bo_ids"))) + structure_ids = set(_strings(universe.get("structure_ids"))) + evidence_ids = set(_strings(universe.get("evidence_indexes"))) + bo_id = _text(row.get("source_bo_id"), 120) + unknown_bo = [] if bo_id in bo_ids else [bo_id] + unknown_structures = sorted(set(_strings(_as_dict(row.get("linked_structures")).get("structure_ids"))) - structure_ids) + unknown_evidence = sorted(set(_strings(row.get("evidence_refs"))) - evidence_ids) + report["reference_gate"]["unknown_bo_refs"].extend(unknown_bo) + report["reference_gate"]["unknown_structure_refs"].extend(unknown_structures) + report["reference_gate"]["unknown_evidence_refs"].extend(unknown_evidence) + claim = _as_dict(row.get("claim_chain_ref")) + if not set(_strings(claim.get("claim_group_ids"))).issubset(group_universes["claim"]): + raise ValueError(f"unknown claim group ref: {bo_id}") + if not set(_strings(claim.get("liability_group_ids"))).issubset(group_universes["liability"]): + raise ValueError(f"unknown liability group ref: {bo_id}") + + + def main() -> None: + _init_mcp_session() + source_outer = _read_json(SOURCE_PACK_PATH) + candidate_outer = _read_json(CANDIDATE_PATH) + manifest_outer = _read_json(EXCEPTION_MANIFEST_PATH) + source = _as_dict(_as_dict(source_outer).get("fact_source_pack")) + bundle = _as_dict(_as_dict(candidate_outer).get("fact_ledger_candidate_bundle")) + manifest = _as_dict(_as_dict(manifest_outer).get("fact_exception_manifest")) + if source.get("schema_version") != "stage1_fact_source_pack.v3" or source.get("status") not in {"READY", "READY_WITH_REVIEW"}: + raise ValueError("fact_source_pack invalid") + if manifest.get("schema_version") != "stage1_fact_exception_manifest.v2" or manifest.get("status") not in {"READY", "READY_NO_EXCEPTIONS", "READY_WITH_REVIEW"}: + raise ValueError("fact_exception_manifest invalid") + source_bo_ids = _strings(_as_dict(source.get("reference_universe")).get("bo_ids")) + if not source_bo_ids: + raise ValueError("upstream BO universe missing") + _, row_validator, row_schema_sha256 = _load_row_schema(bundle) + rows = _validate_rows(bundle, row_validator, set(source_bo_ids), "candidate_items") + run_fingerprint = _text(source.get("run_fingerprint"), 100) + if not run_fingerprint or bundle.get("run_fingerprint") != run_fingerprint or manifest.get("run_fingerprint") != run_fingerprint: + raise ValueError("run_fingerprint mismatch") + source_ref = _as_dict(bundle.get("source_pack_ref")) + manifest_ref = _as_dict(bundle.get("exception_manifest_ref")) + if source_ref.get("path") != SOURCE_PACK_PATH or source_ref.get("schema_version") != "stage1_fact_source_pack.v3" or source_ref.get("sha256") != _digest(source_outer): + raise ValueError("source pack hash/ref mismatch") + if manifest_ref.get("path") != EXCEPTION_MANIFEST_PATH or manifest_ref.get("schema_version") != "stage1_fact_exception_manifest.v2" or manifest_ref.get("sha256") != _digest(manifest_outer): + raise ValueError("exception manifest hash/ref mismatch") + + expected_ids = _strings(manifest.get("expected_exception_ids")) + if expected_ids != _strings(manifest_ref.get("expected_exception_ids")) or len(expected_ids) != manifest.get("exception_count"): + raise ValueError("expected exception IDs mismatch") + if not expected_ids and (_as_list(manifest.get("packs")) or manifest.get("pack_count") != 0): + raise ValueError("no-exception manifest contains packs") + items_by_id, input_ids = _load_exception_inputs(manifest, run_fingerprint) + decisions_by_id, decision_ids = _load_decisions(manifest, run_fingerprint) + missing_ids = sorted(set(expected_ids) - set(decision_ids)) + unknown_ids = sorted(set(decision_ids) - set(expected_ids)) + duplicate_ids = sorted({item for item in decision_ids if decision_ids.count(item) > 1}) + if sorted(expected_ids) != sorted(input_ids) or missing_ids or unknown_ids or duplicate_ids or len(decision_ids) != len(expected_ids): + raise ValueError("exact exception coverage failed") + + rows_by_fact = {row["fact_id"]: row for row in rows} + blocked_items: list[dict[str, Any]] = [] + for exception_id in expected_ids: + item = items_by_id.get(exception_id) + decision = decisions_by_id.get(exception_id) + if not item or not decision: + raise ValueError(f"exception coverage missing: {exception_id}") + row = rows_by_fact.get(_text(item.get("fact_id"), 120)) + if not row or row.get("source_bo_id") != item.get("source_bo_id"): + raise ValueError("exception Fact/BO target mismatch") + blocked = _apply_decision(row, item, decision) + _validate_row_schema(row, row_validator, f"post_patch[{exception_id}]") + if blocked: + blocked_items.append({ + "exception_id": exception_id, + "exception_type": item.get("exception_type"), + "fact_id": item.get("fact_id"), + "source_bo_id": item.get("source_bo_id"), + "reason_code": decision.get("reason_code"), + }) + for index, row in enumerate(rows): + _validate_row_schema(row, row_validator, f"post_patch_items[{index}]") + + compile_gate = _as_dict(bundle.get("compile_gate")) + candidate_bo_ids = [str(row.get("source_bo_id")) for row in rows] + duplicate_bos = sorted({item for item in candidate_bo_ids if candidate_bo_ids.count(item) > 1}) + missing_bos = sorted(set(source_bo_ids) - set(candidate_bo_ids)) + unknown_bos = sorted(set(candidate_bo_ids) - set(source_bo_ids)) + if source_bo_ids != _strings(compile_gate.get("source_bo_ids")) or candidate_bo_ids != _strings(compile_gate.get("candidate_bo_ids")) or missing_bos or unknown_bos or duplicate_bos or compile_gate.get("bo_conservation_pass") is not True: + raise ValueError("candidate/BO conservation failed") + if compile_gate.get("candidate_schema_validation_pass") is not True or compile_gate.get("row_schema_id") != FACT_LEDGER_ROW_SCHEMA_ID or compile_gate.get("row_schema_sha256") != row_schema_sha256: + raise ValueError("candidate row schema compile gate failed") + + deterministic_reviews = _as_list(manifest.get("deterministic_review_items")) + report = { + "schema_version": "stage1_fact_ledger_writer_report.v3", + "status": "READY", + "run_fingerprint": run_fingerprint, + "conservation_gate": { + "source_bo_count": len(source_bo_ids), + "candidate_count": len(candidate_bo_ids), + "final_count": len(rows), + "missing_bo_ids": missing_bos, + "unknown_bo_ids": unknown_bos, + "duplicate_bo_ids": duplicate_bos, + "pass": True, + }, + "schema_gate": { + "schema_path": ROW_SCHEMA_PATH, + "schema_id": FACT_LEDGER_ROW_SCHEMA_ID, + "schema_sha256": row_schema_sha256, + "candidate_validation_pass": True, + "post_patch_validation_pass": True, + "final_pre_write_validation_pass": False, + "post_write_validation_pass": False, + }, + "exception_gate": { + "expected_count": len(expected_ids), + "decision_count": len(decision_ids), + "blocked_count": len(blocked_items), + "missing_exception_ids": missing_ids, + "unknown_exception_ids": unknown_ids, + "duplicate_exception_ids": duplicate_ids, + "coverage_pass": True, + }, + "reference_gate": {"unknown_bo_refs": [], "unknown_structure_refs": [], "unknown_evidence_refs": [], "pass": True}, + "upstream_les_review_gate": {}, + "blocked_review_items": blocked_items, + "row_warnings": [], + "possession_conflict_gate": {"status": "PASS", "conflicts": []}, + "drafting_gates": [], + "gate_review_code_summary": {}, + "evidence_score_change_summary": [], + } + group_universes = _reference_universes(source) + for row in rows: + _validate_row_refs(row, source, group_universes, report) + for key in ("unknown_bo_refs", "unknown_structure_refs", "unknown_evidence_refs"): + report["reference_gate"][key] = sorted(set(report["reference_gate"][key])) + report["reference_gate"]["pass"] = not any(report["reference_gate"][key] for key in ("unknown_bo_refs", "unknown_structure_refs", "unknown_evidence_refs")) + if not report["reference_gate"]["pass"]: + raise ValueError("reference gate failed") + + code_rows: dict[str, list[str]] = {} + for ordinal, row in enumerate(sorted(rows, key=lambda item: item["source_bo_id"]), start=1): + row["fact_id"] = f"F-{ordinal:03d}" + row["skeleton_validation_codes"] = sorted(set(_strings(row.get("skeleton_validation_codes")))) + for code in row["skeleton_validation_codes"]: + code_rows.setdefault(code, []).append(row["fact_id"]) + rows.sort(key=lambda item: item["fact_id"]) + report["gate_review_code_summary"] = {code: ids for code, ids in sorted(code_rows.items())} + conflicts = code_rows.get("POSSESSION_CAUSE_CONFLICT", []) + report["possession_conflict_gate"] = {"status": "HARD_WARNING" if conflicts else "PASS", "conflicts": conflicts} + + blocking_reviews = [item for item in deterministic_reviews if _as_dict(item).get("severity") == "BLOCK_FINAL_DRAFTING" or _as_dict(item).get("review_code") == "LES_BLOCKED_EXCEPTION_REVIEW"] + affected_facts = sorted(set(_text(_as_dict(item).get("fact_id"), 120) for item in deterministic_reviews if _text(_as_dict(item).get("fact_id"), 120))) + affected_bos = sorted(set(_text(_as_dict(item).get("source_bo_id"), 120) for item in deterministic_reviews if _text(_as_dict(item).get("source_bo_id"), 120))) + upstream_status = "BLOCK_FINAL_DRAFTING" if blocking_reviews else ("REVIEW" if deterministic_reviews else "PASS") + report["upstream_les_review_gate"] = { + "review_count": len(deterministic_reviews), + "blocking_count": len(blocking_reviews), + "affected_fact_ids": affected_facts, + "affected_bo_ids": affected_bos, + "status": upstream_status, + } + if blocking_reviews: + report["drafting_gates"].append({ + "gate": "BLOCK_FINAL_DRAFTING_UPSTREAM_LES_REVIEW", + "affected_fact_ids": sorted(set(_text(_as_dict(item).get("fact_id"), 120) for item in blocking_reviews)), + "affected_bo_ids": sorted(set(_text(_as_dict(item).get("source_bo_id"), 120) for item in blocking_reviews)), + }) + notice_rows = code_rows.get("NOTICE_RECEIPT_UNVERIFIED", []) + if notice_rows: + report["drafting_gates"].append({"gate": "BLOCK_FINAL_DRAFTING_NOTICE_RECEIPT", "affected_fact_ids": notice_rows}) + report["evidence_score_change_summary"] = [ + {"fact_id": fact_id, "basis_code": code} + for fact_id, code in sorted(_as_dict(compile_gate.get("evidence_score_basis")).items()) + ] + has_review = bool(deterministic_reviews or blocked_items or report["drafting_gates"] or conflicts) + report["status"] = "READY_WITH_REVIEW" if has_review else "READY" + + for index, row in enumerate(rows): + _validate_row_schema(row, row_validator, f"final_rows[{index}]") + if row.get("source_bo_id") not in set(source_bo_ids): + raise ValueError(f"final_rows[{index}].source_bo_id is outside upstream BO universe") + report["schema_gate"]["final_pre_write_validation_pass"] = True + + _write_json(TARGET_PATH, rows) + reread_rows = _read_json(TARGET_PATH) + if _digest(reread_rows) != _digest(rows) or not isinstance(reread_rows, list): + raise RuntimeError(f"Post-write verification failed: {TARGET_PATH}") + for index, row in enumerate(reread_rows): + if not isinstance(row, dict): + raise ValueError(f"reread_rows[{index}] is not an object") + _validate_row_schema(row, row_validator, f"reread_rows[{index}]") + if row.get("source_bo_id") not in set(source_bo_ids): + raise ValueError(f"reread_rows[{index}].source_bo_id is outside upstream BO universe") + report["schema_gate"]["post_write_validation_pass"] = True + _write_json(REPORT_PATH, report) + _verify_reread(REPORT_PATH, report) + print(_canonical({"fact_ledger_final_writer": { + "schema_version": "stage1_fact_ledger_final_writer.v3", + "status": report["status"], + "written_file": TARGET_PATH, + "writer_report_file": REPORT_PATH, + "run_fingerprint": run_fingerprint, + "fact_count": len(rows), + "row_schema_path": ROW_SCHEMA_PATH, + "row_schema_sha256": row_schema_sha256, + "review_count": len(deterministic_reviews), + "blocked_count": len(blocked_items), + "drafting_gate_count": len(report["drafting_gates"]), + }})) + + + try: + main() + except Exception as exc: + print(_canonical({"fact_ledger_final_writer": { + "schema_version": "stage1_fact_ledger_final_writer.v3", + "status": "FAILED", + "written_file": None, + "error": _safe_error(exc), + }})) + raise