fix: refuse Presence acknowledgement on infrastructure terminals and lost terminal writes

Independent exact-range review (Astra, 55ad78213..15dbc7533) found two P1 paths
where the Host acknowledged a Presence event the model never answered:

- a quota-refused primary whose fallback died with an unknown outcome is worded
  provider_unavailable by the forced rail (the unknown fence outranks the
  refusal source), so a guard keyed on resource_refusal_no_resend alone let the
  host-authored terminal project to completed/silent (HTTP 200) live and on
  replay, and the adapter could drop the event;
- a terminal write that fails after the durable start barrier is swallowed by
  the pipeline, and the in-memory presence_result envelope was returned as the
  answer while the durable row still said RUNNING.

One helper, _terminal_refusal, now owns the verdict the durable row demands
before acknowledgement: a confirmed resource refusal keeps the event
(presence_resources_unavailable); no canonical terminal (RUNNING/INTERRUPTED,
reconciled placeholder, lost terminal write) or any infra_failed execution
axis is presence_attempt_outcome_unknown — 409, empty external speech, the
admitted work_ref preserved (durable metadata, else the execution's own
handoff fact), one serialized owner-only notice, never a retry certificate.
Both guards (live and cached replay) consume it. Two consumer regressions
through the real Host app cover both paths; test fakes that returned an
envelope without writing a terminal now write the durable row the real
pipeline writes. Chapter 12 records the rule.
This commit is contained in:
Ouroboros 2026-09-25 23:51:52 +03:00
parent 15dbc7533e
commit 54c2be5517
8 changed files with 143 additions and 13 deletions

View file

@ -14,7 +14,7 @@ Operation correlation: a named injected message has `operation_ref=<chat_id>:<cl
Presence (`presence_runner.py`): `POST /presence/turn` admits exact provider/account/conversation/thread facts, skill-confined files and a reviewed binding. Auth/grant discovery and pre-turn admission, staging and replay reads share a two-permit async admission off the ASGI loop, leaving shared executor capacity for owner state and controls even during many slow probes. Conversation and active-turn gates hold cross-process leases: queued turns wait as coroutines, executing turns own ContextVars-preserving threads, and an HTTP disconnect does not cancel work or return capacity early. The Host joins a live retry by turn ID only when the canonical conversation, actor and text match; reporting-version and staged-file paths are not event identity. A mismatch returns typed `presence_event_identity_conflict` (409, rejected), never another room's answer. A legacy result whose source cannot be reconstructed refuses replay rather than guessing.
Before `handle_task` can call any model or tool, Presence writes and reads back a source-bound RUNNING row in the existing task-result store. A failed/ambiguous start write is `presence_start_unwritable` (409, retry); a persisted RUNNING/INTERRUPTED row or reconciled placeholder with no terminal is `presence_attempt_outcome_unknown` (409, retry), not permission to regenerate. A cancelled turn stays blocked. A later same-source terminal can outrank an old quarantined row; an unreadable or identity-mismatched row cannot. Provider monetary abandonment does not prove no external effect. Owner-only Main recovery notices serialize same-task replay readers through the existing Presence gate's process/file locks (without taking a turn slot), and require a chat write on a fresh JSONL boundary before the notified stamp; a crash between the two can still repeat a notice. When primary or fallback routes refuse quota, even if a later fallback errors or overflows, `resource_refusal_no_resend` is a retryable `presence_resources_unavailable` (409) without terminal speech: the original event stays with its transport, the gate is released and Main receives a deduplicated notice. An admitted child keeps its `work_ref` on this refusal for `/presence/work` polling; that handoff precludes the no-effect retry proof. A failed row is reusable only with terminal `presence_retry_proof`: first-round temporary refusal, *all* route/rotation operations engine-confirmed `not_started` and ledger-released, no response/tool/handoff, and a future reset. After reset the conversation gate writes `presence_retry_next`, a locked link to a new physical task ID derived from the predecessor ID, retaining the old result. Host joins by original event ID; the live/orphan gate and `turn_ref` use the physical ID. A crash after linking resumes the unused successor; another refusal requires its own proof and reset, not transport-cadence retries. Old/unproven/undated rows stay empty 409 and owner-only; the input is not logged twice. No arbitrary external API exactly-once claim.
Before `handle_task` can call any model or tool, Presence writes and reads back a source-bound RUNNING row in the existing task-result store. A failed/ambiguous start write is `presence_start_unwritable` (409, retry); a persisted RUNNING/INTERRUPTED row or reconciled placeholder with no terminal is `presence_attempt_outcome_unknown` (409, retry), not permission to regenerate. A cancelled turn stays blocked. A later same-source terminal can outrank an old quarantined row; an unreadable or identity-mismatched row cannot. Provider monetary abandonment does not prove no external effect. Owner-only Main recovery notices serialize same-task replay readers through the existing Presence gate's process/file locks (without taking a turn slot), and require a chat write on a fresh JSONL boundary before the notified stamp; a crash between the two can still repeat a notice. When primary or fallback routes refuse quota, even if a later fallback errors or overflows, `resource_refusal_no_resend` is a retryable `presence_resources_unavailable` (409) without terminal speech: the original event stays with its transport, the gate is released and Main receives a deduplicated notice. An admitted child keeps its `work_ref` on this refusal for `/presence/work` polling; that handoff precludes the no-effect retry proof. A failed row is reusable only with terminal `presence_retry_proof`: first-round temporary refusal, *all* route/rotation operations engine-confirmed `not_started` and ledger-released, no response/tool/handoff, and a future reset. After reset the conversation gate writes `presence_retry_next`, a locked link to a new physical task ID derived from the predecessor ID, retaining the old result. Host joins by original event ID; the live/orphan gate and `turn_ref` use the physical ID. A crash after linking resumes the unused successor; another refusal requires its own proof and reset, not transport-cadence retries. Old/unproven/undated rows stay empty 409 and owner-only; the input is not logged twice. No arbitrary external API exactly-once claim. The durable row is the one authority the Host reads back after `handle_task` (`_terminal_refusal`): a host-authored infrastructure terminal (`outcome_axes.execution.status=infra_failed` — a provider death or an unknown outcome behind a refusal, worded `provider_unavailable` by the forced rail) and a row still RUNNING because the terminal write failed are both `presence_attempt_outcome_unknown`, never completed/silent; the in-memory envelope supplies only the admitted `work_ref`.
All entry/receipt paths derive `presence_bindings.conversation_key` (empty thread = `0`). The Presence-local live set protects in-process work from orphan reaping but is not owner-addressable; update drains do not wait for it. The rebuildable last-turn pointer stores completed speech and real deferred `work_ref`; `/presence/work/{work_ref}` polls only correlated work, not the general task API. Turn, receipt and inject budgets are independent. `message`, `silent`, `tool_delivered`, `deferred` are distinct outcomes; deferred requires admitted work, orphaned children replay silent. Promotion preserves ceiling, cost and return context, refuses unusable folders and cannot widen Project/workspace/source. `presence_cancel_work` requires the current binding; related work may originate in any conversation on that binding. Owner chat and consciousness Act+ may initiate an enabled binding; agents share one autobiography.

View file

@ -524,11 +524,35 @@ def _stored_turn(drive_root: Path, task_id: str, identity: str = "") -> 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(

View file

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

View file

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

View file

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

View file

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

View file

@ -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() == []

View file

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