diff --git a/src/agent_runtime/completion.py b/src/agent_runtime/completion.py index 6bf9ef374..f667b6937 100644 --- a/src/agent_runtime/completion.py +++ b/src/agent_runtime/completion.py @@ -155,10 +155,16 @@ def completion_answer(text: str, ledger: EvidenceLedger, decision: CompletionDec def _disclose(answer: str, ledger: EvidenceLedger) -> str: """Append the server's facts for unverified external effects.""" + disclosure = _disclosure(answer, ledger) + return answer.rstrip() + disclosure if disclosure else answer + + +def _disclosure(answer: str, ledger: EvidenceLedger) -> str: + """Build the complete server-owned disclosure independently of prose length.""" summary = ' '.join(ledger.effect_disclosures()) if not summary: - return answer - return (answer.rstrip() + '\n\n' + summary) if answer.strip() else summary + return '' + return ('\n\n' + summary) if answer.strip() else summary def _completion_answer(text: str, ledger: EvidenceLedger, decision: CompletionDecision) -> tuple[str, str]: @@ -351,7 +357,7 @@ def with_completion_gate(func): # When the only change is the server's effect disclosure, the # model's answer events are released unchanged and the disclosure # follows them, so no earlier-round text is dropped. - disclosure = safe_answer[len(filtered_answer):] if safe_answer != filtered_answer else '' + disclosure = _disclosure(filtered_answer, ledger) disclosure_only = bool(disclosure) and not (presentation_replaced or reason or unsafe_draft or filtered_answer != answer) replaced_answer = not disclosure_only and bool( @@ -392,7 +398,7 @@ def with_completion_gate(func): elif disclosure_only and metadata.get('round_texts') and isinstance(metadata['round_texts'], list) \ and isinstance(metadata['round_texts'][-1], str): # Reload renders round_texts: keep the disclosure with them. - metadata['round_texts'] = [*metadata['round_texts'][:-1], metadata['round_texts'][-1] + disclosure] + metadata['round_texts'] = [*metadata['round_texts'][:-1], metadata['round_texts'][-1].rstrip() + disclosure] if provider_error and isinstance(metadata.get('round_texts'), list): # Failed rounds stay as per-round diagnostics, but they are # rendered again on reload. Apply the same statement filter diff --git a/src/agent_runtime/effect_adapters.py b/src/agent_runtime/effect_adapters.py index 03083e47f..3c66ae4f8 100644 --- a/src/agent_runtime/effect_adapters.py +++ b/src/agent_runtime/effect_adapters.py @@ -98,7 +98,7 @@ def _pre_state_text(resource: Any, *, newline: str | None) -> str | None: replaced or undecodable file yields None: a truncated read must never stand in for the whole pre-state. """ - data = _read_whole(resource, _PRE_STATE_LIMIT) + data = _read_whole(resource, _PRE_STATE_LIMIT).data if data is None or len(data) > _PRE_STATE_LIMIT: return None try: @@ -336,33 +336,58 @@ def settle_effect(journal: Any, action: Any, capture: DispatchCapture | None, *, and capture.process.launch is None): _settle_background(log, capture, result) # e.g. an exact kill return - if error is not None or not isinstance(result, dict) or result.get("exit_code") != 0 or result.get("error"): + if error is not None or not isinstance(result, dict): + return + successful = result.get("exit_code") == 0 and not result.get("error") + # A missing-file read reports failure, but can independently establish + # absence. No other failed read is eligible for an observation. + absent_read = (capture.filesystem is not None and capture.filesystem.operation.tool == "read_file" + and capture.filesystem.bindings[0].resource.identity is None) + if not successful and not absent_read: return for fields in _observations(capture, action, result): - log.observe(**fields) + if successful or fields.get("exists") is False: + log.observe(**fields) if capture.process is not None and capture.process.launch is None: _settle_background(log, capture, result) # -- observations ------------------------------------------------------------ -def _read_whole(resource: Any, limit: int) -> bytes | None: - """Re-read the exact admitted source binding; None if it is not stable.""" +@dataclass(frozen=True) +class _WholeFileRead: + data: bytes | None = None + known_absent: bool = False + + +def _read_whole(resource: Any, limit: int) -> _WholeFileRead: + """Read a stable binding, distinguish validated ENOENT from uncertainty. + + Only a binding admitted as absent can prove absence. Disappearance of an + existing identity, replacement, or any validation/access failure is unknown. + """ flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0) | getattr(os, "O_CLOEXEC", 0) try: resource.validate() + if resource.identity is None: + try: + os.lstat(resource.path) + except FileNotFoundError: + resource.validate() + return _WholeFileRead(known_absent=True) + return _WholeFileRead() descriptor = os.open(resource.path, flags) with os.fdopen(descriptor, "rb") as stream: info = os.fstat(stream.fileno()) identity = resource.identity if (not stat.S_ISREG(info.st_mode) or identity is None or (info.st_dev, info.st_ino) != (identity.device, identity.inode)): - return None + return _WholeFileRead() data = stream.read(limit + 1) resource.validate() - except (OSError, ValueError): - return None - return data + except (OSError, ValueError, RuntimeError): + return _WholeFileRead() + return _WholeFileRead(data=data) def _file_observation(capture: DispatchCapture, action: Any) -> dict[str, Any] | None: @@ -373,17 +398,21 @@ def _file_observation(capture: DispatchCapture, action: Any) -> dict[str, Any] | args = json.loads(bound.execution_input) partial = bool(args.get("offset") or args.get("limit")) or ( os.path.splitext(resource.path)[1].lower() in producer._STRUCTURED_DOCUMENT_SUFFIXES) - data = _read_whole(resource, producer.MAX_READ_CHARS * 4) - if data is None: + read = _read_whole(resource, producer.MAX_READ_CHARS * 4) + data = read.data + if data is None and not read.known_absent: return None - if len(data) > producer.MAX_READ_CHARS * 4 or len(data.decode("utf-8", errors="replace")) > producer.MAX_READ_CHARS: + if read.known_absent: + partial = False # ENOENT establishes absence of the whole bound path. + elif len(data) > producer.MAX_READ_CHARS * 4 or len(data.decode("utf-8", errors="replace")) > producer.MAX_READ_CHARS: partial = True # the producer truncated what it read complete = not partial return dict(observation_id=action.action_id + ":observation", resource=resource_ref(resource, binding.role), mechanism=ObservationMechanism.FILESYSTEM_READ, coverage=Coverage.COMPLETE if complete else Coverage.PARTIAL, source_action_id=action.action_id, source_execution_id=action.execution_id or "", - exists=True, content_sha256=hashlib.sha256(data).hexdigest() if complete else "") + exists=not read.known_absent, + content_sha256=hashlib.sha256(data).hexdigest() if complete and data is not None else "") def _observations(capture: DispatchCapture, action: Any, result: dict) -> list[dict[str, Any]]: diff --git a/src/agent_runtime/resource_binding.py b/src/agent_runtime/resource_binding.py index 73658be55..0a1c89586 100644 --- a/src/agent_runtime/resource_binding.py +++ b/src/agent_runtime/resource_binding.py @@ -176,7 +176,7 @@ def resolve_filesystem_operation(operation, *, roots, workspace="", request_id=" raise ValueError("Search root is unresolved") args["path"] = bind(selector, "search_root" if search else "source" if tool == "read_file" else "destination" if tool == "write_file" else "target", - missing=tool == "write_file") + missing=tool in {"write_file", "read_file"}) execution_input = json.dumps(args, sort_keys=True, allow_nan=False) bound = BoundFilesystemOperation(operation, execution_input, tuple(bindings), request_id) bound.validate() diff --git a/src/tool_execution.py b/src/tool_execution.py index e495dd1fb..e5f77af59 100644 --- a/src/tool_execution.py +++ b/src/tool_execution.py @@ -1355,7 +1355,7 @@ from src.agent_runtime.process_resources import ( active_process_operation, bind_process_operation, needs_process_binding, resolve_process_operation, ) from src.browser_identity import ( - native_browser, parse_operation as parse_browser_operation, SESSION_ACTIONS, + native_browser, parse_operation as parse_browser_operation, SESSION_ACTIONS, PAGE_FAILURE, page_unavailable, resolve_browser_operation, bind_browser_operation, revalidate_browser_operation, ) @@ -1642,6 +1642,8 @@ async def execute_tool_block( ) return output except ResourceIdentityError as error: + if native_browser(operation, backend_operation.resource) and str(error) == PAGE_FAILURE: + return f"{transport}: UNSUPPORTED", page_unavailable() return f"{transport}: BLOCKED", { "error": str(error), "exit_code": 1, "blocked": True, "failure_kind": "resource_identity_denied", diff --git a/tests/test_browser_resource_identity.py b/tests/test_browser_resource_identity.py index 7ec9de375..f17452707 100644 --- a/tests/test_browser_resource_identity.py +++ b/tests/test_browser_resource_identity.py @@ -98,6 +98,42 @@ async def test_disabled_page_operations_never_observe_select_or_execute(producer assert old.target_id != producer.target +@pytest.mark.parametrize("reason,expected", [ + (browser.PAGE_FAILURE, browser.PAGE_FAILURE), + ("Unrelated resource identity changed", "resource_identity_denied"), + (browser.PAGE_FAILURE + ": arbitrary detail", "resource_identity_denied"), +]) +async def test_dispatch_boundary_preserves_only_native_browser_page_failure(producer, tmp_path, monkeypatch, reason, expected): + from src import tool_execution + from src.agent_runtime.effect_log import EffectLog + from src.agent_runtime.journal import ActionJournal, bind_journal + + await observed(producer) + producer.calls.clear() + producer.cdp_calls.clear() + journal = ActionJournal() + journal.effects = EffectLog(journal.run_id, directory=tmp_path / "fx") + attempts = [] + + async def unsupported_operation(*args, **kwargs): + attempts.append(1) + if reason == browser.PAGE_FAILURE: + # The legacy server page helper raises the reserved identity error. + await PrivateBrowserTool()._capture_post_click_state() + raise ResourceIdentityError(reason) + + monkeypatch.setattr(tool_execution, "_execute_tool_block_impl", unsupported_operation) + monkeypatch.setattr(tool_execution, "mark_dispatch", lambda: pytest.fail("Unsupported page operation dispatched")) + with bind_journal(journal): + _, result = await dispatch(authority(), "private_browser", '{"action":"session_info"}') + assert result["failure_kind"] == expected + if expected == browser.PAGE_FAILURE: + assert result["executed"] is False and result["retryable"] is False + assert journal.actions[0].execution_id is None + assert journal.effects.history().claims == () + assert attempts == [1] + + @pytest.mark.parametrize("args", [{"action": "batch", "commands": [["click", "@e1"]]}, {"action": "tab"}, {"action": "window"}, {"action": "frame"}, {"action": "connect"}, {"action": "click", "target": "--new-tab"}, {"action": "evaluate", "--cdp": "endpoint"}, diff --git a/tests/test_effect_resource_bindings.py b/tests/test_effect_resource_bindings.py index 834edcd09..b29cc8146 100644 --- a/tests/test_effect_resource_bindings.py +++ b/tests/test_effect_resource_bindings.py @@ -229,6 +229,65 @@ def test_patch_obligations_follow_exact_bindings(run, ws): assert {o.predicate for o in claim.obligations} == {fx.Predicate.CONTENT_SHA256, fx.Predicate.ABSENT} +def test_deleted_file_read_emits_known_absence(run, ws): + (ws / "old.txt").write_text("old\n") + _, result = run("apply_patch", {"patch_text": "*** Begin Patch\n*** Delete File: old.txt\n*** End Patch"}) + assert result["exit_code"] == 0 + _, result = run("read_file", {"path": "old.txt"}) + assert result["exit_code"] == 1 # The producer still reports a missing file. + history = run.journal.effects.history() + observation, = history.observations + assert observation.exists is False and observation.content_sha256 == "" + assert observation.coverage is fx.Coverage.COMPLETE + assert fx.predicate_holds(history.claims[0].obligations[0], observation) is True + assert verdicts(run.journal) == [fx.EffectVerdict.VERIFIED] + + +@pytest.mark.parametrize("failure", ["identity_mismatch", "replaced_path", "post_probe_replacement", "permission", "validation"]) +def test_indeterminate_deleted_file_read_cannot_prove_absence(run, ws, monkeypatch, failure): + from src.agent_runtime import effect_adapters as adapters + from src.agent_runtime.resources import ResourceIdentityError + + target = ws / "old.txt" + target.write_text("old\n") + run("apply_patch", {"patch_text": "*** Begin Patch\n*** Delete File: old.txt\n*** End Patch"}) + if failure == "identity_mismatch": + target.write_text("replacement\n") + original = adapters._read_whole + + def indeterminate(resource, limit): + if failure == "identity_mismatch": + target.unlink() # An existing binding disappearing is an identity failure. + elif failure == "replaced_path": + target.write_text("replacement\n") + elif failure == "post_probe_replacement": + validate = type(resource).validate + calls = [] + + def replace_after_probe(self): + calls.append(1) + if len(calls) == 2: + target.write_text("appeared after ENOENT\n") + return validate(self) + + monkeypatch.setattr(type(resource), "validate", replace_after_probe) + elif failure == "permission": + def denied(path): + raise PermissionError("access denied") + monkeypatch.setattr(adapters.os, "lstat", denied) + else: + def invalid(self): + raise ResourceIdentityError("unresolved binding") + monkeypatch.setattr(type(resource), "validate", invalid) + return original(resource, limit) + + monkeypatch.setattr(adapters, "_read_whole", indeterminate) + run("read_file", {"path": "old.txt"}) + history = run.journal.effects.history() + assert history.observations == () + assert verdicts(run.journal) == [fx.EffectVerdict.UNVERIFIED] + + def test_listing_is_partial_and_does_not_verify_content(run): run("write_file", {"path": "a.txt", "content": "hello\n"}) run("ls", {"path": "."}) @@ -243,7 +302,11 @@ def test_listing_is_partial_and_does_not_verify_content(run): {"action": "snapshot", "page": "t1"}, {"action": "evaluate", "page": "t1", "script": "1"}, ]) -def test_browser_page_operations_stay_fail_closed_with_effects(run, args): +def test_browser_page_operations_stay_fail_closed_with_effects(run, args, monkeypatch): + async def unexpected_dispatch(*args, **kwargs): + pytest.fail("Unsupported page operation reached execution") + + monkeypatch.setattr(tool_execution, "_execute_tool_block_impl", unexpected_dispatch) description, result = run("private_browser", args) assert "UNSUPPORTED" in description assert result["failure_kind"] == "browser_page_authority_unavailable" and result["executed"] is False diff --git a/tests/test_effect_verification_adapters.py b/tests/test_effect_verification_adapters.py index 8af800ef9..db0af801f 100644 --- a/tests/test_effect_verification_adapters.py +++ b/tests/test_effect_verification_adapters.py @@ -436,6 +436,35 @@ DISCLOSURE = ("External operation mcp__server__send_email reported success; any "not independently verified.") +@pytest.mark.parametrize("answer", ["Here is the draft. \n\n", " \n\n"]) +@pytest.mark.parametrize("final", [False, True]) +async def test_streaming_external_disclosure_survives_trailing_whitespace(store, monkeypatch, answer, final): + from src.agent_runtime.completion import completion_answer, with_completion_gate + from src.agent_runtime.journal import current_journal + + expected = [] + + @with_completion_gate + async def stream(messages): + journal = current_journal() + journal.effects = EffectLog(journal.run_id, directory=store) + remote_act(journal, monkeypatch, result={"stdout": "ok", "stderr": "", "exit_code": 0}) + ledger = _ledger(journal, CompletionRequirements()) + expected.append(completion_answer(answer, ledger, ledger.evaluate())[0]) + yield "data: " + json.dumps({"type": "final_response", "content": answer} if final else {"delta": answer}) + "\n\n" + yield "data: " + json.dumps({"type": "metrics", "data": {"round_texts": [answer]}}) + "\n\n" + + events = [json.loads(chunk[6:]) async for chunk in stream([])] + disclosure = next(event["delta"] for event in events if event.get("delta") != answer and "delta" in event) + assert disclosure == ("\n\n" if answer.strip() else "") + DISCLOSURE + visible = "".join(event.get("delta", event.get("content", "")) for event in events) + assert visible.count(DISCLOSURE) == 1 + metrics = next(event["data"] for event in events if event.get("type") == "metrics") + assert metrics["round_texts"] == expected + assert expected[0].count(DISCLOSURE) == 1 + assert metrics["completion_gate"]["answer_replaced"] is False + + def test_reported_external_mutation_cannot_complete_as_satisfied(tmp_path, store, monkeypatch): from src.agent_evidence import EXTERNAL_EFFECT_UNVERIFIED from src.agent_runtime.completion import completion_answer diff --git a/tests/test_resource_identity.py b/tests/test_resource_identity.py index a608486a2..230e9d8a7 100644 --- a/tests/test_resource_identity.py +++ b/tests/test_resource_identity.py @@ -615,10 +615,24 @@ async def test_approved_resource_cannot_migrate_to_another_request(tmp_path): assert result["failure_kind"] == "resource_identity_denied" -async def test_missing_approval_resource_snapshot_cannot_be_reconstructed(tmp_path): +async def test_missing_approval_resource_snapshot_cannot_be_reconstructed(tmp_path, monkeypatch): + grant = authority(tmp_path, "read_file") + def unavailable(*args, **kwargs): + raise PermissionError("Cannot establish the proposal's resource identity") + + with monkeypatch.context() as patch: + patch.setattr("src.agent_runtime.resource_binding.resolve_filesystem_operation", unavailable) + exact, security = approval(grant, "read_file", "missing") + assert exact.pending.resource_operation is None + (tmp_path / "missing").write_text("appeared after proposal") + _, result = await dispatch(grant, "read_file", "missing", exact_approval=exact, security_context=security) + assert result["failure_kind"] == "resource_identity_denied" + + +async def test_approved_absent_read_cannot_bind_a_file_that_appeared(tmp_path): grant = authority(tmp_path, "read_file") exact, security = approval(grant, "read_file", "missing") - assert exact.pending.resource_operation is None + assert exact.pending.resource_operation.bindings[0].resource.identity is None (tmp_path / "missing").write_text("appeared after proposal") _, result = await dispatch(grant, "read_file", "missing", exact_approval=exact, security_context=security) assert result["failure_kind"] == "resource_identity_denied"