fix: retain producer tool payloads beside dispatcher annotations

Preserve exact machine-readable sources through task artifact promotion while keeping warning and outcome projections intact.
This commit is contained in:
Ouroboros 2026-09-23 05:20:52 +03:00
parent 558b1ea2d1
commit a9e65e974f
10 changed files with 286 additions and 36 deletions

View file

@ -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/<name>/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.

View file

@ -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 "")

View file

@ -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
)

View file

@ -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.",

View file

@ -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 {})},
)

View file

@ -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,
)

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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"] == {}