diff --git a/src/agent_runtime/effect_adapters.py b/src/agent_runtime/effect_adapters.py index 98b062962..03083e47f 100644 --- a/src/agent_runtime/effect_adapters.py +++ b/src/agent_runtime/effect_adapters.py @@ -255,18 +255,43 @@ def begin_effect(journal: Any, action: Any) -> DispatchCapture: return capture -def _execution(result: Any, facts: ProducerFacts) -> ExecutionOutcome: +def _server_producer(capture: DispatchCapture) -> bool: + """The backend was a server-owned producer bound by Wave 3 admission. + + Only such producers build their result dictionaries from server state. An + unbound dynamic/registry tool returns whatever it likes, so its keys carry + no lifecycle meaning. The MCP bridge builds only stdout/stderr/exit_code. + """ + return any(bound is not None for bound in (capture.filesystem, capture.owned, capture.process, capture.browser)) + + +def _facts(result: Any, capture: DispatchCapture) -> ProducerFacts: + """Typed producer facts, scoped to what the captured producer can attest.""" + facts = producer_facts(result) + if not _server_producer(capture): + # Reported success or failure is all an untrusted result can say. + facts = ProducerFacts(exit_code=facts.exit_code) + if capture.backend is not None and capture.claim is not None and capture.claim.external: + facts = ProducerFacts(**{**facts.to_dict(), "external": True, "remote_acknowledged": facts.exit_code == 0}) + return facts + + +def _execution(result: Any, facts: ProducerFacts, capture: DispatchCapture) -> ExecutionOutcome: if not isinstance(result, dict): return ExecutionOutcome.INTERRUPTED if facts.timed_out: return ExecutionOutcome.TIMED_OUT - # Only launch-shaped results mean this operation's own work continues: - # the native detached launch, or a bridge's explicit detachment. A listing - # that merely reports some other thing as "running" is not. - if isinstance(result.get("bg_job_id"), str) and facts.exit_code == 0: - return ExecutionOutcome.RUNNING - if result.get("detached") is True: - return ExecutionOutcome.RUNNING + # Only a server process producer can say this operation's own work + # continues: the native detached launch of an exact Wave 3 launch + # reservation, or the host bridge's server-set detachment. Lifecycle keys + # from any other producer (or a listing reporting something else as + # running) do not. + process = capture.process + if process is not None: + if process.launch is not None and isinstance(result.get("bg_job_id"), str) and facts.exit_code == 0: + return ExecutionOutcome.RUNNING + if result.get("detached") is True: + return ExecutionOutcome.RUNNING denied = bool(result.get("blocked") or result.get("approval_required") or facts.failure_kind.endswith("_denied")) if facts.exit_code == 0 and not result.get("error") and not denied: @@ -274,17 +299,20 @@ def _execution(result: Any, facts: ProducerFacts) -> ExecutionOutcome: return ExecutionOutcome.FAILED -def _cleanup(result: Any, facts: ProducerFacts) -> CleanupState: +def _cleanup(result: Any, facts: ProducerFacts, capture: DispatchCapture) -> CleanupState: if not isinstance(result, dict): return CleanupState.UNKNOWN + if facts.external: + # External execution reports no locally observed teardown. + return CleanupState.UNKNOWN + if capture.process is None: + return CleanupState.NOT_APPLICABLE + # Teardown is attested only by the native process/containment producer. if facts.failure_kind == "process_teardown_failed": return CleanupState.FAILED teardown = result.get("teardown") if isinstance(teardown, dict) and type(teardown.get("dead")) is bool: return CleanupState.VERIFIED if teardown["dead"] else CleanupState.FAILED - if facts.external: - # External execution reports no locally observed teardown. - return CleanupState.UNKNOWN return CleanupState.NOT_APPLICABLE @@ -300,11 +328,8 @@ def settle_effect(journal: Any, action: Any, capture: DispatchCapture | None, *, else ExecutionOutcome.INTERRUPTED) facts, cleanup = ProducerFacts(), CleanupState.UNKNOWN else: - facts = producer_facts(result) - if capture.backend is not None and capture.claim.external: - facts = ProducerFacts(**{**facts.to_dict(), "external": True, - "remote_acknowledged": facts.exit_code == 0}) - execution, cleanup = _execution(result, facts), _cleanup(result, facts) + facts = _facts(result, capture) + execution, cleanup = _execution(result, facts, capture), _cleanup(result, facts, capture) log.outcome(effect_id=capture.claim.effect_id, execution=execution, impact=Impact.POSSIBLE, facts=facts, cleanup=cleanup, execution_id=action.execution_id or "") if (execution is ExecutionOutcome.REPORTED_SUCCESS and capture.process is not None @@ -384,8 +409,9 @@ def _observations(capture: DispatchCapture, action: Any, result: dict) -> list[d mechanism=ObservationMechanism.OWNED_RECORD_READ, coverage=Coverage.PARTIAL, exists=True, **base) for i, r in enumerate(capture.owned.resources) if r.record_id != "*"] if capture.process is not None and capture.process.launch is None: + from src.agent_runtime.process_resources import JOB_TOOL job = result.get("job") - if isinstance(job, dict) and len(capture.process.jobs) == 1: + if isinstance(job, dict) and len(capture.process.jobs) == 1 and capture.process.operation.tool == JOB_TOOL: return [dict(observation_id=action.action_id + ":observation", resource=resource_ref(capture.process.jobs[0], "job"), mechanism=ObservationMechanism.JOB_STATE, coverage=Coverage.PARTIAL, exists=True, **base)] @@ -399,8 +425,10 @@ def _settle_background(log: Any, capture: DispatchCapture, result: dict) -> None validated by ``job_from_record`` at admission. Job completion is execution evidence for that claim; it verifies no postcondition. """ + from src.agent_runtime.process_resources import JOB_TOOL job_facts = result.get("job") - if not isinstance(job_facts, dict) or len(capture.process.jobs) != 1: + if (not isinstance(job_facts, dict) or len(capture.process.jobs) != 1 + or capture.process.operation.tool != JOB_TOOL): return settle_background_job(capture.process.jobs[0], job_facts, log=log)