diff --git a/docs/architecture/12-host-service-companions-and-chat-ids.md b/docs/architecture/12-host-service-companions-and-chat-ids.md index 6a7c08e72..a21783b7f 100644 --- a/docs/architecture/12-host-service-companions-and-chat-ids.md +++ b/docs/architecture/12-host-service-companions-and-chat-ids.md @@ -14,7 +14,7 @@ Operation correlation: a named injected message has `operation_ref=: dict[str raise PresenceTurnError("presence_result_unreadable", "source_event_id", turn_ref=task_id) from exc +def _terminal_refusal(stored: Mapping[str, Any]) -> str: + """The typed refusal a durable Presence row demands before its event may be acknowledged. + + The durable terminal cause decides, never a draft or the in-memory envelope. A confirmed + resource refusal keeps the event with the transport (``presence_resources_unavailable``). + A row without a canonical terminal (RUNNING/INTERRUPTED, a reconciled placeholder, a + terminal write that failed after the start barrier) and any host-authored infrastructure + terminal (provider death, an unknown outcome behind a quota refusal, overflow) is an + attempt whose external effect is unproven (``presence_attempt_outcome_unknown``): the + diagnostic reason word may say ``provider_unavailable``, but no model answered this + event, so completed/silent would let the adapter drop it. Empty: the row may answer. + """ + if str(stored.get("reason_code") or "") == "resource_refusal_no_resend": + return "presence_resources_unavailable" + if str(stored.get("status") or "") not in {STATUS_COMPLETED, STATUS_FAILED} or is_reconciled_presence_placeholder(stored): + return "presence_attempt_outcome_unknown" + axes = stored.get("outcome_axes") if isinstance(stored.get("outcome_axes"), dict) else {} + execution = axes.get("execution") if isinstance(axes.get("execution"), dict) else {} + if str(execution.get("status") or stored.get("execution_status") or "") == "infra_failed": + return "presence_attempt_outcome_unknown" + return "" + + def _cached_result(drive_root: Path, task_id: str, identity: str = "") -> PresenceTurnResult | None: if identity: task_id = _retry_target(drive_root, task_id, identity) stored = _stored_turn(drive_root, task_id, identity) - if str(stored.get("reason_code") or "") == "resource_refusal_no_resend": + refusal = _terminal_refusal(stored) + if refusal == "presence_resources_unavailable": metadata = stored.get("metadata") if isinstance(stored.get("metadata"), dict) else {} if identity and isinstance(metadata.get("presence_retry_proof"), dict): proof = metadata["presence_retry_proof"] @@ -545,8 +569,16 @@ def _cached_result(drive_root: Path, task_id: str, identity: str = "") -> Presen _notify_unresolved_turn(drive_root, task_id) raise PresenceTurnError("presence_resources_unavailable", "source_event_id", turn_ref=task_id, work_ref=str(metadata.get("presence_work_ref") or "")) - if str(stored.get("status") or "") not in {"completed", "failed"} or is_reconciled_presence_placeholder(stored): + if refusal and (str(stored.get("status") or "") not in {STATUS_COMPLETED, STATUS_FAILED} + or is_reconciled_presence_placeholder(stored)): return None # a host-lost turn is not a result; the later admission guard refuses regeneration + if refusal: + # A failed infrastructure terminal never earns a retry certificate: unknown is not + # not_started. Retain the event and any already scheduled work; ask the owner once. + metadata = stored.get("metadata") if isinstance(stored.get("metadata"), dict) else {} + _notify_unresolved_turn(drive_root, task_id) + raise PresenceTurnError(refusal, "source_event_id", turn_ref=task_id, + work_ref=str(metadata.get("presence_work_ref") or "")) return presence_result_from_stored(stored, task_id) @@ -1058,16 +1090,20 @@ def run_presence_turn( except (OSError, ValueError) as exc: raise PresenceTurnError("presence_start_unwritable", "source_event_id", turn_ref=task_id) from exc events = agent.handle_task(task) - # The model can have a draft reply before every allowed route refuses quota. - # The durable terminal cause, not that draft or the presence_result envelope, - # determines whether the transport may acknowledge the original event. + # The model can have a draft reply before every allowed route refuses quota, and the + # terminal write can fail after the start barrier. The durable terminal cause, not that + # draft or the presence_result envelope, determines whether the transport may + # acknowledge the original event; a row still RUNNING here is an unproven effect. terminal = _stored_turn(Path(drive_root), task_id, identity) - if str(terminal.get("reason_code") or "") == "resource_refusal_no_resend": + row = next((item for item in events if item.get("type") == "presence_result"), None) + refusal = _terminal_refusal(terminal) + if refusal: _notify_unresolved_turn(Path(drive_root), task_id) metadata = terminal.get("metadata") if isinstance(terminal.get("metadata"), dict) else {} - raise PresenceTurnError("presence_resources_unavailable", "source_event_id", turn_ref=task_id, - work_ref=str(metadata.get("presence_work_ref") or "")) - row = next((item for item in events if item.get("type") == "presence_result"), None) + # Scheduled work survives the refusal: the durable ref when the terminal landed, else the + # host-built handoff fact of this execution (the terminal write itself may have failed). + work_ref = str(metadata.get("presence_work_ref") or (row or {}).get("work_ref") or "") + raise PresenceTurnError(refusal, "source_event_id", turn_ref=task_id, work_ref=work_ref) if not isinstance(row, dict): raise PresenceTurnError("presence_result_missing", "presence_result") result = PresenceTurnResult( diff --git a/tests/test_host_service_api.py b/tests/test_host_service_api.py index 24e03e8c1..5be3d734d 100644 --- a/tests/test_host_service_api.py +++ b/tests/test_host_service_api.py @@ -3,6 +3,8 @@ from types import SimpleNamespace from starlette.testclient import TestClient +from ouroboros.task_results import write_task_result + from ouroboros.gateway.host_service import AUTH_TOKEN_FILENAME, create_host_service_app from ouroboros.event_bus import CHAT_OUTBOUND, publish_event from ouroboros.skill_loader import compute_content_hash, save_enabled, save_review_state, save_skill_grants, SkillReviewState @@ -298,6 +300,7 @@ def test_presence_turn_attachment_refusal_returns_complete_typed_manifest( class Agent: def handle_task(self, task): agent_calls.append(task) + write_task_result(tmp_path, task["id"], "completed", metadata=task["metadata"], result="ok") return [{"type": "presence_result", "outcome": "message", "text": "ok"}] def run_real_presence(**kwargs): @@ -369,6 +372,7 @@ def test_presence_turn_host_passes_attachment_limit_to_canonical_staging_owner( class Agent: def handle_task(self, task): agent_calls.append(task) + write_task_result(tmp_path, task["id"], "completed", metadata=task["metadata"], result="ok") return [{"type": "presence_result", "outcome": "message", "text": "ok"}] def run_real_presence(**kwargs): @@ -441,6 +445,7 @@ def test_presence_turn_host_passes_internal_missing_and_directory_to_staging_own class Agent: def handle_task(self, task): agent_calls.append(task) + write_task_result(tmp_path, task["id"], "completed", metadata=task["metadata"], result="ok") return [{"type": "presence_result", "outcome": "message", "text": "ok"}] def run_real_presence(**kwargs): diff --git a/tests/test_presence_cognitive_baseline.py b/tests/test_presence_cognitive_baseline.py index 4f973528a..db65b946f 100644 --- a/tests/test_presence_cognitive_baseline.py +++ b/tests/test_presence_cognitive_baseline.py @@ -181,6 +181,9 @@ def test_admitted_external_turn_writes_global_knowledge_and_nothing_else(tmp_pat ("run_command", {"cmd": ["true"]}), ) } + from ouroboros.task_results import write_task_result + + write_task_result(data, task["id"], "completed", metadata=task["metadata"], result="Noted.") return [{"type": "presence_result", "outcome": "message", "text": "Noted.", "work_ref": ""}] result = run_presence_turn( diff --git a/tests/test_presence_delivery_mode.py b/tests/test_presence_delivery_mode.py index 022a984b1..a851dea91 100644 --- a/tests/test_presence_delivery_mode.py +++ b/tests/test_presence_delivery_mode.py @@ -7,6 +7,7 @@ from dataclasses import replace import pytest from ouroboros.presence_runner import PresenceTurnGate, run_presence_turn +from ouroboros.task_results import write_task_result from ouroboros.utils import atomic_write_json from tests.test_presence_runner import _admission, _event @@ -19,6 +20,8 @@ def test_receipt_mode_defers_only_outgoing_log_until_transport_confirmation(tmp_ class Agent: def handle_task(self, task): captured.update(task) + # The durable terminal is the authority the Host reads back; the envelope alone never answers. + write_task_result(tmp_path / "data", task["id"], "completed", metadata=task["metadata"], result="The reply") return [{"type": "presence_result", "outcome": outcome, "text": "The reply", "work_ref": "work-1"}] result = run_presence_turn( diff --git a/tests/test_presence_orphan_replay.py b/tests/test_presence_orphan_replay.py index b50696275..7c9485689 100644 --- a/tests/test_presence_orphan_replay.py +++ b/tests/test_presence_orphan_replay.py @@ -147,8 +147,14 @@ def test_only_the_orphan_placeholder_of_a_presence_turn_reopens(tmp_path, monkey before = task_result_path(tmp_path, task_id).read_bytes() assert reopen_reconciled_presence_placeholder(tmp_path, task_id) is False assert task_result_path(tmp_path, task_id).read_bytes() == before - if seed is not _non_presence_orphan: - assert _cached_result(tmp_path, task_id) is not None # a real terminal replays + if seed is _ordinary_failure: + assert _cached_result(tmp_path, task_id) is not None # the model's own terminal replays + elif seed is _outcome_failure: + # A host infrastructure terminal answered nothing: it is never reopened as a placeholder, and + # replay keeps the event with the transport instead of acknowledging it as silent. + with pytest.raises(PresenceTurnError) as refused: + _cached_result(tmp_path, task_id) + assert refused.value.code == "presence_attempt_outcome_unknown" write_task_result(tmp_path, task_id, STATUS_COMPLETED, result="Late answer") assert load_task_result(tmp_path, task_id)["status"] == STATUS_FAILED # sticky terminal unchanged @@ -605,7 +611,8 @@ def test_a_refused_retry_asks_the_owner_once_and_never_the_correspondent(tmp_pat assert {k: v for k, v in after.items() if k not in stamp} == {k: v for k, v in before.items() if k not in stamp} class Quiet: - def handle_task(self, _task): + def handle_task(self, task): + write_task_result(tmp_path, task["id"], STATUS_COMPLETED, metadata=task["metadata"], result="") return [{"type": "presence_result", "outcome": "silent", "text": "", "work_ref": ""}] fresh = run_presence_turn(**{**kwargs, "event": replace(kwargs["event"], source_event_id="telegram:bot-1:43"), diff --git a/tests/test_presence_refusals.py b/tests/test_presence_refusals.py index e5e8a3f84..8c7c21d8c 100644 --- a/tests/test_presence_refusals.py +++ b/tests/test_presence_refusals.py @@ -552,3 +552,70 @@ def test_ambiguous_start_write_never_regenerates_after_error(tmp_path, monkeypat again = asyncio.run(_turn(app, binding, "event")) assert again.status_code == 409 and json.loads(again.body)["code"] == "presence_attempt_outcome_unknown" assert invoked == [] and ctx.presence_turns.live() == [] + + +def test_unknown_outcome_fallback_terminal_never_acknowledges_the_event(tmp_path, monkeypatch): + """A quota-refused primary whose fallback died with an unknown outcome is not a silent answer. + + The forced rail words that terminal ``provider_unavailable`` (the unknown fence outranks the + refusal source), so the guard cannot key on the resource-refusal word alone: the durable + infrastructure terminal keeps the event with the transport on the first call and on replay, + preserves already scheduled work, and asks the owner once — never a retry certificate. + """ + child_id = "scheduled-after-quota" + invoked, notices = [], [] + monkeypatch.setattr("ouroboros.presence_runner._write_unresolved_notice", + lambda _root, task_id: notices.append(task_id)) + + class Agent: + def handle_task(self, task): + invoked.append(task["id"]) + write_task_result(tmp_path, task["id"], "failed", result="[PROVIDER_UNAVAILABLE] host text", + metadata={**task["metadata"], "presence_work_ref": child_id}, + reason_code="provider_unavailable", terminal_origin="host_notice", + outcome_axes={"execution": {"status": "infra_failed", + "reason_code": "provider_unavailable", + "source": "provider_outcome_unknown_no_resend"}}) + return [{"type": "presence_result", "outcome": "silent", "text": "", "work_ref": child_id}] + + app, binding, ctx = _presence_app(tmp_path, lambda **kwargs: run_presence_turn( + repo_dir=tmp_path, drive_root=tmp_path, agent_factory=lambda **_kw: Agent(), + gate=PresenceTurnGate(1), **kwargs)) + for _ in range(2): + response = asyncio.run(_turn(app, binding, "event")) + body = json.loads(response.body) + assert response.status_code == 409 and body["code"] == "presence_attempt_outcome_unknown" + assert body["disposition"] == "retry" and not body.get("text") + assert body["work_ref"] == child_id + assert not ctx.presence_turns.live() and not any(ctx._inflight.values()) + assert invoked == [presence_turn_task_id(binding, "event")] + assert notices == [presence_turn_task_id(binding, "event")] * 2 + stored = load_task_result(tmp_path, presence_turn_task_id(binding, "event")) + assert "presence_retry_proof" not in (stored.get("metadata") or {}) # unknown is not not_started + + +def test_lost_terminal_write_after_start_barrier_never_acknowledges_the_event(tmp_path, monkeypatch): + """The in-memory envelope is not authority: a RUNNING row after handle_task is an unproven effect. + + The pipeline logs and swallows a failed terminal write; the Host must then refuse with the + scheduled work preserved from the execution's own handoff fact, and a retry must not regenerate. + """ + child_id = "scheduled-before-terminal-loss" + invoked = [] + + class Agent: + def handle_task(self, task): + invoked.append(task["id"]) # the terminal write failed after the durable start + return [{"type": "presence_result", "outcome": "message", "text": "answer", "work_ref": child_id}] + + app, binding, ctx = _presence_app(tmp_path, lambda **kwargs: run_presence_turn( + repo_dir=tmp_path, drive_root=tmp_path, agent_factory=lambda **_kw: Agent(), + gate=PresenceTurnGate(1), **kwargs)) + first = asyncio.run(_turn(app, binding, "event")) + body = json.loads(first.body) + assert first.status_code == 409 and body["code"] == "presence_attempt_outcome_unknown" + assert body["disposition"] == "retry" and not body.get("text") and body["work_ref"] == child_id + assert load_task_result(tmp_path, presence_turn_task_id(binding, "event"))["status"] == "running" + again = asyncio.run(_turn(app, binding, "event")) + assert again.status_code == 409 and json.loads(again.body)["code"] == "presence_attempt_outcome_unknown" + assert invoked == [presence_turn_task_id(binding, "event")] and ctx.presence_turns.live() == [] diff --git a/tests/test_presence_runner.py b/tests/test_presence_runner.py index d7606e7ef..ca7e95f1b 100644 --- a/tests/test_presence_runner.py +++ b/tests/test_presence_runner.py @@ -23,6 +23,12 @@ from ouroboros.presence_runner import ( PresenceTurnGate, run_presence_turn, ) +from ouroboros.task_results import write_task_result + + +def _terminal(drive_root, task, reply): + """The real pipeline's durable terminal: the Host reads it back, the envelope alone never answers.""" + write_task_result(pathlib.Path(drive_root), task["id"], "completed", metadata=task["metadata"], result=reply) def _admission() -> PresenceAdmission: @@ -84,6 +90,7 @@ def test_runner_builds_bounded_fresh_task_and_logs_shared_dialogue(tmp_path): class Agent: def handle_task(self, task): captured.update(task) + _terminal(data, task, "Hi") return [{"type": "presence_result", "outcome": "message", "text": "Hi", "work_ref": ""}] result = run_presence_turn( @@ -136,6 +143,7 @@ def test_presence_initial_attachment_rejection_defaults_to_partial_staging(tmp_p class Agent: def handle_task(self, task): seen_tasks.append(task) + _terminal(data, task, "ok") return [{"type": "presence_result", "outcome": "message", "text": "ok"}] result = run_presence_turn( @@ -299,6 +307,7 @@ def test_presence_turn_is_live_for_liveness_readers_but_never_an_owner_target(mo seen["with_main"] = observe(task["id"]) finally: registry.unregister("main-turn") + _terminal(tmp_path, task, "") return [{"type": "presence_result", "outcome": "silent", "text": "", "work_ref": ""}] try: