mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
P4.3: address direct reaper and crash notices by the project binding
The reaper incident notice and the worker-crash toast are DIRECT sends: nothing re-addresses them downstream, so a task converted into a project mid-run was told about its timeout or crash in the chat it was born in, while the handler-mediated task_done path already resolves the binding first. _incident_chat_id takes the supervisor queue module as an optional ctx and puts the bound project chat (own binding, then parent, then root) ahead of the row chat, with the owner chat still the absent-binding fallback; both call sites in reap_timed_out_task pass it. task_reaper.py stays at exactly 1600 lines: the two imports ride the existing module-level lines, which also removes the lazy import inside the function. worker_health resolves the same binding for the crash toast, its single consumer. _cascade_delivery_row_locked tests chat_id with `is not None`, so a descendant homed in the hidden partition (chat 0) is a usable routing row again instead of being skipped, and _routing_project_address does the same for a project row. The truthiness lint gains one alternative for the form both of those used - the id read straight off a mapping inside the condition, where no local named chat_id exists for the older alternatives to see. A comparison that merely reads the value stays out; after the two fixes the new form has no occurrences and no allowlist row, so it constrains new code only. Issues I30c and I30d; disposition R12 (plain import, net-zero file).
This commit is contained in:
parent
1058e5824f
commit
ce2b9b15bc
7 changed files with 167 additions and 17 deletions
|
|
@ -522,11 +522,11 @@ def _cascade_delivery_row_locked(q: Any, task_id: str) -> Dict[str, Any]:
|
|||
(Moved verbatim from ``task_lifecycle.py`` at its module-size boundary.)
|
||||
"""
|
||||
for task in q.PENDING:
|
||||
if isinstance(task, dict) and q._is_descendant_of(task, task_id) and task.get("chat_id"):
|
||||
if isinstance(task, dict) and q._is_descendant_of(task, task_id) and task.get("chat_id") is not None:
|
||||
return dict(task)
|
||||
for meta in q.RUNNING.values():
|
||||
task = meta.get("task") if isinstance(meta, dict) else None
|
||||
if isinstance(task, dict) and q._is_descendant_of(task, task_id) and task.get("chat_id"):
|
||||
if isinstance(task, dict) and q._is_descendant_of(task, task_id) and task.get("chat_id") is not None:
|
||||
return dict(task)
|
||||
return {}
|
||||
|
||||
|
|
|
|||
|
|
@ -39,7 +39,7 @@ def _routing_project_address(ctx: Any, target: str, status: str) -> Dict[str, An
|
|||
|
||||
binding = project_binding_for_task(ctx.DRIVE_ROOT, target) or {}
|
||||
project = get_project(ctx.DRIVE_ROOT, str(binding.get("project_id") or ""))
|
||||
if project and project.get("chat_id") and project.get("lifecycle") not in {"deleting", "deleted"}:
|
||||
if project and project.get("chat_id") is not None and project.get("lifecycle") not in {"deleting", "deleted"}:
|
||||
return {"project_id": project["id"], "project_chat_id": int(project["chat_id"])}
|
||||
except Exception:
|
||||
log.debug("Routing destination projection unavailable", exc_info=True)
|
||||
|
|
|
|||
|
|
@ -20,8 +20,8 @@ from typing import Any, Dict, Optional
|
|||
|
||||
from ouroboros.outcomes import EXECUTION_INFRA_FAILED, terminal_outcome_axes
|
||||
from ouroboros.utils import append_jsonl, utc_now_iso
|
||||
from supervisor.events import HOST_NARRATION
|
||||
from supervisor.message_bus import send_with_budget
|
||||
from supervisor.events import HOST_NARRATION, _bound_project_chat_id
|
||||
from supervisor.message_bus import notification_chat_route, send_with_budget
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
|
@ -1006,10 +1006,11 @@ def _emit_cancel_suppressed_retry_task_done(
|
|||
return False
|
||||
|
||||
|
||||
def _incident_chat_id(task: Any, owner_chat_id: int) -> Optional[int]:
|
||||
"""C4: an incident notice belongs to the TASK'S OWN chat; the owner chat is
|
||||
only the absent-binding fallback (the same precedence queue.py already uses
|
||||
for grace episodes).
|
||||
def _incident_chat_id(task: Any, owner_chat_id: int, ctx: Any = None) -> Optional[int]:
|
||||
"""C4: an incident notice belongs to the TASK'S OWN chat, and a DURABLE project binding
|
||||
outranks even that (this send is DIRECT, so nothing re-addresses it downstream, while a
|
||||
task converted into a project mid-run keeps its origin chat on the row); the owner chat
|
||||
is only the absent-binding fallback (the same precedence queue.py uses for grace episodes).
|
||||
|
||||
Routed through the ONE notification normalizer, so membership decides instead
|
||||
of truthiness: chat **0 is the Skill Review panel** and a task bound there
|
||||
|
|
@ -1017,12 +1018,11 @@ def _incident_chat_id(task: Any, owner_chat_id: int) -> Optional[int]:
|
|||
AND then refused to send, because `if owner_chat_id:` drops 0 as well). A
|
||||
negative (A2A/internal) chat is suppressed and falls through to the owner
|
||||
fallback; ``None`` means there is no deliverable route at all."""
|
||||
from supervisor.message_bus import notification_chat_route
|
||||
|
||||
row = task if isinstance(task, dict) else {}
|
||||
return notification_chat_route(
|
||||
task.get("chat_id") if isinstance(task, dict) else None,
|
||||
# A 0/absent owner chat is "not configured", not the panel — only an
|
||||
# explicit TASK binding routes to 0.
|
||||
_bound_project_chat_id(ctx, row.get("id"), row.get("parent_task_id"), row.get("root_task_id")) or None,
|
||||
row.get("chat_id"),
|
||||
# A 0/absent owner chat is "not configured", not the panel — only an explicit TASK binding routes to 0.
|
||||
owner_chat_id or None,
|
||||
)
|
||||
|
||||
|
|
@ -1355,7 +1355,7 @@ def reap_timed_out_task(job: Dict[str, Any]) -> None:
|
|||
return
|
||||
if not confirmed:
|
||||
_hold_wedged_worker(task_id, task_type, worker_id, terminal_reason, runtime_sec,
|
||||
_incident_chat_id(task, owner_chat_id))
|
||||
_incident_chat_id(task, owner_chat_id, _q))
|
||||
return
|
||||
|
||||
workers_mod._reconcile_confirmed_dead_review_owner(
|
||||
|
|
@ -1550,7 +1550,7 @@ def reap_timed_out_task(job: Dict[str, Any]) -> None:
|
|||
log.debug("Reaper: failed to log task_terminal_timeout for %s", task_id, exc_info=True)
|
||||
|
||||
# Notification failures cannot hold the slot; route to its task's chat.
|
||||
incident_chat_id = _incident_chat_id(task, owner_chat_id)
|
||||
incident_chat_id = _incident_chat_id(task, owner_chat_id, _q)
|
||||
if incident_chat_id is not None:
|
||||
try:
|
||||
if requeued:
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ from ouroboros.outcomes import (
|
|||
EXECUTION_INFRA_FAILED,
|
||||
terminal_outcome_axes,
|
||||
)
|
||||
from supervisor.log_addressing import resolve_project_chat
|
||||
from supervisor.queue import _queue_lock
|
||||
|
||||
|
||||
|
|
@ -289,7 +290,16 @@ def _recover_crashed_task_without_terminal(job: dict, queue: Any) -> None:
|
|||
# Signal crashes are terminal infrastructure failures for every task type.
|
||||
is_crash_signal = isinstance(exitcode, int) and exitcode < 0
|
||||
crash_signal = -exitcode if is_crash_signal else None
|
||||
chat_id = _pool().coerce_chat_identity(task.get("chat_id"), 0)
|
||||
# The crash toast is a DIRECT send - nothing re-addresses it downstream - so
|
||||
# the task's durable project binding has to win here, or a task converted
|
||||
# into a project mid-run is told about its crash in the chat it was born in.
|
||||
chat_id = _pool().coerce_chat_identity(
|
||||
resolve_project_chat(
|
||||
_pool().DRIVE_ROOT, task_id, task.get("parent_task_id"), task.get("root_task_id")
|
||||
)
|
||||
or task.get("chat_id"),
|
||||
0,
|
||||
)
|
||||
attempt = int(task.get("_attempt") or 1)
|
||||
replay_unsafe = (not getattr(w, "active_capacity", True)
|
||||
or has_owner_wait_checkpoint(meta, attempt))
|
||||
|
|
|
|||
|
|
@ -201,6 +201,40 @@ class TestReaperIncidentChat:
|
|||
assert _incident_chat_id({}, 1) == 1
|
||||
assert _incident_chat_id(None, 1) == 1
|
||||
|
||||
def test_project_binding_beats_the_row_chat(self, tmp_path):
|
||||
"""A reaper incident is a DIRECT send, so the binding must win HERE: a
|
||||
task converted into a project mid-run keeps its origin chat on the row."""
|
||||
from types import SimpleNamespace
|
||||
|
||||
from ouroboros.projects_registry import bind_task_to_project
|
||||
from supervisor.task_reaper import _incident_chat_id
|
||||
|
||||
bind_task_to_project(tmp_path, "root-inc", "reap-proj", 5150, origin={"absent": "system"})
|
||||
ctx = SimpleNamespace(DRIVE_ROOT=tmp_path)
|
||||
|
||||
assert _incident_chat_id({"id": "root-inc", "chat_id": 1}, 1, ctx) == 5150
|
||||
# A child is never bound itself; it inherits the room through lineage.
|
||||
assert _incident_chat_id(
|
||||
{"id": "child-inc", "root_task_id": "root-inc", "chat_id": 1}, 1, ctx
|
||||
) == 5150
|
||||
# An unbound task keeps today's order: its own chat, owner as fallback.
|
||||
assert _incident_chat_id({"id": "unbound-inc", "chat_id": 7}, 1, ctx) == 7
|
||||
assert _incident_chat_id({"id": "unbound-inc"}, 1, ctx) == 1
|
||||
|
||||
def test_a2a_binding_falls_through_to_the_task_chat(self, tmp_path):
|
||||
"""Suppression still applies to the new first candidate: a synthetic
|
||||
(negative) bound chat is skipped, it does not silence the notice."""
|
||||
from types import SimpleNamespace
|
||||
|
||||
from ouroboros.projects_registry import bind_task_to_project
|
||||
from supervisor.task_reaper import _incident_chat_id
|
||||
|
||||
bind_task_to_project(tmp_path, "a2a-inc", "a2a-proj", -1001, origin={"absent": "system"})
|
||||
ctx = SimpleNamespace(DRIVE_ROOT=tmp_path)
|
||||
|
||||
assert _incident_chat_id({"id": "a2a-inc", "chat_id": 7}, 1, ctx) == 7
|
||||
assert _incident_chat_id({"id": "a2a-inc"}, 1, ctx) == 1
|
||||
|
||||
def test_negative_task_chat_never_reaches_a_human_stream(self):
|
||||
from supervisor.task_reaper import _incident_chat_id
|
||||
|
||||
|
|
@ -209,3 +243,35 @@ class TestReaperIncidentChat:
|
|||
# all — reported as None, not as the panel.
|
||||
assert _incident_chat_id({"chat_id": -42}, 0) is None
|
||||
assert _incident_chat_id({"chat_id": "junk"}, -3) is None
|
||||
|
||||
|
||||
class TestCascadeDeliveryRow:
|
||||
"""A cancel cascade whose root already left the live maps borrows a live
|
||||
descendant's row for its lineage chat. Reading that chat for TRUTH skipped
|
||||
every descendant homed in the hidden partition (chat 0), so the cascade
|
||||
silently had no row and the notice went nowhere."""
|
||||
|
||||
def _queue(self, pending, running):
|
||||
from supervisor import queue as real_queue
|
||||
|
||||
class _Q:
|
||||
PENDING = pending
|
||||
RUNNING = running
|
||||
_is_descendant_of = staticmethod(real_queue._is_descendant_of)
|
||||
|
||||
return _Q
|
||||
|
||||
def test_a_descendant_homed_in_the_partition_is_returned(self):
|
||||
from supervisor.cancel_publication import _cascade_delivery_row_locked
|
||||
|
||||
child = {"id": "kid", "root_task_id": "root-9", "chat_id": 0}
|
||||
assert _cascade_delivery_row_locked(self._queue([child], {}), "root-9") == child
|
||||
assert _cascade_delivery_row_locked(
|
||||
self._queue([], {"kid": {"task": child}}), "root-9"
|
||||
) == child
|
||||
|
||||
def test_a_descendant_without_a_chat_is_still_skipped(self):
|
||||
from supervisor.cancel_publication import _cascade_delivery_row_locked
|
||||
|
||||
child = {"id": "kid", "root_task_id": "root-9"}
|
||||
assert _cascade_delivery_row_locked(self._queue([child], {}), "root-9") == {}
|
||||
|
|
|
|||
|
|
@ -30,10 +30,16 @@ POISONED_RECORD = re.compile(r'(?:"chat_id":|\bchat_id=)\s*[^,\n]*\bor\s+None')
|
|||
# hidden partition, and `if chat_id:` guards a send the same wrong way.
|
||||
# ``owner_chat_id`` is deliberately exempt — a 0/absent OWNER chat means "no
|
||||
# owner chat is configured", never the panel, so testing it for truth is honest.
|
||||
# The third alternative is the same habit written WITHOUT a local: the id read
|
||||
# straight off a mapping inside the condition (`... and task.get("chat_id"):`),
|
||||
# which is how a cascade silently skipped a descendant homed in the partition.
|
||||
# A comparison that merely reads the value (`int(t.get("chat_id") or 0) == x`)
|
||||
# decides no route and stays out.
|
||||
_CHAT_NAME = r"(?!owner_chat_id\b)(?:[A-Za-z_]*_)?chat_id"
|
||||
TRUTHY_ROUTE = re.compile(
|
||||
rf"^\s*if (?:not {_CHAT_NAME}\b|(?:[^:\n]*\band )?{_CHAT_NAME}\s*:)"
|
||||
rf"|^\s*if not [^:\n]*\bor not {_CHAT_NAME}\b"
|
||||
rf"|^\s*if [^:\n]*\.get\(\s*[\"']chat_id[\"']\s*\)\s*(?::|and\b)"
|
||||
)
|
||||
|
||||
# (repo-relative path, exact stripped line) -> (occurrences, why it stays)
|
||||
|
|
@ -120,6 +126,28 @@ def test_no_new_truthiness_route_for_a_chat_id():
|
|||
)
|
||||
|
||||
|
||||
def test_the_lint_sees_the_mapping_read_form():
|
||||
"""The widened alternative, pinned by the two lines that motivated it.
|
||||
|
||||
Both defects read the id straight off a mapping inside the condition, so no
|
||||
local named ``chat_id`` existed for the first two alternatives to see.
|
||||
"""
|
||||
caught = (
|
||||
' if isinstance(task, dict) and q._is_descendant_of(task, task_id) and task.get("chat_id"):',
|
||||
' if project and project.get("chat_id") and project.get("lifecycle") not in {"deleting"}:',
|
||||
' if not task.get("chat_id"):',
|
||||
)
|
||||
ignored = (
|
||||
' if isinstance(task, dict) and task.get("chat_id") is not None:',
|
||||
' if project and project.get("chat_id") is not None and project.get("id"):',
|
||||
' if int(t.get("chat_id") or 0) == chat_id:',
|
||||
)
|
||||
for line in caught:
|
||||
assert TRUTHY_ROUTE.search(line), line
|
||||
for line in ignored:
|
||||
assert not TRUTHY_ROUTE.search(line), line
|
||||
|
||||
|
||||
def test_no_record_stores_the_hidden_partition_as_absent():
|
||||
hits = _hits(POISONED_RECORD)
|
||||
assert not hits, (
|
||||
|
|
|
|||
|
|
@ -664,6 +664,52 @@ def test_signal_crash_is_terminal_no_retry(tmp_path):
|
|||
}
|
||||
|
||||
|
||||
def test_crash_toast_for_a_bound_task_goes_to_its_project_chat(tmp_path):
|
||||
"""The crash toast is a DIRECT send, so the durable project binding has to
|
||||
win over the chat the task was born in: a task converted into a project
|
||||
mid-run would otherwise be told about its crash in Main."""
|
||||
import supervisor.workers as W
|
||||
import queue as _queue
|
||||
|
||||
from ouroboros.projects_registry import bind_task_to_project
|
||||
|
||||
bind_task_to_project(tmp_path, "sig02", "crash-proj", 5151, origin={"absent": "system"})
|
||||
task = _make_task(task_id="sig02", attempt=1, chat_id=1) # born in Main
|
||||
worker = _make_worker(busy_task_id="sig02", exitcode=-11)
|
||||
|
||||
W.DRIVE_ROOT = tmp_path
|
||||
(tmp_path / "logs").mkdir(parents=True, exist_ok=True)
|
||||
W.QUEUE_MAX_RETRIES = 1
|
||||
W.WORKERS = {0: worker}
|
||||
W.RUNNING = {
|
||||
"sig02": {
|
||||
"task": task,
|
||||
"started_at": time.time() - 5,
|
||||
"last_heartbeat_at": time.time() - 5,
|
||||
"attempt": 1,
|
||||
}
|
||||
}
|
||||
W._LAST_SPAWN_TIME = 0
|
||||
W.CRASH_TS = []
|
||||
|
||||
import supervisor.queue as sq
|
||||
|
||||
incident_notice = MagicMock()
|
||||
with patch.object(sq, "enqueue_task", side_effect=lambda t, front=False: None), \
|
||||
patch.object(sq, "persist_queue_snapshot", MagicMock()), \
|
||||
patch("supervisor.workers.respawn_worker"), \
|
||||
patch("supervisor.workers.load_state", return_value={}), \
|
||||
patch("ouroboros.task_results.load_task_result", return_value=None), \
|
||||
patch("ouroboros.task_results.write_task_result", side_effect=lambda *a, **k: None), \
|
||||
patch("supervisor.workers.get_event_q", return_value=_queue.Queue()), \
|
||||
patch("supervisor.workers.send_with_budget", incident_notice), \
|
||||
patch("supervisor.message_bus.get_bridge", return_value=None):
|
||||
_run_health_and_reap()
|
||||
|
||||
incident_notice.assert_called_once()
|
||||
assert incident_notice.call_args[0][0] == 5151
|
||||
|
||||
|
||||
def test_deep_self_review_crash_emits_task_done_event(tmp_path):
|
||||
"""deep_self_review crash must emit task_done so the UI live card closes."""
|
||||
import supervisor.workers as W
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue