diff --git a/CHANGELOG.md b/CHANGELOG.md index e32a06bc..d554d726 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,6 +23,8 @@ This repo also publishes GitHub Releases. This file is the repo-root release sur stamp downgrades, deadline-interruptible summary lineage expansion, delta refs rebuilt after response-cap eviction, spend-ledger completeness on chunked backfills, and the benchmark evidence trail (`bench/`, F20–F37). +- Added nested-default-JSON-bounded, tool-extracted `lcm_expand_query` evidence provenance so successful and degraded answers retain synthesis-context identities, occurrences, paths, and excerpts while explicitly distinguishing locator coverage from unverified replay, semantic entailment, and caller authorization. + ## v0.20.0 - 2026-07-23 Release focus: Lossless-Claw parity plus the merged cross-session recall and temporal retrieval stack. diff --git a/docs/retrieval-tools.md b/docs/retrieval-tools.md index 3aa44faf..1c1c9fcb 100644 --- a/docs/retrieval-tools.md +++ b/docs/retrieval-tools.md @@ -36,7 +36,7 @@ is the active context engine. | `lcm_load_session` | Load one ordered raw-message transcript page for an explicit `session_id`. This is not search: it returns raw rows in `store_id` order, bounded by `limit`, with per-message content bounded by `max_content_chars`, and continues with `after_store_id` from `next_cursor`. Set `include_exact_ref=true` when rows will feed exact citation or computation; the default response stays byte-compatible. | | `lcm_describe` | Inspect the current-session DAG or preview an `externalized_ref` without loading full content. | | `lcm_expand` | Recover source messages, child summaries, or externalized payloads with pagination. Use `store_id` to fetch a single raw message regardless of session, suitable for drilling into a cross-session `lcm_grep` result. In `store_id` mode, `include_exact_ref=true` adds the exact returned slice without changing default bytes. | -| `lcm_expand_query` | Answer a question using expanded current-session LCM context while returning a bounded answer. | +| `lcm_expand_query` | Answer a question using expanded current-session LCM context while returning a bounded answer plus bounded, tool-extracted `evidence_provenance` for the context supplied to synthesis. | | `lcm_status` | Show runtime health, context pressure, config, source lineage, and lifecycle stats. | | `lcm_inspect` | Read-only operator inventory for current-session lineage, message/frontier metadata, fresh tail, externalized refs/readability, compaction skip/no-op reasons, and matched ignore/stateless patterns. It returns metadata only; use `lcm_load_session`/`lcm_expand` when you need content. | | `lcm_doctor` | Run database, FTS, lifecycle, config, and context-pressure diagnostics. | @@ -243,6 +243,68 @@ the same 20,000-character response ceiling used by the retrieval tools. {"period": "date:2026-07-15", "scope": "global"} ``` +### `lcm_expand_query` evidence provenance + +Every successful, no-match, or structured degraded `lcm_expand_query` response +includes an additive `evidence_provenance` object. On paths that run or attempt +synthesis, the tool builds it from the exact bounded context blocks supplied to +the auxiliary model; no-match returns an explicit empty bundle. It is not +generated by that model. + +The bundle includes up to 24 unique summary, raw-message, and externalized- +payload identities in deterministic first-seen context order. Repeated identities +retain `context_occurrence_count` and up to eight distinct source-path records. +Each record preserves up to eight production-shaped `{node_id, source_index}` +hops plus the producer's original `depth` and `truncated` state, rather than +silently discarding recursive occurrence information. + +The nested object is capped at 10,000 characters using the same default +`json.dumps` representation used by the final tool response (`ensure_ascii=true`). +The cap applies to `evidence_provenance` only, not to the complete +`lcm_expand_query` response. Whole trailing items are removed until the nested +object fits. String metadata is capped at 256 characters; any shortening reports +the original character count and SHA-256 digest. A compact fixed-envelope +fallback preserves the limit for pathological inputs. + +Each quote is capped at 500 characters. `quote_chars_before_provenance_cap` is +the length of the context representation before this 500-character provenance +cap—not necessarily the durable source's original length if retrieval already +sliced it. `quote_truncated_by_provenance_cap` describes only this provenance +cap; `context_truncated` separately describes synthesis-budget omission. + +Items retain available bounded node/store/ref, session, source, role, +content-offset, occurrence/path, and direct `lcm_expand` arguments. If hydrated +content differs from its durable `transcript_content`, both representations are +recorded because both were visible in the synthesis context. + +Response-level and item-level fields separate four claims: + +- `locator_coverage` says whether every unique bounded synthesis-context identity + has locator arguments present (`none`, `partial`, or `complete`); it does not + claim those locators still resolve; +- `locator_replay_status` is `unverified` when locator arguments are present and + `not_available` otherwise; +- `semantic_entailment` is `not_verified` for a completed answer and + `not_applicable` when synthesis did not run or failed; +- `identifiers_are_authority` is always `false`: node/store/session identifiers + locate evidence but do not authenticate a caller or grant access. + +`locator_replay_safety` is `not_guaranteed`: current scalar node/store/ref and +offset locators do not provide durable staleness, revision, or database-generation +detection. Locator presence is therefore not proof of exact replay. + +This PR adds no authorization mechanism. `node_id` and `externalized_ref` retain +their current-session checks; `store_id` remains an intentional cross-session +locator. Hosts/callers must authorize every expansion before invoking it and must +not treat identifiers as capabilities. + +The provenance layer does not claim that a fluent answer is correct, that it +resolved contradictory sources, or that every answer clause is supported. +Existing top-level answer, match, pagination, and degraded-response fields are +unchanged. Unknown-field-tolerant callers remain compatible; strict response +schemas and fixed-size logging/gateway consumers must be updated for the new +nested object. + When temporal rollups are enabled and `ready` rollups cover the **entire** requested window, the response includes their ids and `ready` status in `provenance.rollups`. If any day in the window lacks a ready rollup (missing or diff --git a/schemas.py b/schemas.py index 25c53eef..41f48542 100644 --- a/schemas.py +++ b/schemas.py @@ -1146,6 +1146,7 @@ "query matching summaries/raw messages to expand or explicit node_ids to inspect. Uses the expansion path " "instead of the summarization path so retrieval/synthesis can use a different model or timeout. " "When expanding parent summary nodes, it recursively descends the DAG under the context budget to include leaf evidence where possible. " + "The response includes a nested, default-JSON-bounded, tool-extracted evidence provenance object with locator coverage, while marking locator replay and semantic entailment as unverified and identifiers as non-authoritative. This adds no authorization: node_id and externalized_ref retain current-session checks, store_id remains an intentional cross-session locator, and hosts must authorize expansion before invoking it. " "Prefer this for questions about the active conversation after compaction; for cross-session recall, use session_search first." ), "parameters": { diff --git a/tests/test_expand_query_provenance.py b/tests/test_expand_query_provenance.py new file mode 100644 index 00000000..a097512a --- /dev/null +++ b/tests/test_expand_query_provenance.py @@ -0,0 +1,644 @@ +"""Evidence-provenance contract tests for lcm_expand_query.""" + +import json + +import pytest + +import hermes_lcm.tools as lcm_tools +from hermes_lcm.config import LCMConfig +from hermes_lcm.dag import SummaryNode +from hermes_lcm.engine import LCMEngine +from hermes_lcm.schemas import LCM_EXPAND_QUERY + + +@pytest.fixture +def engine(tmp_path): + config = LCMConfig(database_path=str(tmp_path / "expand-query-provenance.db")) + instance = LCMEngine(config=config) + instance._session_id = "test-session" + instance.context_length = 200_000 + instance.threshold_tokens = int(instance.context_length * config.context_threshold) + try: + yield instance + finally: + instance.shutdown() + + +def _add_message_summary(engine, content="The rollout target is Tuesday at 09:00 UTC."): + store_id = engine._store.append( + "test-session", + {"role": "user", "content": content}, + source="cli", + ) + node_id = engine._dag.add_node( + SummaryNode( + session_id="test-session", + depth=0, + summary="Rollout timing was discussed.", + token_count=8, + source_token_count=12, + source_ids=[store_id], + source_type="messages", + created_at=1, + ) + ) + return store_id, node_id + + +def test_expand_query_returns_tool_extracted_evidence_without_claiming_entailment(engine, monkeypatch): + store_id, node_id = _add_message_summary(engine) + monkeypatch.setattr( + lcm_tools, + "_synthesize_expansion_answer", + lambda **kwargs: "The rollout is Friday and the moon is cheese.", + ) + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + {"prompt": "When is rollout?", "node_ids": [node_id], "context_max_tokens": 1000}, + ) + ) + + assert result["answer"] == "The rollout is Friday and the moon is cheese." + evidence = result["evidence_provenance"] + assert evidence["synthesis_status"] == "completed" + assert evidence["semantic_entailment"] == "not_verified" + assert evidence["quote_origin"] == "tool_extracted_from_synthesis_context" + assert evidence["retrieval_scope"] == {"kind": "current_session", "session_ids": ["test-session"]} + assert evidence["identifiers_are_authority"] is False + assert evidence["locator_replay_safety"] == "not_guaranteed" + assert evidence["locator_coverage"] == "complete" + + summary_item = next(item for item in evidence["items"] if item["source_type"] == "summary") + message_item = next(item for item in evidence["items"] if item["source_type"] == "raw_message") + assert summary_item["node_id"] == node_id + assert summary_item["session_id"] == "test-session" + assert summary_item["quote"] == "Rollout timing was discussed." + assert summary_item["locator_present"] is True + assert summary_item["locator_replay_status"] == "unverified" + assert summary_item["expand_args"] == {"node_id": node_id} + assert message_item["store_id"] == store_id + assert message_item["session_id"] == "test-session" + assert message_item["quote"] == "The rollout target is Tuesday at 09:00 UTC." + assert message_item["locator_present"] is True + assert message_item["locator_replay_status"] == "unverified" + assert message_item["expand_args"] == {"store_id": store_id, "content_offset": 0} + + expanded = json.loads(engine.handle_tool_call("lcm_expand", message_item["expand_args"])) + assert expanded["content"].startswith(message_item["quote"]) + + +def test_expand_query_session_snapshot_survives_rebind_during_synthesis(engine, monkeypatch): + _, node_id = _add_message_summary(engine) + + def synthesize(**kwargs): + engine._session_id = "rebound-session" + return "snapshot answer" + + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", synthesize) + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + { + "prompt": "When is rollout?", + "node_ids": [node_id], + "context_max_tokens": 1000, + }, + ) + ) + + evidence = result["evidence_provenance"] + assert engine.current_session_id == "rebound-session" + assert evidence["retrieval_scope"] == { + "kind": "current_session", + "session_ids": ["test-session"], + } + assert all(item["session_id"] == "test-session" for item in evidence["items"]) + + +def test_expand_query_preserves_evidence_when_synthesis_fails(engine, monkeypatch): + store_id, node_id = _add_message_summary(engine) + + def timeout(**kwargs): + raise TimeoutError("provider timeout") + + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", timeout) + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + {"prompt": "When is rollout?", "node_ids": [node_id], "context_max_tokens": 1000}, + ) + ) + + assert result["degraded"] is True + assert "answer" not in result + assert result["evidence_provenance"]["synthesis_status"] == "failed" + assert result["evidence_provenance"]["semantic_entailment"] == "not_applicable" + assert any(item.get("store_id") == store_id for item in result["evidence_provenance"]["items"]) + + +def test_expand_query_represents_empty_evidence_without_running_synthesis(engine, monkeypatch): + called = False + + def synthesize(**kwargs): + nonlocal called + called = True + return "should not run" + + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", synthesize) + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + {"prompt": "What mentions a missing token?", "query": "NEVER_PRESENT"}, + ) + ) + + assert called is False + assert result["answer"] == "No matching summaries or raw messages found in the current session." + assert result["evidence_provenance"]["items"] == [] + assert result["evidence_provenance"]["locator_coverage"] == "none" + assert result["evidence_provenance"]["synthesis_status"] == "not_run" + assert result["evidence_provenance"]["semantic_entailment"] == "not_applicable" + + +def test_expand_query_evidence_reports_partial_context_separately_from_entailment(engine, monkeypatch): + _, node_id = _add_message_summary(engine, content="abcdef") + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", lambda **kwargs: "bounded answer") + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + { + "prompt": "What raw detail?", + "node_ids": [node_id], + "max_tokens": 5, + "context_max_tokens": 1, + }, + ) + ) + + assert result["context_truncated"] is True + assert result["evidence_provenance"]["context_truncated"] is True + assert result["evidence_provenance"]["semantic_entailment"] == "not_verified" + + +def test_expand_query_evidence_represents_distinct_transcript_content(): + evidence = lcm_tools._build_expand_query_evidence( + [ + { + "type": "messages", + "messages": [ + { + "store_id": 17, + "session_id": "test-session", + "role": "tool", + "source": "cli", + "content": "hydrated external payload", + "transcript_content": "[externalized ref=payload-17]", + "content_source": "externalized_payload", + "content_offset": 0, + "externalized": {"ref": "payload-17"}, + } + ], + } + ], + session_id="test-session", + context_truncated=False, + synthesis_status="completed", + ) + + assert evidence["locator_coverage"] == "complete" + assert evidence["context_unique_item_count"] == 2 + assert evidence["context_occurrence_count"] == 2 + by_source = {item["content_source"]: item for item in evidence["items"]} + assert by_source["externalized_payload"]["quote"] == "hydrated external payload" + assert by_source["externalized_payload"]["expand_args"] == { + "externalized_ref": "payload-17", + "content_offset": 0, + } + assert by_source["transcript_content"]["quote"] == "[externalized ref=payload-17]" + assert by_source["transcript_content"]["expand_args"] == { + "store_id": 17, + "content_offset": 0, + } + + +def test_expand_query_provenance_covers_recursive_descendants(engine, monkeypatch): + store_id, leaf_id = _add_message_summary(engine, content="recursive leaf evidence") + middle_id = engine._dag.add_node( + SummaryNode( + session_id="test-session", + depth=1, + summary="middle summary", + token_count=4, + source_token_count=12, + source_ids=[leaf_id], + source_type="nodes", + created_at=2, + ) + ) + parent_id = engine._dag.add_node( + SummaryNode( + session_id="test-session", + depth=2, + summary="parent summary", + token_count=4, + source_token_count=16, + source_ids=[middle_id], + source_type="nodes", + created_at=3, + ) + ) + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", lambda **kwargs: "recursive answer") + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + {"prompt": "What is the leaf evidence?", "node_ids": [parent_id], "context_max_tokens": 1000}, + ) + ) + + items = result["evidence_provenance"]["items"] + assert any(item.get("node_id") == parent_id for item in items) + middle_item = next(item for item in items if item.get("node_id") == middle_id) + leaf_item = next( + item + for item in items + if item.get("node_id") == leaf_id and item["source_type"] == "summary" + ) + message_item = next(item for item in items if item.get("store_id") == store_id) + assert message_item["quote"] == "recursive leaf evidence" + assert middle_item["source_paths"] == [ + { + "path": [{"node_id": parent_id, "source_index": 0}], + "depth": 1, + "truncated": False, + } + ] + assert leaf_item["source_paths"] == [ + { + "path": [ + {"node_id": parent_id, "source_index": 0}, + {"node_id": middle_id, "source_index": 0}, + ], + "depth": 2, + "truncated": False, + } + ] + assert message_item["source_paths"] == [ + { + "path": [ + {"node_id": parent_id, "source_index": 0}, + {"node_id": middle_id, "source_index": 0}, + {"node_id": leaf_id, "source_index": 0}, + ], + "depth": 3, + "truncated": False, + } + ] + + +def test_expand_query_raw_window_provenance_expands_from_exact_offset(engine, monkeypatch): + content = "prefix-noise " * 80 + "TAILMATCH exact evidence" + store_id = engine._store.append("test-session", {"role": "user", "content": content}, source="cli") + stored = engine._store.get(store_id) + stored["search_rank"] = 1 + stored["snippet"] = "TAILMATCH exact evidence" + monkeypatch.setattr(engine._store, "search", lambda *args, **kwargs: [stored]) + monkeypatch.setattr(engine._dag, "search", lambda *args, **kwargs: []) + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", lambda **kwargs: "raw answer") + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + { + "prompt": "What does TAILMATCH say?", + "query": "TAILMATCH", + "max_tokens": 20, + "context_max_tokens": 8, + }, + ) + ) + + item = next(item for item in result["evidence_provenance"]["items"] if item.get("store_id") == store_id) + assert item["content_source"] == "raw_search_hit" + assert item["content_offset"] == content.index("TAILMATCH") + assert "TAILMATCH" in item["quote"] + expanded = json.loads(engine.handle_tool_call("lcm_expand", item["expand_args"])) + assert expanded["content"].startswith(item["quote"]) + + +def test_expand_query_externalized_provenance_is_directly_expandable(tmp_path, monkeypatch): + config = LCMConfig( + database_path=str(tmp_path / "externalized-provenance.db"), + large_output_externalization_enabled=True, + large_output_externalization_threshold_chars=200, + ) + instance = LCMEngine(config=config, hermes_home=str(tmp_path / "hermes")) + instance._session_id = "test-session" + content = "EXTERNALIZED RAW DETAIL " + ("abcdef" * 100) + try: + instance._serialize_messages([{"role": "tool", "tool_call_id": "call_ext", "content": content}]) + ref = next((tmp_path / "hermes" / "lcm-large-outputs").glob("*.json")).name + placeholder = f"[GC'd externalized tool output: tool_call_id=call_ext; chars={len(content)}; ref={ref}]" + store_id = instance._store.append( + "test-session", + {"role": "tool", "tool_call_id": "call_ext", "content": placeholder}, + ) + node_id = instance._dag.add_node( + SummaryNode( + session_id="test-session", + depth=0, + summary="externalized payload summary", + token_count=10, + source_token_count=200, + source_ids=[store_id], + source_type="messages", + created_at=0, + ) + ) + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", lambda **kwargs: "externalized answer") + + result = json.loads( + instance.handle_tool_call( + "lcm_expand_query", + {"prompt": "What externalized detail exists?", "node_ids": [node_id], "context_max_tokens": 500}, + ) + ) + items = result["evidence_provenance"]["items"] + payload_item = next(item for item in items if item.get("content_source") == "externalized_payload") + transcript_item = next(item for item in items if item.get("content_source") == "transcript_content") + assert payload_item["source_type"] == "externalized_payload" + assert payload_item["externalized_ref"] == ref + assert payload_item["expand_args"] == {"externalized_ref": ref, "content_offset": 0} + assert payload_item["quote"].startswith("EXTERNALIZED RAW DETAIL") + assert transcript_item["source_type"] == "raw_message" + assert transcript_item["quote"] == placeholder + assert transcript_item["expand_args"] == {"store_id": store_id, "content_offset": 0} + expanded = json.loads(instance.handle_tool_call("lcm_expand", payload_item["expand_args"])) + assert expanded["content"].startswith(payload_item["quote"]) + finally: + instance.shutdown() + + +def test_expand_query_evidence_bounds_dedupes_and_preserves_conflicts(): + messages = [] + for store_id in range(1, 27): + content = ("A" if store_id == 1 else "B" if store_id == 2 else f"evidence {store_id}") + if store_id == 1: + content = content * 700 + messages.append( + { + "store_id": store_id, + "session_id": "test-session", + "role": "user", + "source": "cli", + "content": content, + "content_source": "message", + "content_offset": 0, + } + ) + messages.append(dict(messages[0])) + + evidence = lcm_tools._build_expand_query_evidence( + [{"type": "raw_messages", "messages": messages}], + session_id="test-session", + context_truncated=False, + synthesis_status="completed", + ) + + assert evidence["context_unique_item_count"] == 26 + assert evidence["context_occurrence_count"] == 27 + assert 2 <= len(evidence["items"]) <= 24 + assert evidence["items_truncated"] is True + assert evidence["locator_coverage"] == "partial" + assert evidence["serialized_char_limit"] == 10_000 + assert len(json.dumps(evidence)) <= evidence["serialized_char_limit"] + assert evidence["items"][0]["quote_chars_before_provenance_cap"] == 700 + assert evidence["items"][0]["quote_truncated_by_provenance_cap"] is True + assert len(evidence["items"][0]["quote"]) == 500 + assert evidence["items"][1]["quote"] == "B" + + +def test_expand_query_blank_synthesis_retains_failed_provenance(engine, monkeypatch): + _, node_id = _add_message_summary(engine) + monkeypatch.setattr(lcm_tools, "_synthesize_expansion_answer", lambda **kwargs: " ") + + result = json.loads( + engine.handle_tool_call( + "lcm_expand_query", + {"prompt": "When is rollout?", "node_ids": [node_id]}, + ) + ) + + assert result["degraded"] is True + assert result["evidence_provenance"]["synthesis_status"] == "failed" + assert result["evidence_provenance"]["semantic_entailment"] == "not_applicable" + assert result["evidence_provenance"]["items"] + + +def test_expand_query_schema_describes_evidence_boundary(): + description = LCM_EXPAND_QUERY["description"] + + assert "locator coverage" in description.lower() + assert "semantic entailment" in description.lower() + + +def test_expand_query_evidence_cap_uses_final_ascii_escaping(): + messages = [ + { + "store_id": index, + "session_id": "unicode-session", + "role": "user", + "source": "cli", + "content": ("漢字🙂" * 300) + str(index), + "content_source": "message", + "content_offset": 0, + } + for index in range(1, 25) + ] + + evidence = lcm_tools._build_expand_query_evidence( + [{"type": "raw_messages", "messages": messages}], + session_id="unicode-session", + context_truncated=False, + synthesis_status="completed", + ) + + assert evidence["serialization"] == { + "scope": "evidence_provenance_only", + "json_ensure_ascii": True, + } + assert len(json.dumps(evidence)) <= evidence["serialized_char_limit"] + assert len(evidence["items"]) < len(messages) + assert evidence["items_truncated"] is True + + +def test_expand_query_evidence_bounds_oversized_scope_and_item_metadata(): + long_session = "session-" + ("x" * 20_000) + long_ref = "ref-" + ("界" * 10_000) + evidence = lcm_tools._build_expand_query_evidence( + [ + { + "type": "raw_messages", + "messages": [ + { + "store_id": 7, + "session_id": long_session, + "role": "role-" + ("r" * 2_000), + "source": "source-" + ("s" * 2_000), + "content": "bounded payload", + "content_source": "externalized_payload", + "content_offset": 0, + "externalized": {"ref": long_ref}, + } + ], + } + ], + session_id=long_session, + context_truncated=False, + synthesis_status="completed", + ) + + assert len(json.dumps(evidence)) <= evidence["serialized_char_limit"] + assert evidence["metadata_truncated"] is True + assert len(evidence["retrieval_scope"]["session_ids"][0]) == 256 + scope_meta = evidence["retrieval_scope"]["metadata_truncation"]["session_ids[0]"] + assert scope_meta["original_chars"] == len(long_session) + assert len(scope_meta["sha256"]) == 64 + item = evidence["items"][0] + assert item["locator_present"] is False + assert item["locator_replay_status"] == "not_available" + assert "expand_args" not in item + assert item["metadata_truncation"]["externalized_ref"]["original_chars"] == len(long_ref) + + +def test_expand_query_locator_presence_does_not_claim_verified_replay(): + evidence = lcm_tools._build_expand_query_evidence( + [ + { + "type": "raw_messages", + "messages": [ + { + "store_id": 999_999_999, + "session_id": "missing-session", + "content": "context can carry a stale locator", + "content_source": "message", + "content_offset": 0, + } + ], + } + ], + session_id="missing-session", + context_truncated=False, + synthesis_status="completed", + ) + + item = evidence["items"][0] + assert item["locator_present"] is True + assert item["locator_replay_status"] == "unverified" + assert evidence["locator_coverage"] == "complete" + assert evidence["locator_replay_safety"] == "not_guaranteed" + + +def test_expand_query_dedupe_preserves_production_paths_and_occurrences(): + truncated_path = [ + {"node_id": node_id, "source_index": node_id - 3} + for node_id in range(3, 11) + ] + blocks = [ + { + "type": "child_nodes", + "node_id": 1, + "children": [ + { + "node_id": 99, + "source_index": 0, + "summary": "shared descendant", + } + ], + }, + { + "type": "descendant_child_nodes", + "node_id": 20, + "source_path": truncated_path, + "source_path_depth": 10, + "source_path_truncated": True, + "children": [ + { + "node_id": 99, + "source_index": 3, + "summary": "shared descendant", + } + ], + }, + ] + + evidence = lcm_tools._build_expand_query_evidence( + blocks, + session_id="test-session", + context_truncated=False, + synthesis_status="completed", + ) + + assert evidence["context_unique_item_count"] == 1 + assert evidence["context_occurrence_count"] == 2 + item = evidence["items"][0] + assert item["context_occurrence_count"] == 2 + assert item["source_paths"] == [ + { + "path": [{"node_id": 1, "source_index": 0}], + "depth": 1, + "truncated": False, + }, + { + "path": [ + *truncated_path[1:], + {"node_id": 20, "source_index": 3}, + ], + "depth": 11, + "truncated": True, + }, + ] + + +def test_expand_query_child_message_paths_include_final_dag_edge(): + evidence = lcm_tools._build_expand_query_evidence( + [ + { + "type": "child_messages", + "node_id": 7, + "source_path": [{"node_id": 1, "source_index": 0}], + "source_path_depth": 1, + "messages": [ + { + "store_id": 42, + "source_index": 2, + "session_id": "test-session", + "content": "hydrated child evidence", + "transcript_content": "stored child transcript", + } + ], + } + ], + session_id="test-session", + context_truncated=False, + synthesis_status="completed", + ) + + assert evidence["context_unique_item_count"] == 2 + for item in evidence["items"]: + assert item["source_paths"] == [ + { + "path": [ + {"node_id": 1, "source_index": 0}, + {"node_id": 7, "source_index": 2}, + ], + "depth": 2, + "truncated": False, + } + ] diff --git a/tools.py b/tools.py index 766445ea..981a1ecc 100644 --- a/tools.py +++ b/tools.py @@ -3,6 +3,7 @@ from __future__ import annotations import copy +import hashlib import json import logging import re @@ -139,9 +140,15 @@ def _require_engine(kwargs: Dict[str, Any]) -> "LCMEngine | None": return engine if engine is not None else None -def _get_session_node(engine: "LCMEngine", node_id: int): +def _get_session_node( + engine: "LCMEngine", + node_id: int, + *, + session_id: str | None = None, +): + session_id = engine.current_session_id if session_id is None else session_id node = engine._dag.get_node(node_id) - if node is None or node.session_id != engine.current_session_id: + if node is None or node.session_id != session_id: return None return node @@ -1187,9 +1194,11 @@ def _expand_message_sources( source_limit: int | None = None, content_offset: int = 0, hydrate_externalized_content: bool = False, + session_id: str | None = None, ) -> tuple[list[dict[str, Any]], dict[str, Any]]: from .tokens import count_tokens + session_id = engine.current_session_id if session_id is None else session_id total_sources = len(node.source_ids) source_offset = min(max(0, source_offset), total_sources) remaining_source_count = max(0, total_sources - source_offset) @@ -1232,7 +1241,7 @@ def _expand_message_sources( ref_payload = _get_externalized_payload( engine, ref, - allowed_session_ids={engine.current_session_id, stored.get("session_id", "")}, + allowed_session_ids={session_id, stored.get("session_id", "")}, ) if ref_payload is not None and ref_payload.get("kind") != "ingest_payload": externalized = ref_payload @@ -1249,7 +1258,7 @@ def _expand_message_sources( "source_index": source_index, "session_id": stored.get("session_id", ""), "source": stored.get("source") or "", - "from_current_session": stored.get("session_id", "") == engine.current_session_id, + "from_current_session": stored.get("session_id", "") == session_id, "role": stored["role"], "content": sliced["content"], "content_chars": sliced["content_chars"], @@ -1330,9 +1339,11 @@ def _expand_child_nodes( *, source_offset: int = 0, source_limit: int | None = None, + session_id: str | None = None, ) -> tuple[list[dict[str, Any]], dict[str, Any]]: from .tokens import count_tokens + session_id = engine.current_session_id if session_id is None else session_id total_sources = len(node.source_ids) source_offset = min(max(0, source_offset), total_sources) remaining_source_count = max(0, total_sources - source_offset) @@ -1344,7 +1355,7 @@ def _expand_child_nodes( children: list[tuple[int, Any]] = [] for relative_index, child_id in enumerate(selected_source_ids): child = engine._dag.get_node(child_id) - if child is None or child.session_id != engine.current_session_id: + if child is None or child.session_id != session_id: continue children.append((source_offset + relative_index, child)) @@ -1416,6 +1427,7 @@ def _collect_descendant_evidence_blocks( visited_node_ids: set[int] | None = None, source_path: list[dict[str, int]] | None = None, remaining_node_visits: list[int] | None = None, + session_id: str | None = None, ) -> list[dict[str, Any]]: if max_tokens <= 0 or node.source_type != "nodes": return [] @@ -1429,6 +1441,7 @@ def _collect_descendant_evidence_blocks( # walk an unbounded number of nodes while normal deep summaries still # reach their leaf evidence. remaining_node_visits = [max(64, int(max_tokens) * 4)] + session_id = engine.current_session_id if session_id is None else session_id if remaining_node_visits[0] <= 0: return [] @@ -1447,7 +1460,7 @@ def _collect_descendant_evidence_blocks( stack.append((current, current_path, current_visited, source_index + 1)) child_id = current.source_ids[source_index] child = engine._dag.get_node(child_id) - if child is None or child.session_id != engine.current_session_id: + if child is None or child.session_id != session_id: continue child_node_id = int(child.node_id) if child_node_id in current_visited: @@ -1462,6 +1475,7 @@ def _collect_descendant_evidence_blocks( child, max_tokens=remaining_tokens, hydrate_externalized_content=hydrate_externalized_content, + session_id=session_id, ) if messages or pagination.get("has_more"): block = { @@ -1479,7 +1493,12 @@ def _collect_descendant_evidence_blocks( continue if child.source_type == "nodes": - children, pagination = _expand_child_nodes(engine, child, max_tokens=remaining_tokens) + children, pagination = _expand_child_nodes( + engine, + child, + max_tokens=remaining_tokens, + session_id=session_id, + ) if children or pagination.get("has_more"): block = { "type": "descendant_child_nodes", @@ -1504,9 +1523,11 @@ def _collect_context_blocks_for_node( max_tokens: int, *, hydrate_externalized_content: bool = False, + session_id: str | None = None, ) -> list[dict[str, Any]]: from .tokens import count_tokens + session_id = engine.current_session_id if session_id is None else session_id summary, summary_truncated = _truncate_text_to_token_budget(node.summary, max_tokens) blocks: list[dict[str, Any]] = [ { @@ -1527,6 +1548,7 @@ def _collect_context_blocks_for_node( node, max_tokens=remaining_tokens, hydrate_externalized_content=hydrate_externalized_content, + session_id=session_id, ) if messages or pagination.get("has_more"): block = { @@ -1537,7 +1559,12 @@ def _collect_context_blocks_for_node( } blocks.append(block) elif node.source_type == "nodes": - children, pagination = _expand_child_nodes(engine, node, max_tokens=remaining_tokens) + children, pagination = _expand_child_nodes( + engine, + node, + max_tokens=remaining_tokens, + session_id=session_id, + ) if children or pagination.get("has_more"): blocks.append( { @@ -1556,6 +1583,7 @@ def _collect_context_blocks_for_node( node, max_tokens=descendant_tokens, hydrate_externalized_content=hydrate_externalized_content, + session_id=session_id, ) ) @@ -1678,6 +1706,459 @@ def _context_content_token_count(blocks: list[dict[str, Any]]) -> int: return total +_EXPAND_QUERY_EVIDENCE_QUOTE_MAX_CHARS = 500 +_EXPAND_QUERY_EVIDENCE_MAX_ITEMS = 24 +_EXPAND_QUERY_EVIDENCE_MAX_SERIALIZED_CHARS = 10_000 +_EXPAND_QUERY_EVIDENCE_METADATA_MAX_CHARS = 256 +_EXPAND_QUERY_EVIDENCE_SOURCE_PATH_MAX_HOPS = 8 +_EXPAND_QUERY_EVIDENCE_SOURCE_PATH_MAX_PATHS = 8 + + +def _expand_query_evidence_serialized_chars(evidence: dict[str, Any]) -> int: + """Measure the exact JSON representation used by the tool response.""" + + return len(json.dumps(evidence)) + + +def _bounded_evidence_text(value: Any) -> tuple[str, dict[str, Any] | None]: + """Bound metadata while retaining explicit length and digest disclosure.""" + + text = str(value or "") + if len(text) <= _EXPAND_QUERY_EVIDENCE_METADATA_MAX_CHARS: + return text, None + return text[:_EXPAND_QUERY_EVIDENCE_METADATA_MAX_CHARS], { + "original_chars": len(text), + "sha256": hashlib.sha256(text.encode("utf-8")).hexdigest(), + } + + +def _normalized_evidence_source_path( + value: Any, + *, + source_path_depth: Any = None, + source_path_truncated: Any = None, +) -> dict[str, Any] | None: + if not isinstance(value, (list, tuple)): + return None + result: list[dict[str, int]] = [] + for part in value[-_EXPAND_QUERY_EVIDENCE_SOURCE_PATH_MAX_HOPS:]: + if not isinstance(part, dict): + continue + node_id = part.get("node_id") + source_index = part.get("source_index") + if isinstance(node_id, int) and isinstance(source_index, int): + result.append({"node_id": node_id, "source_index": source_index}) + if not result: + return None + depth = ( + source_path_depth + if isinstance(source_path_depth, int) and source_path_depth >= len(result) + else len(result) + ) + return { + "path": result, + "depth": depth, + "truncated": bool(source_path_truncated or depth > len(result)), + } + + +def _evidence_source_path_with_final_edge( + block: dict[str, Any], + source_index: Any, +) -> tuple[Any, Any, Any]: + """Return a block path extended through its selected child/message edge.""" + + source_path = block.get("source_path") + path = list(source_path) if isinstance(source_path, (list, tuple)) else [] + depth = block.get("source_path_depth") + if not isinstance(depth, int) or depth < len(path): + depth = len(path) + node_id = block.get("node_id") + if isinstance(node_id, int) and isinstance(source_index, int): + path.append({"node_id": node_id, "source_index": source_index}) + depth += 1 + return path, depth, block.get("source_path_truncated") + + +def _build_expand_query_evidence( + context_blocks: list[dict[str, Any]], + *, + session_id: str, + context_truncated: bool, + synthesis_status: str, +) -> dict[str, Any]: + """Return bounded, tool-extracted provenance for expansion synthesis context. + + The bundle identifies excerpts that were actually sent to the auxiliary + model. It deliberately does not claim that the resulting answer is + semantically entailed by those excerpts. + """ + + synthesis_status = ( + synthesis_status + if synthesis_status in {"completed", "failed", "not_run"} + else "unknown" + ) + candidates: list[dict[str, Any]] = [] + seen_items: dict[str, dict[str, Any]] = {} + + def put_bounded_text(item: dict[str, Any], field: str, value: Any) -> bool: + bounded, truncation = _bounded_evidence_text(value) + item[field] = bounded + if truncation is not None: + item.setdefault("metadata_truncation", {})[field] = truncation + return True + return False + + def add_candidate( + identity: str, + item: dict[str, Any], + *, + source_path: Any = None, + source_path_depth: Any = None, + source_path_truncated: Any = None, + ) -> None: + path = _normalized_evidence_source_path( + source_path, + source_path_depth=source_path_depth, + source_path_truncated=source_path_truncated, + ) + existing = seen_items.get(identity) + if existing is not None: + existing["context_occurrence_count"] += 1 + source_paths = existing.setdefault("source_paths", []) + if path and path not in source_paths: + if len(source_paths) < _EXPAND_QUERY_EVIDENCE_SOURCE_PATH_MAX_PATHS: + source_paths.append(path) + else: + existing["source_paths_truncated"] = True + return + put_bounded_text(item, "evidence_id", identity) + item["context_occurrence_count"] = 1 + if path: + item["source_paths"] = [path] + seen_items[identity] = item + candidates.append(item) + + def quote_payload(content: Any) -> dict[str, Any]: + quote = str(content or "") + return { + "quote": quote[:_EXPAND_QUERY_EVIDENCE_QUOTE_MAX_CHARS], + "quote_chars_before_provenance_cap": len(quote), + "quote_truncated_by_provenance_cap": len(quote) + > _EXPAND_QUERY_EVIDENCE_QUOTE_MAX_CHARS, + } + + for block_index, block in enumerate(context_blocks): + if not isinstance(block, dict): + continue + block_type = block.get("type") + + if block_type == "summary": + node_id = block.get("node_id") + identity = ( + f"node:{node_id}" + if isinstance(node_id, int) + else f"context:{block_index}:summary" + ) + item: dict[str, Any] = { + "source_type": "summary", + "node_id": node_id if isinstance(node_id, int) else None, + "locator_present": isinstance(node_id, int), + "locator_replay_status": "unverified" if isinstance(node_id, int) else "not_available", + **quote_payload(block.get("summary")), + } + put_bounded_text(item, "session_id", block.get("session_id") or session_id) + if isinstance(node_id, int): + item["expand_args"] = {"node_id": node_id} + add_candidate( + identity, + item, + source_path=block.get("source_path"), + source_path_depth=block.get("source_path_depth"), + source_path_truncated=block.get("source_path_truncated"), + ) + + if block_type in {"child_nodes", "descendant_child_nodes"}: + for child_index, child in enumerate(block.get("children", []) or []): + if not isinstance(child, dict): + continue + node_id = child.get("node_id") + identity = ( + f"node:{node_id}" + if isinstance(node_id, int) + else f"context:{block_index}:child:{child_index}" + ) + item = { + "source_type": "summary", + "node_id": node_id if isinstance(node_id, int) else None, + "parent_node_id": ( + block.get("node_id") + if isinstance(block.get("node_id"), int) + else None + ), + "source_index": ( + child.get("source_index") + if isinstance(child.get("source_index"), int) + else None + ), + "locator_present": isinstance(node_id, int), + "locator_replay_status": "unverified" if isinstance(node_id, int) else "not_available", + **quote_payload(child.get("summary")), + } + put_bounded_text( + item, + "session_id", + child.get("session_id") or block.get("session_id") or session_id, + ) + if isinstance(node_id, int): + item["expand_args"] = {"node_id": node_id} + source_path, source_path_depth, source_path_truncated = ( + _evidence_source_path_with_final_edge( + block, + child.get("source_index"), + ) + ) + add_candidate( + identity, + item, + source_path=source_path, + source_path_depth=source_path_depth, + source_path_truncated=source_path_truncated, + ) + + if block_type not in {"messages", "child_messages", "raw_messages"}: + continue + for message_index, message in enumerate(block.get("messages", []) or []): + if not isinstance(message, dict): + continue + store_id = message.get("store_id") + content_offset = message.get("content_offset") + if not isinstance(content_offset, int) or content_offset < 0: + content_offset = 0 + content_source = str(message.get("content_source") or "message") + externalized = message.get("externalized") or {} + externalized_ref = externalized.get("ref") if isinstance(externalized, dict) else None + if content_source == "externalized_payload" and externalized_ref: + identity = f"externalized:{externalized_ref}:{content_offset}" + source_type = "externalized_payload" + elif isinstance(store_id, int): + identity = f"store:{store_id}:{content_offset}:{content_source}" + source_type = "raw_message" + else: + identity = f"context:{block_index}:message:{message_index}" + source_type = "raw_message" + + item = { + "source_type": source_type, + "store_id": store_id if isinstance(store_id, int) else None, + "node_id": ( + block.get("node_id") + if isinstance(block.get("node_id"), int) + else None + ), + "parent_node_id": ( + block.get("parent_node_id") + if isinstance(block.get("parent_node_id"), int) + else None + ), + "source_index": ( + message.get("source_index") + if isinstance(message.get("source_index"), int) + else None + ), + "content_offset": content_offset, + **quote_payload(message.get("content")), + } + put_bounded_text(item, "session_id", message.get("session_id") or "") + put_bounded_text(item, "source", message.get("source") or "") + put_bounded_text(item, "role", message.get("role")) + put_bounded_text(item, "content_source", content_source) + ref_truncated = put_bounded_text(item, "externalized_ref", externalized_ref) + if content_source == "externalized_payload": + if externalized_ref and not ref_truncated: + item["expand_args"] = { + "externalized_ref": str(externalized_ref), + "content_offset": content_offset, + } + elif isinstance(store_id, int): + item["expand_args"] = {"store_id": store_id, "content_offset": content_offset} + item["locator_present"] = "expand_args" in item + item["locator_replay_status"] = ( + "unverified" if item["locator_present"] else "not_available" + ) + source_path = block.get("source_path") + source_path_depth = block.get("source_path_depth") + source_path_truncated = block.get("source_path_truncated") + if block_type in {"messages", "child_messages"}: + source_path, source_path_depth, source_path_truncated = ( + _evidence_source_path_with_final_edge( + block, + message.get("source_index"), + ) + ) + add_candidate( + identity, + item, + source_path=source_path, + source_path_depth=source_path_depth, + source_path_truncated=source_path_truncated, + ) + + # Externalized/hydrated context blocks can contain both the text + # exposed as ``content`` and the durable transcript representation + # exposed as ``transcript_content``. Both fields are serialized into + # the auxiliary model prompt, so traceability cannot be called + # complete unless both are represented when they differ. + transcript_content = message.get("transcript_content") + if transcript_content is not None and str(transcript_content) != str( + message.get("content") or "" + ): + transcript_identity = ( + f"store:{store_id}:0:transcript_content" + if isinstance(store_id, int) + else f"context:{block_index}:message:{message_index}:transcript_content" + ) + transcript_item = { + "source_type": "raw_message", + "store_id": store_id if isinstance(store_id, int) else None, + "node_id": ( + block.get("node_id") + if isinstance(block.get("node_id"), int) + else None + ), + "parent_node_id": ( + block.get("parent_node_id") + if isinstance(block.get("parent_node_id"), int) + else None + ), + "source_index": ( + message.get("source_index") + if isinstance(message.get("source_index"), int) + else None + ), + "content_offset": 0, + "locator_present": isinstance(store_id, int), + "locator_replay_status": "unverified" if isinstance(store_id, int) else "not_available", + **quote_payload(transcript_content), + } + put_bounded_text(transcript_item, "session_id", message.get("session_id") or "") + put_bounded_text(transcript_item, "source", message.get("source") or "") + put_bounded_text(transcript_item, "role", message.get("role")) + put_bounded_text(transcript_item, "content_source", "transcript_content") + if isinstance(store_id, int): + transcript_item["expand_args"] = { + "store_id": store_id, + "content_offset": 0, + } + add_candidate( + transcript_identity, + transcript_item, + source_path=source_path, + source_path_depth=source_path_depth, + source_path_truncated=source_path_truncated, + ) + + items = candidates[:_EXPAND_QUERY_EVIDENCE_MAX_ITEMS] + bounded_scope_session, scope_truncation = _bounded_evidence_text(session_id) + retrieval_scope: dict[str, Any] = { + "kind": "current_session", + "session_ids": [bounded_scope_session], + } + if scope_truncation is not None: + retrieval_scope["metadata_truncation"] = {"session_ids[0]": scope_truncation} + evidence = { + "retrieval_scope": retrieval_scope, + "identifiers_are_authority": False, + "locator_replay_safety": "not_guaranteed", + "synthesis_status": synthesis_status, + "semantic_entailment": "not_verified" if synthesis_status == "completed" else "not_applicable", + "quote_origin": "tool_extracted_from_synthesis_context", + "context_truncated": bool(context_truncated), + "context_unique_item_count": len(candidates), + "context_occurrence_count": sum( + int(item.get("context_occurrence_count") or 1) for item in candidates + ), + "serialized_char_limit": _EXPAND_QUERY_EVIDENCE_MAX_SERIALIZED_CHARS, + "serialization": { + "scope": "evidence_provenance_only", + "json_ensure_ascii": True, + }, + "quote_max_chars": _EXPAND_QUERY_EVIDENCE_QUOTE_MAX_CHARS, + "items": items, + } + + def refresh_bounds() -> None: + items_truncated = len(items) < len(candidates) + located_items = sum(1 for item in items if item.get("locator_present")) + if not candidates: + locator_coverage = "none" + elif items_truncated or located_items < len(candidates): + locator_coverage = "partial" + else: + locator_coverage = "complete" + evidence["locator_coverage"] = locator_coverage + evidence["items_truncated"] = items_truncated + evidence["quotes_truncated_by_provenance_cap"] = sum( + 1 for item in items if item.get("quote_truncated_by_provenance_cap") + ) + evidence["metadata_truncated"] = bool( + retrieval_scope.get("metadata_truncation") + or any(item.get("metadata_truncation") for item in items) + ) + + refresh_bounds() + while ( + items + and _expand_query_evidence_serialized_chars(evidence) + > _EXPAND_QUERY_EVIDENCE_MAX_SERIALIZED_CHARS + ): + items.pop() + refresh_bounds() + if ( + _expand_query_evidence_serialized_chars(evidence) + > _EXPAND_QUERY_EVIDENCE_MAX_SERIALIZED_CHARS + ): + # Bounded metadata makes this unreachable under normal inputs. Keep a + # compact fail-safe so even a pathological fixed envelope cannot violate + # the public nested-object contract. + evidence = { + "retrieval_scope": { + "kind": "current_session", + "session_ids": [], + "session_id_sha256": hashlib.sha256( + str(session_id or "").encode("utf-8") + ).hexdigest(), + "session_id_omitted_due_to_size": True, + }, + "identifiers_are_authority": False, + "locator_replay_safety": "not_guaranteed", + "synthesis_status": synthesis_status, + "semantic_entailment": ( + "not_verified" if synthesis_status == "completed" else "not_applicable" + ), + "quote_origin": "tool_extracted_from_synthesis_context", + "context_truncated": bool(context_truncated), + "context_unique_item_count": len(candidates), + "context_occurrence_count": sum( + int(item.get("context_occurrence_count") or 1) for item in candidates + ), + "serialized_char_limit": _EXPAND_QUERY_EVIDENCE_MAX_SERIALIZED_CHARS, + "serialization": { + "scope": "evidence_provenance_only", + "json_ensure_ascii": True, + }, + "quote_max_chars": _EXPAND_QUERY_EVIDENCE_QUOTE_MAX_CHARS, + "locator_coverage": "none" if not candidates else "partial", + "items_truncated": bool(candidates), + "quotes_truncated_by_provenance_cap": 0, + "metadata_truncated": True, + "fixed_envelope_compacted": True, + "items": [], + } + return evidence + + def _synthesize_expansion_answer( *, prompt: str, @@ -5505,6 +5986,7 @@ def lcm_expand_query(args: Dict[str, Any], **kwargs) -> str: engine = _require_engine(kwargs) if engine is None: return json.dumps({"error": "LCM engine not initialized"}) + session_id = engine.current_session_id prompt = str(args.get("prompt") or "").strip() if not prompt: @@ -5543,12 +6025,12 @@ def _parse_int_arg(name: str, default: int) -> tuple[int | None, str | None]: parsed_node_id = int(node_id) except (TypeError, ValueError): return json.dumps({"error": "node_ids must contain only integers"}) - node = _get_session_node(engine, parsed_node_id) + node = _get_session_node(engine, parsed_node_id, session_id=session_id) if node is not None: nodes.append(node) elif query: - nodes = engine._dag.search(query, session_id=engine.current_session_id, limit=max_results) - raw_results = engine._store.search(query, session_id=engine.current_session_id, limit=max_results) + nodes = engine._dag.search(query, session_id=session_id, limit=max_results) + raw_results = engine._store.search(query, session_id=session_id, limit=max_results) else: return json.dumps({"error": "Provide either query or node_ids"}) @@ -5561,6 +6043,12 @@ def _parse_int_arg(name: str, default: int) -> tuple[int | None, str | None]: "node_ids": [], "matches": [], "raw_matches": [], + "evidence_provenance": _build_expand_query_evidence( + [], + session_id=session_id, + context_truncated=False, + synthesis_status="not_run", + ), } ) @@ -5573,6 +6061,7 @@ def _parse_int_arg(name: str, default: int) -> tuple[int | None, str | None]: node, max_tokens=remaining_context_tokens, hydrate_externalized_content=True, + session_id=session_id, ) context_blocks.extend(node_blocks) context_budget_used += _context_content_token_count(node_blocks) @@ -5715,6 +6204,12 @@ def _degraded_payload(reason: str, *, include_timeout: bool = False) -> str: "node_ids": node_ids, "matches": matches, "raw_matches": raw_matches, + "evidence_provenance": _build_expand_query_evidence( + context_blocks, + session_id=session_id, + context_truncated=context_truncated, + synthesis_status="failed", + ), } if include_timeout: payload["timeout_seconds"] = timeout @@ -5755,6 +6250,12 @@ def _degraded_payload(reason: str, *, include_timeout: bool = False) -> str: "node_ids": node_ids, "matches": matches, "raw_matches": raw_matches, + "evidence_provenance": _build_expand_query_evidence( + context_blocks, + session_id=session_id, + context_truncated=context_truncated, + synthesis_status="completed", + ), } )