Bind a timeout retry to its predecessor's Project at admission

A Project-bound root that timed out was retried under a new physical id
which copied the owner-message origin ref but carried NO durable binding,
so the retried work painted a Main card offering "Turn into project" — a
second convertible unit for one piece of work — until some later implicit
act adopted it. docs/DEVELOPMENT.md described that as deliberate; it was
the #900 implementer's size-cap deferral, and the owner accepted that a
retry inherits the old id's binding on 2026-09-14.

The reaper now binds the successor inside its own retry admission
transaction (task_reaper._run_retry_admission_transaction ->
worker_promotion.bind_retry_to_origin_project), so admission and bind are
one step:

- the claim lock moves OUTSIDE the queue lock, the order an implicit claim
  already uses (project_id_for_origin's live-task tie-break takes the
  queue lock under it), so nothing inverts and a concurrent conversion of
  the same owner message cannot interleave;
- the bind lands only after cancellation has lost the boundary. A binding
  is immutable, so a bound-but-never-admitted retry id would answer
  project_id_for_task forever; a retry suppressed by a cancelled or
  already-terminal root is therefore never bound;
- the predecessor's own binding answers first (the SSOT for that task's
  project), then the origin-keyed lookup for a root that was never bound
  itself while another task id of the same message was; the new row reuses
  the predecessor's stored origin BY VALUE, so the retry joins that
  message's one convertible unit instead of starting a second;
- a refused bind or an unreadable store leaves the retry unbound, admits
  it anyway, and discloses project_binding_failed on the reaper's own
  drive (P1) instead of holding up the retry.

supervisor/task_reaper.py sat exactly on the 1600-line hard cap, so the
call site was paid for in place: "the retry id never ran, drop its copied
inputs" existed three times in _enqueue_retry and is now one
_discard_retry_inputs owner (1600 -> 1598 lines).

Docs flip with the behaviour: docs/DEVELOPMENT.md (identity-of-the-work
section) and docs/ARCHITECTURE.md (the origin membership list and the
implicit-claim list now name the reaper's retry admission).

Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
This commit is contained in:
Ouroboros 2026-09-15 16:56:30 +03:00
parent ddd5af5e5f
commit 913034c60e
5 changed files with 362 additions and 42 deletions

File diff suppressed because one or more lines are too long

View file

@ -293,12 +293,21 @@ exercises that seam. For fuzzy entities use the LLM-first pattern
The same captured reference is also the IDENTITY OF THE WORK, not just its
provenance: a new task id minted from the same owner message (a promoted root,
a mid-run scope call) must INHERIT that origin's project binding
(`projects_registry.project_id_for_origin`, keyed by value on chat id +
client message id), never re-derive project membership from its own id — one
convertible unit per message, not one per task id. A timeout retry COPIES the
origin ref onto its new id but is deliberately NOT bound at clone time; it joins
that origin's project when a later implicit act adopts it.
a mid-run scope call, the timeout retry that replaces a dead attempt) must
INHERIT that origin's project binding (`projects_registry.project_id_for_origin`,
keyed by value on chat id + client message id), never re-derive project
membership from its own id — one convertible unit per message, not one per task
id. A timeout retry is bound at RETRY ADMISSION, inside the same transaction and
under the same claim lock that admits it
(`worker_promotion.bind_retry_to_origin_project`, called from the reaper's
`_run_retry_admission_transaction`): the predecessor's own binding answers first,
then that origin's, and the new row reuses the predecessor's stored origin by
value, so the retried work stays in its room instead of painting a Main card that
offers to turn it into a project. The bind lands ONLY once cancellation can no
longer win the boundary — a binding is immutable, so a bound-but-never-admitted
retry id would answer `project_id_for_task` forever — and a retry suppressed by a
cancelled or already-terminal root is therefore never bound. Enforced by
`tests/test_retry_project_binding.py`.
One named exception inside role (b): a verification RECEIPT with no earlier
ingress point is reconciled by ONE TYPED IDENTITY KEY, matching on the key's

View file

@ -434,8 +434,10 @@ def _run_retry_admission_transaction(
salvage_note: str,
custody_audit: Optional[Dict[str, Any]] = None,
) -> tuple[Dict[str, str], str]:
"""Publish one retry admission under the queue -> cancel lock order."""
"""Publish one retry admission under the claim -> queue -> cancel lock order."""
from ouroboros.projects_registry import origin_claim_lock
from supervisor.cancel_publication import _custody_disclosure_fields
from supervisor.worker_promotion import bind_retry_to_origin_project
from ouroboros.task_results import (
STATUS_CANCELLED,
@ -461,7 +463,11 @@ def _run_retry_admission_transaction(
# Match assignment's lock order (queue -> cancel projection). Holding
# only the projection lock while enqueue_task takes the queue lock would
# invert _drop_cancelled_pending and deadlock with ordinary dispatch.
with q._queue_lock:
# The claim lock is OUTERMOST because an implicit claim already takes it
# before the queue lock (project_id_for_origin's live-task tie-break), and
# it makes this admission and its project bind ONE transaction against a
# concurrent conversion of the same owner message.
with origin_claim_lock(), q._queue_lock:
with cancellation_projection_lock(q.DRIVE_ROOT):
intents = active_intents(q.DRIVE_ROOT, strict=True)
if not isinstance(intents, dict):
@ -653,6 +659,10 @@ def _run_retry_admission_transaction(
suppression["retry_status"] = str(
retry_terminal.get("status") or ""
)
else:
# Cancellation lost the boundary and the successor is
# durable: the retry inherits its predecessor's room.
bind_retry_to_origin_project(q.DRIVE_ROOT, task, task_id, retry_task_id)
except Exception:
admission_block = "cancel_intent_authority_unreadable"
log.error(
@ -665,6 +675,18 @@ def _run_retry_admission_transaction(
return suppression, admission_block
def _discard_retry_inputs(q: Any, task: Dict[str, Any], task_id: str, retry_task_id: str) -> None:
"""Drop the inputs copied onto a retry id that will never run."""
if not retry_task_id or retry_task_id == task_id:
return
from ouroboros.artifacts import task_artifact_dir_path
from ouroboros.owner_mailbox import cleanup_task_mailbox
drive = q._task_drive_for_task(task, task_id)
cleanup_task_mailbox(drive, retry_task_id)
shutil.rmtree(task_artifact_dir_path(drive, retry_task_id), ignore_errors=True)
def _enqueue_retry(
q: Any,
task: Dict[str, Any],
@ -687,11 +709,7 @@ def _enqueue_retry(
admission wins, a later SINGLE cancel resolves the complete durable chain
under the same projection lock and targets the physical leaf.
"""
from ouroboros.task_results import (
STATUS_FAILED,
load_task_result,
write_task_result,
)
from ouroboros.task_results import STATUS_FAILED, load_task_result, write_task_result
from supervisor.cancel_publication import _custody_disclosure_fields
retried = dict(task)
@ -701,11 +719,8 @@ def _enqueue_retry(
retried["timeout_retry_from"] = task_id
retried["timeout_retry_at"] = utc_now_iso()
if retry_task_id and retry_task_id != task_id:
from ouroboros.artifacts import (
handoff_task_attachments_for_retry,
task_artifact_dir_path,
)
from ouroboros.owner_mailbox import cleanup_task_mailbox, copy_owner_mailbox_for_retry
from ouroboros.artifacts import handoff_task_attachments_for_retry
from ouroboros.owner_mailbox import copy_owner_mailbox_for_retry
task_drive = q._task_drive_for_task(task, task_id)
replacements, attachment_error = handoff_task_attachments_for_retry(
@ -726,10 +741,7 @@ def _enqueue_retry(
"Reaper: %s handoff failed for retry %s -> %s: %s",
failure, task_id, retry_task_id, attachment_error,
)
shutil.rmtree(
task_artifact_dir_path(task_drive, retry_task_id), ignore_errors=True,
)
cleanup_task_mailbox(task_drive, retry_task_id)
_discard_retry_inputs(q, task, task_id, retry_task_id)
outcome = terminal_outcome_axes(
lifecycle=STATUS_FAILED,
execution=EXECUTION_INFRA_FAILED,
@ -775,14 +787,7 @@ def _enqueue_retry(
)
if suppression:
if retry_task_id and retry_task_id != task_id:
cleanup_task_mailbox(q._task_drive_for_task(task, task_id), retry_task_id)
shutil.rmtree(
task_artifact_dir_path(
q._task_drive_for_task(task, task_id), retry_task_id,
),
ignore_errors=True,
)
_discard_retry_inputs(q, task, task_id, retry_task_id)
reason = (
"cancel_pending_retry_suppressed"
if suppression.get("kind") == "cancel_intent"
@ -791,14 +796,7 @@ def _enqueue_retry(
return False, attempt, reason, suppression
if not admission_block:
return True, attempt + 1, terminal_reason, {}
if retry_task_id and retry_task_id != task_id:
cleanup_task_mailbox(q._task_drive_for_task(task, task_id), retry_task_id)
shutil.rmtree(
task_artifact_dir_path(
q._task_drive_for_task(task, task_id), retry_task_id,
),
ignore_errors=True,
)
_discard_retry_inputs(q, task, task_id, retry_task_id)
blocked_reason = f"{terminal_reason}_retry_admission_blocked"
outcome = terminal_outcome_axes(

View file

@ -72,19 +72,24 @@ def _origin_from_task_record(task_id: str) -> Optional[dict]:
def _report_binding_failure(
task_id: str, project_id: str, exc: Exception, *, path: str, reason: str = "",
drive_root: Any = None,
) -> None:
"""A failed durable bind is LOUD (BIBLE P1: silent linkage loss is memory
loss): warning log + typed events.jsonl row; the task itself keeps running.
``reason`` names a REFUSAL that never reached the bind (today only
``reason`` names a REFUSAL that never reached the bind (today
``project_scope_conflict``: the task is already bound elsewhere, so no second
project is created); the row is otherwise the same shape a raising bind writes.
project is created; ``project_binding_unreadable``: the store could not be
read at all); the row is otherwise the same shape a raising bind writes.
``drive_root`` names the store the failure belongs to for a caller that does
not own the pool root — the reaper passes the queue's drive.
"""
# A refusal carries no live traceback, so only a real bind failure logs one.
log.warning("%s for %s/%s (%s)", reason or "bind_task_to_project failed",
task_id, project_id, path, exc_info=not reason)
try:
append_jsonl(_pool().DRIVE_ROOT / "logs" / "events.jsonl", {
root = _pool().DRIVE_ROOT if drive_root is None else drive_root
append_jsonl(root / "logs" / "events.jsonl", {
"ts": utc_now_iso(),
"type": "project_binding_failed",
"task_id": str(task_id or ""),
@ -316,6 +321,68 @@ def _admit_project_scope(
return None
def bind_retry_to_origin_project(
drive_root: Any, task: dict, task_id: str, retry_task_id: str,
) -> str:
"""Carry a timed-out root's Project onto the retry that replaces it.
A retry is the SAME work under a new physical id, so it belongs to the room
the work already has (DEVELOPMENT "the captured reference is the identity of
the work"). Before this, the new id copied the origin ref but carried no
binding, and the retried work painted a Main card offering "Turn into
project" until some later implicit act adopted it. The predecessor's own
binding answers first — it is the SSOT for that task's project — and the
origin-keyed lookup covers a root that was never bound itself while another
task id of the same owner message was. The new row reuses the predecessor's
stored origin BY VALUE, so the retry joins that message's one convertible
unit instead of starting a second.
The reaper calls this INSIDE its retry admission transaction, under the same
``origin_claim_lock`` every implicit claim holds, and only once cancellation
can no longer win the boundary: ``bind_task_to_project`` is immutable, so a
bound-but-never-admitted retry id would answer ``project_id_for_task``
forever. Creates no project and never raises — a refused bind (a project that
stopped accepting them) or an unreadable store leaves the retry unbound and
is disclosed as ``project_binding_failed``. Returns the project the retry was
bound to, "" when there was nothing to inherit.
"""
tid = str(retry_task_id or "").strip()
origin_id = str(task_id or "").strip()
if not tid or not origin_id or tid == origin_id:
return ""
from ouroboros.projects_registry import (
bind_task_to_project,
project_binding_for_task,
project_id_for_origin,
)
origin = _origin_from_mapping(task, absent="mid_task_no_origin")
try:
predecessor = project_binding_for_task(drive_root, origin_id) or {}
pid = str(predecessor.get("project_id") or "") or str(
project_id_for_origin(drive_root, origin.get("ref"), strict=True) or ""
)
except Exception as exc:
_report_binding_failure(tid, "", exc, path="timeout_retry_admission",
reason="project_binding_unreadable", drive_root=drive_root)
return ""
if not pid:
return ""
if isinstance(predecessor.get("source_ref"), dict):
origin = {"ref": dict(predecessor["source_ref"])}
if isinstance(predecessor.get("source_text"), str):
origin["text"] = predecessor["source_text"]
elif predecessor.get("origin_absent"):
origin = {"absent": str(predecessor["origin_absent"])}
try:
bind_task_to_project(drive_root, tid, pid, origin=origin)
except Exception as exc:
_report_binding_failure(tid, pid, exc, path="timeout_retry_admission",
drive_root=drive_root)
return ""
return pid
def promote_chat_to_task(evt: dict, ctx: Any) -> dict:
"""Enqueue a first-class pooled owner task from a conversation-lane promote.
The task carries the originating ``chat_id`` (its live card and replies

View file

@ -0,0 +1,246 @@
"""A timeout retry keeps the Project its predecessor's work already has.
The incident: a Project-bound root that timed out was retried under a new physical
id which copied the owner-message origin ref but carried no durable binding, so the
retried work painted a Main card offering "Turn into project" — a second convertible
unit for one piece of work — until some later implicit act adopted it.
The reaper now binds the successor inside its own admission transaction
(``task_reaper._run_retry_admission_transaction`` ->
``worker_promotion.bind_retry_to_origin_project``), under the claim lock every
implicit claim holds, and only once cancellation has lost that boundary: a binding
is immutable, so a bound-but-never-admitted retry id would answer
``project_id_for_task`` forever.
"""
from __future__ import annotations
import json
import types
from ouroboros import cancel_intents as ci
from ouroboros.contracts.chat_id_policy import project_chat_id
from ouroboros.project_dialogue import build_owner_message_ref
from ouroboros.projects_registry import (
all_task_project_bindings,
begin_project_deletion,
bind_task_to_project,
create_project,
project_binding_for_task,
project_id_for_task,
)
from ouroboros.task_results import STATUS_CANCELLED, STATUS_RUNNING, write_task_result
from tests._cancel_intents_shared import qenv as _qenv
qenv = _qenv
OWNER_TEXT = "the importer keeps timing out, please fix it"
def _owner_ref(client_message_id: str = "msg-retry-1", chat_id: int = 1) -> dict:
return build_owner_message_ref(
chat_id=chat_id,
client_message_id=client_message_id,
ts="2026-09-15T09:30:00+00:00",
text=OWNER_TEXT,
)
def _patch_retry_input_handoff(monkeypatch):
"""Attachment/mailbox carry-over is exercised by the reaper's own suites; these
tests are about what the admission transaction binds."""
monkeypatch.setattr(
"ouroboros.artifacts.handoff_task_attachments_for_retry",
lambda *_args, **_kwargs: ({}, ""),
)
monkeypatch.setattr(
"ouroboros.owner_mailbox.copy_owner_mailbox_for_retry",
lambda *_args, **_kwargs: True,
)
def _root_task(task_id: str, *, ref: dict | None = None) -> dict:
"""A pooled root exactly as the queue carries it. A root converted post-hoc
("Turn into project") carries NO ``project_id`` on its row — the durable
binding is the only truth about its Project, and it is the shape the retry
used to lose.
"""
task = {
"id": task_id,
"type": "task",
"chat_id": 1,
"depth": 0,
"root_task_id": task_id,
"parent_task_id": "",
"delegation_role": "root",
}
if ref is not None:
task["origin_message_ref"] = dict(ref)
task["origin_message_text"] = OWNER_TEXT
return task
def _bind_root(drive, task_id: str, pid: str, *, ref: dict | None) -> dict:
create_project(drive, pid, name="Importer room", origin="owner_ui")
origin = (
{"ref": dict(ref), "text": OWNER_TEXT} if ref is not None
else {"absent": "mid_task_no_origin"}
)
return bind_task_to_project(drive, task_id, pid, origin=origin)
def _retry(qenv, task: dict, old_id: str, new_id: str):
from supervisor import task_reaper as tr
return tr._enqueue_retry(
qenv.q,
task,
task_id=old_id,
retry_task_id=new_id,
attempt=1,
terminal_reason="idle_timeout",
recon_fields={},
)
def _events(drive) -> list[dict]:
path = drive / "logs" / "events.jsonl"
if not path.exists():
return []
return [
json.loads(line)
for line in path.read_text(encoding="utf-8").splitlines() if line.strip()
]
def test_retry_of_a_bound_root_is_in_the_project_before_its_first_round(qenv, monkeypatch):
"""(a) The successor is bound while it is still PENDING — no worker has picked
it up, so the owner never sees the retried work as an unclaimed Main card."""
_patch_retry_input_handoff(monkeypatch)
old_id, new_id, pid = "bound-old", "bound-new", "importer-room"
ref = _owner_ref()
root = _root_task(old_id, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_RUNNING, result="working")
predecessor = _bind_root(qenv.drive, old_id, pid, ref=ref)
requeued, new_attempt, _reason, suppression = _retry(qenv, root, old_id, new_id)
assert (requeued, new_attempt, suppression) == (True, 2, {})
assert [row["id"] for row in qenv.q.PENDING] == [new_id]
assert qenv.q.RUNNING == {}
assert project_id_for_task(qenv.drive, new_id) == pid
# The UI census (/api/state) is what paints the card's room.
assert all_task_project_bindings(qenv.drive)[new_id] == {
"project_id": pid, "chat_id": project_chat_id(pid),
}
# The retry joins the SAME owner-message unit: its origin is the predecessor's
# stored one, carried by value rather than re-derived.
binding = project_binding_for_task(qenv.drive, new_id)
assert binding["source_ref"] == predecessor["source_ref"]
assert binding["source_text"] == predecessor["source_text"]
def test_retry_adopts_the_project_a_sibling_of_the_same_owner_message_holds(qenv, monkeypatch):
"""The origin-keyed half: the timed-out root was never bound itself, but the
direct turn that received the same owner message was. One message, one room."""
_patch_retry_input_handoff(monkeypatch)
old_id, new_id, pid = "adopt-old", "adopt-new", "sibling-room"
ref = _owner_ref("msg-retry-sibling")
root = _root_task(old_id, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_RUNNING, result="working")
_bind_root(qenv.drive, "turn-that-received-it", pid, ref=ref)
assert project_id_for_task(qenv.drive, old_id) == ""
requeued, _attempt, _reason, suppression = _retry(qenv, root, old_id, new_id)
assert (requeued, suppression) == (True, {})
assert project_id_for_task(qenv.drive, new_id) == pid
def test_retry_of_a_cancelled_root_is_never_bound(qenv, monkeypatch):
"""(b) A root that already settled as cancelled has no successor to place: the
admission transaction suppresses the retry, so nothing durable claims its id."""
_patch_retry_input_handoff(monkeypatch)
old_id, new_id, pid = "cancelled-old", "cancelled-new", "cancelled-room"
ref = _owner_ref("msg-retry-cancelled")
root = _root_task(old_id, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_RUNNING, result="working")
_bind_root(qenv.drive, old_id, pid, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_CANCELLED, result="owner stopped it")
requeued, _attempt, reason, suppression = _retry(qenv, root, old_id, new_id)
assert requeued is False
assert reason == "terminal_result_retry_suppressed"
assert suppression["kind"] == "terminal_result"
assert qenv.q.PENDING == []
assert project_binding_for_task(qenv.drive, new_id) is None
assert new_id not in all_task_project_bindings(qenv.drive)
# The predecessor keeps its own room; only the successor is absent.
assert project_id_for_task(qenv.drive, old_id) == pid
def test_cancellation_winning_the_admission_leaves_no_retry_binding(qenv, monkeypatch):
"""(c) An intent recorded before the admission boundary suppresses the successor
inside the same locked transition, so the immutable bind never lands."""
_patch_retry_input_handoff(monkeypatch)
old_id, new_id, pid = "race-old", "race-new", "race-room"
ref = _owner_ref("msg-retry-race")
root = _root_task(old_id, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_RUNNING, result="working")
_bind_root(qenv.drive, old_id, pid, ref=ref)
ci.request_cancel(qenv.drive, old_id, reason="stop before retry")
requeued, _attempt, reason, suppression = _retry(qenv, root, old_id, new_id)
assert requeued is False
assert reason == "cancel_pending_retry_suppressed"
assert suppression == {"kind": "cancel_intent", "target": old_id}
assert qenv.q.PENDING == []
assert project_binding_for_task(qenv.drive, new_id) is None
def test_first_implicit_promote_from_a_retry_lands_in_the_origins_project(qenv, monkeypatch):
"""(d) What the owner actually sees next: work the retry promotes with no
explicit target goes to the room, not to Main."""
from ouroboros.tools.control_routing import _inherited_project_scope
_patch_retry_input_handoff(monkeypatch)
old_id, new_id, pid = "promote-old", "promote-new", "promote-room"
ref = _owner_ref("msg-retry-promote")
root = _root_task(old_id, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_RUNNING, result="working")
_bind_root(qenv.drive, old_id, pid, ref=ref)
monkeypatch.setattr("ouroboros.config.DATA_DIR", qenv.drive)
requeued, _attempt, _reason, _suppression = _retry(qenv, root, old_id, new_id)
assert requeued is True
# The retry worker's own scope copy is empty (a post-hoc conversion never
# reached the dead attempt), so only the durable binding can answer.
ctx = types.SimpleNamespace(task_id=new_id, task_metadata={}, project_id="")
assert _inherited_project_scope(ctx) == pid
def test_a_room_that_stopped_accepting_bindings_discloses_and_admits_the_retry(qenv, monkeypatch):
"""A refused bind is LOUD but never blocks the retry: the successor is still
admitted, and the failure is a typed row on the reaper's OWN drive."""
_patch_retry_input_handoff(monkeypatch)
old_id, new_id, pid = "fenced-old", "fenced-new", "fenced-room"
ref = _owner_ref("msg-retry-fenced")
root = _root_task(old_id, ref=ref)
write_task_result(qenv.drive, old_id, STATUS_RUNNING, result="working")
_bind_root(qenv.drive, old_id, pid, ref=ref)
begin_project_deletion(qenv.drive, pid)
requeued, _attempt, _reason, suppression = _retry(qenv, root, old_id, new_id)
assert (requeued, suppression) == (True, {})
assert [row["id"] for row in qenv.q.PENDING] == [new_id]
assert project_binding_for_task(qenv.drive, new_id) is None
failures = [
row for row in _events(qenv.drive)
if row.get("type") == "project_binding_failed" and row.get("task_id") == new_id
]
assert [row["bind_path"] for row in failures] == ["timeout_retry_admission"]
assert failures[0]["project_id"] == pid