From a9e65e974fba4c75fdff5bca3f7b817773097c4e Mon Sep 17 00:00:00 2001 From: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com> Date: Wed, 23 Sep 2026 05:20:52 +0300 Subject: [PATCH] fix: retain producer tool payloads beside dispatcher annotations Preserve exact machine-readable sources through task artifact promotion while keeping warning and outcome projections intact. --- docs/architecture/13-external-skills-layer.md | 2 + ouroboros/loop_tool_execution.py | 32 ++++- ouroboros/observability.py | 16 +-- ouroboros/size_ratchet_manifest.py | 1 + ouroboros/tools/extension_dispatch.py | 43 +++--- ouroboros/tools/tool_result.py | 27 +++- tests/test_reference_book_budgets.py | 3 +- tests/test_registry_core.py | 2 + tests/test_tool_producer_promotion.py | 65 +++++++++ tests/test_tool_producer_sources.py | 131 ++++++++++++++++++ 10 files changed, 286 insertions(+), 36 deletions(-) create mode 100644 tests/test_tool_producer_promotion.py create mode 100644 tests/test_tool_producer_sources.py diff --git a/docs/architecture/13-external-skills-layer.md b/docs/architecture/13-external-skills-layer.md index bb2c4441a..b07508ed1 100644 --- a/docs/architecture/13-external-skills-layer.md +++ b/docs/architecture/13-external-skills-layer.md @@ -16,4 +16,6 @@ The optional `model_experience` manifest prose — what the skill adds to the mo Marketplace installs are bounded archives staged privately and landed atomically with per-file hash checks. Inert resource bytes, PDF/PPTX included, need no suffix allowlist; archive validity, extraction paths, checksum and actual loader support remain separate facts, and Cyber policy findings are warnings. Install metadata drives isolated dependencies (`marketplace/install_specs.py`); manual instructions remain guidance, not execution. Fresh review and hash-covered declarations authorize literal build/check argv; verified resources and package caches survive under `state/skills//dependency_cache`, and `deps.json` records resolved package metadata, bytes, outputs and diagnostics. Delivery is distinct from a declared executable check: absent checks remain unknown, and declared-output drift invalidates installed readiness. Large resources stay outside the payload and the Git patch. Marketplace update and adopt rollback (`install.PayloadRollbackSnapshot`, verified `rolled_back`): §6 Skills and extensions; script runtimes and their Go/Deno specifics: `tools/skill_exec.py`, `docs/CREATING_SKILLS.md`. +Tool results retain two views when the host adds route or safety notes. `tools/tool_result.py` captures `producer_text` before composition and keeps `host_annotations` outside the bounded metadata map; the existing `text` projection and typed status remain authoritative for model/review behavior. Builtin, extension and MCP dispatch use this same composer. `loop_tool_execution.py` records both views in observability and stores producer bytes separately through the existing write-once task source store. `PRODUCER_RESULT_SOURCE_JSON` is the parser-readable source, distinct from the complete annotated `FULL_RESULT_SOURCE_JSON` used for review. Notes remain visible after truncation; a failed source write is disclosed rather than presented as a usable file. Unannotated results and historical retained files keep their existing representation. + Extensions import through staged trees (`_stage_extension_import_tree` under `__extension_imports/`), so concurrent workers cannot remove a peer's live import and stale trees stay reclaimable. Per-call child processes may proxy tools/routes/WS/UI/settings/companion descriptors, while persistent subscriptions and supervised tasks require in-process or companion lifecycle — reported through the generic capability matrix, never inferred from a platform name. Isolated children run the same staged loader in a private base with a scrubbed env, so native crashes cannot kill `server.py`. In-process extensions are more powerful, which is why namespacing, declared permissions, per-skill tracking and atomic unload are an executable contract rather than convention. Skill repair enqueues an ordinary managed task with an exact selected resource and revision admission (`skill_repair_admission.py`); a legacy `skill_repair` selector remains readable but selects no reduced profile. Shell, browser and delegation remain normal task capabilities; an installed payload needs no mandatory Git copy, and review never forces unload merely because the caller is repairing it. diff --git a/ouroboros/loop_tool_execution.py b/ouroboros/loop_tool_execution.py index a17ab7497..5bfbcfa8a 100644 --- a/ouroboros/loop_tool_execution.py +++ b/ouroboros/loop_tool_execution.py @@ -397,11 +397,13 @@ def _persist_truncated_tool_source( tool_call_id: str, result: Any, tool_args: Optional[Dict[str, Any]] = None, + *, + force: bool = False, ) -> Dict[str, Any]: """Write a generic over-limit result to this actor's existing artifact root.""" text = str(result) - if ( + if not force and ( _should_skip_tool_result_truncation(tool_name, tool_args) or len(text) <= _tool_result_limit(tool_name) ): @@ -732,6 +734,9 @@ def _execute_single_tool( "round_id": correlation.get("round_id"), "args": args, "result": result, + **({"producer_result": tool_result.producer_text, + "host_annotations": list(tool_result.host_annotations)} + if tool_result.producer_text is not None else {}), "tool_ok": tool_ok, "semantic_ok": not is_error, "result_meta": result_meta, @@ -1311,7 +1316,8 @@ def _maybe_auto_attach_image( typed = exec_result.get("tool_result") if isinstance(typed, ToolResult) and typed.code == "TOOL_REPORTED_FAILURE": return - raw = exec_result.get("result") + raw = (typed.producer_text if isinstance(typed, ToolResult) + and typed.producer_text is not None else exec_result.get("result")) if not isinstance(raw, str) or '"auto_attach_image"' not in raw: return observation = None @@ -1402,6 +1408,25 @@ def process_tool_results( tool_args=exec_result.get("tool_args"), source_ref=result_source_ref, ) + typed = exec_result.get("tool_result") + producer_ref = {} + if isinstance(typed, ToolResult) and typed.producer_text is not None: + if ctx is not None: + producer_ref = _persist_truncated_tool_source( + ctx, fn_name, str(exec_result["tool_call_id"]) + ".producer", + typed.producer_text, force=True, + ) + # The complete annotated source above remains review evidence. This + # separate source is for parsers; neither bytes nor notes live in meta. + if result_partial and typed.host_annotations: + truncated_result += "\n\n" + "\n\n".join(typed.host_annotations) + truncated_result += ( + "\nPRODUCER_RESULT_SOURCE_JSON=" + json.dumps(producer_ref, ensure_ascii=False) + + "\nUnannotated tool data for programmatic reading; host notes and outcome still apply." + if producer_ref else + "\nPRODUCER_RESULT_SOURCE_UNAVAILABLE=true" + "\nHost notes remain in the result; no clean producer file was retained." + ) messages.append({ "role": "tool", @@ -1448,6 +1473,9 @@ def process_tool_results( ), } if result_partial else {}), **(exec_result.get("result_meta") or {}), + **({"producer_source_ref": producer_ref, + "host_annotations": list(typed.host_annotations)} + if isinstance(typed, ToolResult) and typed.producer_text is not None else {}), }) if fn_name == "task_acceptance_review" and not is_error: raw = str(exec_result.get("result") or "") diff --git a/ouroboros/observability.py b/ouroboros/observability.py index 2025e9928..d7bb77d64 100644 --- a/ouroboros/observability.py +++ b/ouroboros/observability.py @@ -572,7 +572,7 @@ _PUBLISHED_CHILD_REF_FIELDS = frozenset( } ) _SOURCE_HANDLES_SUBDIR = "source_handles" -_TASK_SOURCE_MARKER = "FULL_RESULT_SOURCE_JSON=" +_TASK_SOURCE_MARKERS = ("FULL_RESULT_SOURCE_JSON=", "PRODUCER_RESULT_SOURCE_JSON=") _SERVICE_REF_TOOLS = frozenset({"service_logs", "stop_service"}) @@ -857,17 +857,17 @@ def _rewrite_task_source_markers( task_id: str, state: Dict[str, Any], ) -> str: - """Rewrite only Phase3B's explicit actor-source envelope inside tool text.""" + """Promote the host's full-view and clean-producer source envelopes.""" rewritten_lines: List[str] = [] for line in str(text).splitlines(keepends=True): body = line.rstrip("\r\n") - newline = line[len(body):] - if not body.startswith(_TASK_SOURCE_MARKER): + marker = next((prefix for prefix in _TASK_SOURCE_MARKERS if body.startswith(prefix)), None) + if marker is None: rewritten_lines.append(line) continue try: - ref = json.loads(body[len(_TASK_SOURCE_MARKER):]) + ref = json.loads(body[len(marker):]) except (TypeError, ValueError): rewritten_lines.append(line) continue @@ -878,14 +878,14 @@ def _rewrite_task_source_markers( parent_root, child_root, task_id, ref, state ) rewritten_lines.append( - _TASK_SOURCE_MARKER + marker + json.dumps( promoted, ensure_ascii=False, sort_keys=True, separators=(",", ":"), ) - + newline + + line[len(body):] ) return "".join(rewritten_lines) @@ -940,7 +940,7 @@ def _rewrite_child_ref_tree( _rewrite_child_ref_tree(item, parent_root, child_root, task_id, state) for item in value ] - if isinstance(value, str) and _TASK_SOURCE_MARKER in value: + if isinstance(value, str) and any(marker in value for marker in _TASK_SOURCE_MARKERS): return _rewrite_task_source_markers( value, parent_root, child_root, task_id, state ) diff --git a/ouroboros/size_ratchet_manifest.py b/ouroboros/size_ratchet_manifest.py index 06ed35ee9..38a075d3f 100644 --- a/ouroboros/size_ratchet_manifest.py +++ b/ouroboros/size_ratchet_manifest.py @@ -160,6 +160,7 @@ BAND_PATHS = { "ouroboros/tools/skill_preflight.py": "Entered the band by absorbing the upstream classic-script validator, widget entry/grammar findings and the preflight schema description beside the v7 typed _run_check rework they attach to.", "ouroboros/tools/skill_publish.py": "Entered the band from 952 lines: publish now writes the OuroborosHub publication receipt at pr_opened through the shared locked-update seam and maps the receipt from the validated serialized form (hubflow sprint, receipt-as-only-stored-fact design).", "ouroboros/tools/subagent_integration.py": "Existing native result integration owner handles source patches and complete file artifacts under one disposition and target authority.", + "ouroboros/tools/tool_result.py": "Typed result composition owns producer payload and dispatch annotations together while preserving the status and metadata contract.", "ouroboros/usage_compaction.py": "Entered the band from 971 lines: the C6 round-4 fixes homed here \u2014 dir-fd/O_NOFOLLOW anchoring of the archive writer and reader (a link planted after any path check cannot receive or serve monetary history) and the swap's last-instant snapshot re-proof inside the atomic replace \u2014 defenses that belong beside the compaction pass they defend.", "ouroboros/workspace_executor.py": None, "scripts/claudexor_platform_smoke.py": "The managed Claudexor platform smoke owns a multi-platform fixture, lifecycle receipt, and cleanup proof; keeping this runner in the documented band preserves the release gate without moving those checks into product runtime.", diff --git a/ouroboros/tools/extension_dispatch.py b/ouroboros/tools/extension_dispatch.py index 478748fa1..264bbe7b6 100644 --- a/ouroboros/tools/extension_dispatch.py +++ b/ouroboros/tools/extension_dispatch.py @@ -13,7 +13,8 @@ from ouroboros.tools.tool_context import ToolContext from ouroboros.tools.tool_result import ( ToolResult, ToolStatus, - _compose_execute_result, + _compose_execute_result_result, + _replace_tool_result, _structured_failure, ) @@ -98,11 +99,10 @@ def _dispatch_mcp_tool_result( return ToolResult(status="error", code="TOOL_ERROR", text=text) if not safety_msg: return result - text = _compose_execute_result(result.text, "", safety_msg) - meta = {**dict(result.meta), "safety_warning": True} - if result.code == "OK": - return ToolResult(status="ok", code="SAFETY_WARNING", text=text, meta=meta) - return ToolResult(status=result.status, code=result.code, text=text, meta=meta) + return _replace_tool_result( + _compose_execute_result_result(name, result, "", safety_msg), + meta_updates={"safety_warning": True}, + ) def _extension_result( @@ -138,21 +138,18 @@ def _extension_completion(result: str, safety_msg: str) -> ToolResult: structured check is the adapter's, so there is exactly one implementation of what a self-reported failure is.""" reported_failure = _structured_failure(result) + base = _extension_result( + "error" if reported_failure else "ok", + "TOOL_REPORTED_FAILURE" if reported_failure else "OK", + result, + dispatched=True, + ) if safety_msg: - # #447 H1: the warning TRAILS the payload — line 1 belongs to the - # extension, so a structured {"ok": false} answer (and any first-line - # marker) stays readable to every text-only consumer downstream. - text = f"{result}\n\n{safety_msg}" - return _extension_result( - "error" if reported_failure else "ok", - "TOOL_REPORTED_FAILURE" if reported_failure else "SAFETY_WARNING", - text, - safety_warning=True, - dispatched=True, + return _replace_tool_result( + _compose_execute_result_result("", base, "", safety_msg), + meta_updates={"safety_warning": True}, ) - if reported_failure: - return _extension_result("error", "TOOL_REPORTED_FAILURE", result, dispatched=True) - return _extension_result("ok", "OK", result, dispatched=True) + return base def _generation_digest_for(ext_tool: Dict[str, Any]) -> str: @@ -208,11 +205,9 @@ def _dispatch_extension_tool_result( or not digest ): return result - return ToolResult( - status=result.status, - code=result.code, - text=result.text, - meta={**dict(result.meta), "extension_generation": digest, + return _replace_tool_result( + result, + meta_updates={"extension_generation": digest, **({"content_hash": content_hash} if content_hash else {})}, ) diff --git a/ouroboros/tools/tool_result.py b/ouroboros/tools/tool_result.py index 70978da6b..0ed02e943 100644 --- a/ouroboros/tools/tool_result.py +++ b/ouroboros/tools/tool_result.py @@ -496,12 +496,18 @@ TOOL_CODE_SPECS: Mapping[str, ToolCodeSpec] = MappingProxyType( @dataclass(frozen=True) class ToolResult: - """Internal result; ``text`` remains the complete model-facing projection.""" + """Internal result; ``text`` includes notes, ``producer_text`` never does. + + Producer text is captured only when the host adds annotations, before any + composition. It is not bounded metadata and is never reconstructed from text. + """ status: ToolStatus code: str text: str meta: Mapping[str, Any] = field(default_factory=dict) + producer_text: str | None = None + host_annotations: tuple[str, ...] = () def __post_init__(self) -> None: if not isinstance(self.code, str) or not _CODE_RE.fullmatch(self.code): @@ -513,6 +519,12 @@ class ToolResult: raise ValueError(f"status {self.status!r} does not match {self.code} ({spec.status!r})") if not isinstance(self.text, str): raise TypeError("tool result text must be a string") + if self.producer_text is not None and not isinstance(self.producer_text, str): + raise TypeError("tool producer text must be a string or None") + if not isinstance(self.host_annotations, tuple) or any( + not isinstance(note, str) for note in self.host_annotations + ): + raise TypeError("host annotations must be a tuple of strings") raw_meta = dict(self.meta or {}) if any(not isinstance(key, str) for key in raw_meta): raise ValueError("tool result meta keys must be strings") @@ -566,6 +578,8 @@ def _replace_tool_result( code=selected_code, text=result.text if text is None else text, meta=meta, + producer_text=result.producer_text, + host_annotations=result.host_annotations, ) @@ -967,6 +981,14 @@ def _compose_execute_result_result( else LegacyTextResultAdapter.from_text(tool_name, base) ) text = _compose_execute_result(base_result.text, route_note, safety_msg) + notes = tuple(note for note in (route_note, safety_msg) if note) + source = { + "producer_text": ( + base_result.producer_text if base_result.producer_text is not None + else base_result.text if notes else None + ), + "host_annotations": base_result.host_annotations + notes, + } meta = dict(base_result.meta) if route_note: meta["route_note"] = True @@ -983,6 +1005,7 @@ def _compose_execute_result_result( code=base_result.code, text=text, meta=meta, + **source, ) if base_result.code == "OK": return ToolResult( @@ -990,6 +1013,7 @@ def _compose_execute_result_result( code="SAFETY_WARNING", text=text, meta=meta, + **source, ) meta["safety_warning"] = True return ToolResult( @@ -997,4 +1021,5 @@ def _compose_execute_result_result( code=base_result.code, text=text, meta=meta, + **source, ) diff --git a/tests/test_reference_book_budgets.py b/tests/test_reference_book_budgets.py index eb78edf20..12026f41b 100644 --- a/tests/test_reference_book_budgets.py +++ b/tests/test_reference_book_budgets.py @@ -119,7 +119,8 @@ CHAPTER_BYTE_BUDGETS: dict[str, int] = { # 7764 -> 8600 (#1195): the fresh selected-subject + immutable peer projection # execution check (`skill_peer_inventory.py`, `skill_conflicts.py`) replaces # whole-inventory hashing; the chapter had no description of that seam to swap out. - "docs/architecture/13-external-skills-layer.md": 8600, + # Dispatcher producer/annotation separation and its retained-source lifetime. + "docs/architecture/13-external-skills-layer.md": 9500, "docs/development/01-role-and-authority.md": 2437, "docs/development/02-naming-and-boundaries.md": 36372, # 22873 -> 23100: one new invariant (notifications ring for live events diff --git a/tests/test_registry_core.py b/tests/test_registry_core.py index 06728d6a9..710166a4e 100644 --- a/tests/test_registry_core.py +++ b/tests/test_registry_core.py @@ -305,6 +305,8 @@ def test_registry_uses_typed_required_root_not_note_or_tool_name(tmp_path, monke code="OK", text="OK\n\n⚠️ AUTO_ROUTED_TO_ACTIVE_WORKSPACE: benign additive note", meta={"route_note": True}, + producer_text="OK", + host_annotations=("⚠️ AUTO_ROUTED_TO_ACTIVE_WORKSPACE: benign additive note",), ) assert len(calls) == 1 diff --git a/tests/test_tool_producer_promotion.py b/tests/test_tool_producer_promotion.py new file mode 100644 index 000000000..d4b30596a --- /dev/null +++ b/tests/test_tool_producer_promotion.py @@ -0,0 +1,65 @@ +"""Clean tool sources survive when only a model-request closure is published.""" + +from __future__ import annotations + +import json +from pathlib import Path +from types import SimpleNamespace + +import pytest + +from ouroboros.artifacts import read_actor_source_bytes +from ouroboros.headless import copy_child_task_result, prepare_task_drive, prune_headless_task_drives +from ouroboros.loop_tool_execution import process_tool_results +from ouroboros.observability import persist_call, read_blob_ref +from ouroboros.task_results import STATUS_COMPLETED, write_task_result +from ouroboros.tools.core import _read_file +from ouroboros.tools.extension_dispatch import _extension_completion +from ouroboros.tools.tool_context import ToolContext + + +def _marker_ref(text: str, marker: str) -> dict: + return json.loads(next(line[len(marker):] for line in text.splitlines() if line.startswith(marker))) + + +@pytest.mark.parametrize("large", [False, True]) +def test_clean_source_in_model_request_survives_child_copyback_and_pruning(tmp_path, large): + parent = tmp_path / "canonical" + task_id = "producer-copyback" + child = prepare_task_drive(parent, task_id, "empty") + assert child is not None + repo = tmp_path / "repo" + repo.mkdir() + ctx = ToolContext(repo_dir=repo, drive_root=child, task_id=task_id) + payload = json.dumps({"ok": False, "text": "雪" * (20000 if large else 1)}, ensure_ascii=False) + warning = "⚠️ SAFETY_WARNING: the returned content still needs inspection." + typed = _extension_completion(payload, warning) + messages, trace = [], {"tool_calls": []} + process_tool_results( + [{"fn_name": "ext_fixture", "tool_call_id": "fixture-call", "result": typed.text, + "tool_result": typed, "is_error": True, "tool_args": {}, "args_for_log": {}}], + messages, trace, lambda *a, **kw: None, SimpleNamespace(_ctx=ctx), + ) + source = _marker_ref(messages[0]["content"], "PRODUCER_RESULT_SOURCE_JSON=") + request = persist_call(child, task_id=task_id, call_id="producer-request", call_type="llm_request", + payload={"messages": messages}) + # The actual collector retains call refs, not the in-memory tool result row. + # Publish only the model request to prove marker-based dependency custody. + write_task_result(child, task_id, STATUS_COMPLETED, result="done", artifact_status="ready", + trace_refs={"llm_call_refs": [{"request_ref": request["manifest_ref"]}]}) + copied = copy_child_task_result(parent, {"id": task_id, "drive_root": str(child)}) + assert copied is not None and copied["child_ref_promotion"]["status"] == "complete" + request_ref = copied["trace_refs"]["llm_call_refs"][0]["request_ref"] + manifest = json.loads(Path(request_ref["path"]).read_text(encoding="utf-8")) + promoted = read_blob_ref(parent, manifest["full_payload_ref"])["messages"][0]["content"] + assert warning in promoted + assert _marker_ref(promoted, "PRODUCER_RESULT_SOURCE_JSON=") == source + prune_headless_task_drives(parent, retention_days=0, now=4_000_000_000.0) + assert not child.exists() + assert read_actor_source_bytes(parent, task_id, source) == payload.encode("utf-8") + canonical_ctx = ToolContext(repo_dir=repo, drive_root=parent, task_id=task_id) + assert "雪" in _read_file(canonical_ctx, **source["read"]["arguments"]) + assert json.loads(read_actor_source_bytes(parent, task_id, source))["ok"] is False + if large: + full = _marker_ref(promoted, "FULL_RESULT_SOURCE_JSON=") + assert read_actor_source_bytes(parent, task_id, full).decode("utf-8") == typed.text diff --git a/tests/test_tool_producer_sources.py b/tests/test_tool_producer_sources.py new file mode 100644 index 000000000..944e9d87b --- /dev/null +++ b/tests/test_tool_producer_sources.py @@ -0,0 +1,131 @@ +"""Host notes stay visible without corrupting data retained for a program.""" + +from __future__ import annotations + +import json +from types import SimpleNamespace + +import pytest + +from ouroboros import artifacts +from ouroboros import loop_tool_execution as execution +from ouroboros.tools import extension_dispatch +from ouroboros.tools.core import _read_file +from ouroboros.tools.tool_context import ToolContext +from ouroboros.tools.tool_result import ToolResult, _compose_execute_result_result + + +WARNING = "⚠️ SAFETY_WARNING: inspect the returned data before acting." +PAYLOAD = json.dumps({"ok": False, "body": "Unicode: 雪\n\n---\nSAFETY_WARNING is data", "rows": [None, True]}) + + +@pytest.mark.parametrize("route", ["builtin", "extension", "mcp"]) +@pytest.mark.parametrize("warning", ["", WARNING]) +def test_dispatch_preserves_payload_and_failure_without_parsing_notes(monkeypatch, route, warning): + base = ToolResult(status="error", code="TOOL_REPORTED_FAILURE", text=PAYLOAD) + if route == "builtin": + result = _compose_execute_result_result("fixture", base, "", warning) + elif route == "extension": + result = extension_dispatch._extension_completion(PAYLOAD, warning) + else: + monkeypatch.setattr("ouroboros.safety.check_safety", lambda *a, **k: (True, warning)) + monkeypatch.setattr("ouroboros.mcp_client._call_mcp_tool_result", lambda *a, **k: base) + result = extension_dispatch._dispatch_mcp_tool_result(SimpleNamespace(), "mcp_demo", {}) + assert result.status == "error" and result.code == "TOOL_REPORTED_FAILURE" + assert result.text == PAYLOAD + ("\n\n" + warning if warning else "") + assert result.producer_text == (PAYLOAD if warning else None) + assert result.host_annotations == ((warning,) if warning else ()) + assert json.loads(result.producer_text or result.text)["ok"] is False + + +def test_nested_notes_and_generation_stamp_keep_original_payload(monkeypatch): + result = extension_dispatch._extension_completion(PAYLOAD, WARNING) + result = _compose_execute_result_result("fixture", result, "route note", "") + monkeypatch.setattr(extension_dispatch, "_dispatch_extension_tool_untagged", lambda *a: result) + stamped = extension_dispatch._dispatch_extension_tool_result( + SimpleNamespace(), "ext_demo", {"extension_generation": "generation", "content_hash": "hash"}, {}, + ) + assert stamped.producer_text == PAYLOAD + assert stamped.host_annotations == (WARNING, "route note") + assert stamped.meta["extension_generation"] == "generation" + assert stamped.text == PAYLOAD + "\n\n" + WARNING + "\n\nroute note" + + +@pytest.mark.parametrize("actor", ["presence", "project", "local_readonly_subagent"]) +@pytest.mark.parametrize("large", [False, True]) +def test_loop_retains_parseable_source_and_annotated_review_source(tmp_path, monkeypatch, actor, large): + repo = tmp_path / "repo" + repo.mkdir() + drive = tmp_path / actor / "data" + (drive / "logs").mkdir(parents=True) + ctx = ToolContext( + repo_dir=repo, drive_root=drive, task_id="producer-source", + budget_drive_root=tmp_path / "canonical", + workspace_root=tmp_path / "project" if actor != "presence" else None, + workspace_mode="external" if actor != "presence" else "", + task_depth=1 if actor == "local_readonly_subagent" else 0, + ) + ctx.task_metadata = {"task_id": ctx.task_id} + if actor == "presence": + ctx.task_metadata.update(_presence_turn=True, presence={"profile": {}}) + elif actor == "project": + ctx.task_metadata["project_id"] = "project-demo" + else: + ctx.task_constraint = {"mode": "local_readonly_subagent"} + ctx.task_metadata["parent_task_id"] = "parent" + payload = json.dumps({"ok": False, "text": "雪" * (20000 if large else 1)}, ensure_ascii=False) + result = extension_dispatch._extension_completion(payload, WARNING) + registry = SimpleNamespace(_ctx=ctx, CODE_TOOLS=frozenset(), execute_result=lambda *a: result) + captured = [] + monkeypatch.setattr(execution, "persist_call", lambda *a, **k: captured.append(k["payload"]) or {}) + row = execution._execute_single_tool( + registry, {"id": "call-json", "function": {"name": "ext_demo", "arguments": "{}"}}, + drive / "logs", ctx.task_id, + ) + messages, trace = [], {"tool_calls": []} + execution.process_tool_results([row], messages, trace, lambda *a, **k: None, registry) + recorded = trace["tool_calls"][0] + ref = recorded["producer_source_ref"] + exact = artifacts.read_actor_source_bytes(drive, ctx.task_id, ref) + assert exact == payload.encode("utf-8") + assert json.loads(exact)["ok"] is False + assert "雪" in _read_file(ctx, **ref["read"]["arguments"]) + assert WARNING in messages[0]["content"] + assert "PRODUCER_RESULT_SOURCE_JSON=" in messages[0]["content"] + assert recorded["is_error"] is True + assert captured[0]["producer_result"] == payload + assert captured[0]["result"] == result.text + assert captured[0]["host_annotations"] == [WARNING] + if large: + full_ref = recorded["result_source_ref"] + assert artifacts.read_actor_source_bytes(drive, ctx.task_id, full_ref).decode("utf-8") == result.text + assert full_ref != ref + + +def test_clean_small_result_keeps_existing_projection(tmp_path): + ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="clean-source") + result = extension_dispatch._extension_completion(PAYLOAD, "") + messages, trace = [], {"tool_calls": []} + execution.process_tool_results( + [{"fn_name": "ext_demo", "tool_call_id": "clean", "result": PAYLOAD, + "tool_result": result, "is_error": True, "tool_args": {}, "args_for_log": {}}], + messages, trace, lambda *a, **k: None, SimpleNamespace(_ctx=ctx), + ) + assert messages[0]["content"] == PAYLOAD + assert "producer_source_ref" not in trace["tool_calls"][0] + + +def test_missing_source_is_disclosed_without_hiding_warning_or_failure(tmp_path, monkeypatch): + monkeypatch.setattr(artifacts, "store_actor_source_bytes", lambda *a, **k: (_ for _ in ()).throw(OSError("disk full"))) + ctx = ToolContext(repo_dir=tmp_path, drive_root=tmp_path, task_id="source-failure") + result = extension_dispatch._extension_completion(PAYLOAD, WARNING) + messages, trace = [], {"tool_calls": []} + execution.process_tool_results( + [{"fn_name": "ext_demo", "tool_call_id": "failed", "result": result.text, + "tool_result": result, "is_error": True, "tool_args": {}, "args_for_log": {}}], + messages, trace, lambda *a, **k: None, SimpleNamespace(_ctx=ctx), + ) + assert result.text in messages[0]["content"] + assert "PRODUCER_RESULT_SOURCE_UNAVAILABLE=true" in messages[0]["content"] + assert trace["tool_calls"][0]["is_error"] is True + assert trace["tool_calls"][0]["producer_source_ref"] == {}