diff --git a/gateway/AGENTS.md b/gateway/AGENTS.md index df8b2dfd5e..4984974996 100644 --- a/gateway/AGENTS.md +++ b/gateway/AGENTS.md @@ -69,7 +69,11 @@ required. Notify-off drains these without waking. Transport failures are retried durable completion owners/transports do not spend delivery attempts. Profile-namespaced process events keep their owning adapter, including raw progress/final notices. Adapter acceptance is at-least-once admission, not proof that a model turn or final outbound reply completed. API-server -async completions remain durable delivery rows, never autonomous new model turns. +async completions remain durable delivery rows, never autonomous new model turns. Require the +`admit_internal_event` receipt for completion/watch injection: a handler returning None is not +acceptance. Refused admission refunds every claimed batch sibling without spending an attempt; +actual delivery errors keep their bounded retry policy. Recognized raw API routes resolve after +persisted messaging origins and defer quietly when unavailable; malformed routes still warn. Cron deliveries are NOT mirrored into the target gateway session — they land in their own cron session with a header/footer frame so the main conversation's role alternation stays intact diff --git a/gateway/run_notifications.py b/gateway/run_notifications.py index 060786947e..e2a9baaed7 100644 --- a/gateway/run_notifications.py +++ b/gateway/run_notifications.py @@ -33,9 +33,22 @@ _IMAGE_EXTS = {'.jpg', '.jpeg', '.png', '.webp', '.gif'} _DURABLE_CLAIM_OPS = { "drop": ("drop_completion_delivery", "Could not drop durable completion claim"), "release": ("release_completion_delivery", "Could not release durable completion claim"), + "defer": ("defer_completion_delivery", "Could not defer unadmitted completion claim"), + "complete": ("complete_completion_delivery", "Could not acknowledge durable completion claim"), } +def _raw_process_event_session_id(evt: dict) -> str: + """Recognize API routes, not malformed structured or partial messaging routes.""" + session_key = str(evt.get("session_key") or "").strip() + platform = str(evt.get("platform") or "").strip().lower() + if session_key.startswith("agent:") or platform not in {"", "api_server"}: + return "" + if not platform and any(evt.get(field) for field in ("chat_id", "chat_type", "thread_id")): + return "" + return str(evt.get("origin_session_id") or session_key or "").strip() + + class GatewayNotificationsMixin: """Process/completion/update notifications, media delivery and async-delegation delivery for GatewayRunner.""" @@ -781,6 +794,10 @@ class GatewayNotificationsMixin: chat_type = str(evt.get("chat_type") or derived.get("chat_type") or "").strip().lower() chat_id = str(evt.get("chat_id") or derived.get("chat_id") or "").strip() if not platform_name or not chat_type or not chat_id: + # Raw API keys legitimately have no messaging source. Resolve persisted + # origins first, then leave this recognized route to the API dispatcher. + if _raw_process_event_session_id(evt): + return None logger.warning( "Synthetic event source unresolvable: " "session_key=%r platform=%r chat_type=%r chat_id=%r evt_type=%s", @@ -896,27 +913,25 @@ class GatewayNotificationsMixin: return _transport.adapter return next((a for p, a in adapters.items() if p.value == platform_name), None) - async def _inject_watch_notification(self, synth_text: str, evt: dict) -> Optional[bool]: + async def _inject_watch_notification( + self, synth_text: str, evt: dict, *, raise_not_accepted: bool = False, + ) -> Optional[bool]: """Inject a watch/completion notification as a synthetic message event. Routing comes from the queued event, never the active foreground message. Returns ``True`` on adapter acceptance, ``False`` on retryable adapter failure, ``None`` with no gateway route. Not transactional: a crash after acceptance can replay (at-least-once). """ - from gateway.run import _parse_session_key - from gateway.wake import adapter_supports_push + from gateway.wake import WakeNotAccepted, adapter_supports_push, admit_internal_event source = await asyncio.to_thread(self._build_process_event_source, evt) if not source: # API-server sessions bind the RAW X-Hermes-Session-Id key, not a structured ``agent:...`` key. - raw_sid = str(evt.get("origin_session_id") or "").strip() - _sk = str(evt.get("session_key") or "").strip() - if not raw_sid and _sk and _parse_session_key(_sk) is None: - raw_sid = _sk + raw_sid = _raw_process_event_session_id(evt) if raw_sid: adapter = self.adapters.get(Platform.API_SERVER) if adapter is not None and not adapter_supports_push(adapter): return await self._self_post_api_server(adapter, synth_text, raw_sid, evt) - logger.warning( + logger.debug( "Deferring watch notification for raw session %s: no api_server adapter to self-post through", raw_sid, ) @@ -937,6 +952,9 @@ class GatewayNotificationsMixin: return await self._self_post_api_server(adapter, synth_text, raw_sid, evt) try: metadata = {} + session_key = str(evt.get("session_key") or "").strip() + if session_key.startswith("agent:"): + metadata["gateway_session_key"] = session_key parent_session_id = str(evt.get("parent_session_id") or "").strip() if parent_session_id: metadata["gateway_session_id"] = parent_session_id @@ -953,8 +971,13 @@ class GatewayNotificationsMixin: _prime = getattr(adapter, "prime_routing_cache", None) if callable(_prime): _prime(synth_event) - await adapter.handle_message(synth_event) + await admit_internal_event(adapter, synth_event) return True + except WakeNotAccepted: + # Durable callers refund the claim; ordinary watch callers just requeue. + if raise_not_accepted: + raise + return False except Exception as e: logger.error("Watch notification injection error: %s", e) return False @@ -1051,11 +1074,10 @@ class GatewayNotificationsMixin: import tools.async_delegation as _ad getattr(_ad, fn_name)(delegation_id, claim_id) except Exception: - logger.debug(fail_msg, exc_info=True) + logger.log(logging.WARNING if kind == "complete" else logging.DEBUG, fail_msg, exc_info=True) async def _completion_delivery_ready(self, evt: dict) -> bool: """Unavailable owners/transports must not spend a durable delivery attempt.""" - from gateway.run import _parse_session_key from gateway.wake import adapter_supports_push parent_session_id = str(evt.get("parent_session_id") or "").strip() @@ -1069,10 +1091,7 @@ class GatewayNotificationsMixin: platform = source.platform.value if hasattr(source.platform, "value") else str(source.platform) adapter = self._resolve_injection_adapter(platform, source) else: - session_key = str(evt.get("session_key") or "").strip() - raw_sid = str(evt.get("origin_session_id") or "").strip() - if not raw_sid and session_key and _parse_session_key(session_key) is None: - raw_sid = session_key + raw_sid = _raw_process_event_session_id(evt) adapter = self.adapters.get(Platform.API_SERVER) if raw_sid else None if adapter is not None and adapter_supports_push(adapter): return False @@ -1151,41 +1170,49 @@ class GatewayNotificationsMixin: claim.proceed, claim.early_result = False, False return claim - async def _deliver_completion_notification(self, synth_text: str, evt: dict) -> Optional[bool]: - """Deliver once per live gateway, or return False for a retry. + async def _deliver_completion_notification( + self, synth_text: str, evt: dict, *, sibling_claims=(), + ) -> Optional[bool]: + """Acknowledge one admitted batch, refund refusals, or release failed deliveries. - ``True``: adapter accepted; ``False``: injection failed, claim released for retry; ``None``: - another same-lifecycle caller owns/delivered it, or no route. No cross-process exactly-once. + True means adapter admission, not model execution; None means deduplicated or + terminal. False remains retryable. Claims are settled together for every sibling. """ + from gateway.wake import WakeNotAccepted identity = self._completion_delivery_identity(evt) - claim = await self._preflight_completion_delivery(evt) - if not claim.proceed: - return claim.early_result - if identity is not None and self._completion_identity_seen(identity, claim=True): - return None - accepted = False + claim = self._CompletionClaim() + accepted = identity_claimed = refused = False try: - injection_result = await self._inject_watch_notification(synth_text, evt) + claim = await self._preflight_completion_delivery(evt) + if not claim.proceed: + return claim.early_result + if identity is not None: + if self._completion_identity_seen(identity, claim=True): + return None + identity_claimed = True + injection_result = await self._inject_watch_notification(synth_text, evt, raise_not_accepted=True) if injection_result is not True: return injection_result accepted = True if identity is not None: with self._completion_delivery_lock: self._mark_completions_delivered_locked((identity,)) - # The durable row is the authoritative replay state — ack it after adapter acceptance. - if claim.claim_id: - try: - from tools.async_delegation import complete_completion_delivery - complete_completion_delivery(claim.delegation_id, claim.claim_id) - except Exception as exc: - logger.warning("Could not acknowledge durable async completion %s: %s", claim.delegation_id, exc) return True + except WakeNotAccepted: + refused = True + return False finally: - if identity is not None and not accepted: + if identity_claimed and not accepted: with self._completion_delivery_lock: self._completion_deliveries_inflight.discard(identity) - if claim.claim_id and not accepted: - self._settle_durable_claim("release", claim.delegation_id, claim.claim_id) + operation = "complete" if accepted else "defer" if refused else "release" + if claim.claim_id: + self._settle_durable_claim(operation, claim.delegation_id, claim.claim_id) + for sibling, claim_id in sibling_claims: + if claim_id: + self._settle_durable_claim(operation, sibling["delegation_id"], claim_id) + if accepted and sibling_claims: + self._record_coalesced_completion_siblings([event for event, _claim_id in sibling_claims]) @staticmethod def _event_route_key(evt: dict, fields: tuple[str, ...]) -> tuple[str, ...]: @@ -1335,14 +1362,6 @@ class GatewayNotificationsMixin: if parsed.get("thread_id"): evt["thread_id"] = parsed["thread_id"] - @staticmethod - def _settle_sibling_claims(siblings: list[tuple[dict, str]], fn, fail_msg: str) -> None: - for evt, claim_id in siblings: - try: - fn(evt, claim_id) - except Exception: - logger.debug(fail_msg, exc_info=True) - async def _deliver_async_delegation_group(self, group: list[dict]) -> Optional[bool]: """Deliver a same-session batch of async completions as ONE turn: the primary carries the consolidated text of every sibling THIS runner claimed (siblings owned elsewhere are excluded; @@ -1369,7 +1388,7 @@ class GatewayNotificationsMixin: for evt, _text in deliverable: if not await self._completion_delivery_ready(evt): return False - from tools.async_delegation import claim_event_delivery, complete_event_delivery, release_event_delivery + from tools.async_delegation import claim_event_delivery primary_evt, primary_text = deliverable[0] blocks = [primary_text] siblings: list[tuple[dict, str]] = [] @@ -1389,22 +1408,13 @@ class GatewayNotificationsMixin: "response. If a result does not change the current conclusion, absorb it silently.]" ) consolidated = "\n\n".join([header, *blocks]) - delivered: Optional[bool] = False - try: - delivered = await self._deliver_completion_notification(consolidated, primary_evt) - finally: - if delivered is True: - self._settle_sibling_claims( - siblings, complete_event_delivery, "Could not acknowledge coalesced durable completion", - ) - self._record_coalesced_completion_siblings([evt for evt, _claim_id in siblings]) - else: - # Not delivered: release every sibling claim so a retry or another consumer can take it. - self._settle_sibling_claims(siblings, release_event_delivery, "Could not release coalesced durable claim") - if delivered is None: - # Primary dropped/owned elsewhere: requeue just the siblings for the next tick. - for evt, _claim_id in siblings: - _pr.completion_queue.put(evt) + delivered = await self._deliver_completion_notification( + consolidated, primary_evt, sibling_claims=siblings, + ) + if delivered is None: + # Primary dropped/owned elsewhere: retry the unadmitted siblings. + for evt, _claim_id in siblings: + _pr.completion_queue.put(evt) return delivered async def _async_delegation_watcher(self, interval: float = 2.0) -> None: diff --git a/tests/gateway/test_background_process_notifications.py b/tests/gateway/test_background_process_notifications.py index de86b5ae85..031b4af2d3 100644 --- a/tests/gateway/test_background_process_notifications.py +++ b/tests/gateway/test_background_process_notifications.py @@ -23,6 +23,15 @@ from gateway.run import GatewayRunner, _parse_session_key # Helpers # --------------------------------------------------------------------------- +class AdmittingHandler(AsyncMock): + """Fake transport whose successful insertion issues the production receipt.""" + + async def _execute_mock_call(self, event, *args, **kwargs): + result = await super()._execute_mock_call(event, *args, **kwargs) + event._gateway_accepted = True + return result + + class _FakeRegistry: """Return pre-canned sessions, then None once exhausted.""" @@ -51,7 +60,7 @@ def _build_runner(monkeypatch, tmp_path, mode: str) -> GatewayRunner: monkeypatch.setattr(gateway_run, "_hermes_home", tmp_path) runner = GatewayRunner(GatewayConfig()) - adapter = SimpleNamespace(send=AsyncMock(), handle_message=AsyncMock()) + adapter = SimpleNamespace(send=AsyncMock(), handle_message=AdmittingHandler()) runner.adapters[Platform.TELEGRAM] = adapter return runner @@ -524,7 +533,7 @@ async def test_inject_watch_notification_raw_session_key_self_posts(monkeypatch, runner = _build_runner(monkeypatch, tmp_path, "all") api_adapter = SimpleNamespace( supports_async_delivery=False, - handle_message=AsyncMock(), + handle_message=AdmittingHandler(), _host="127.0.0.1", _port=8642, _api_key="k", _model_name="m", ) runner.adapters[Platform.API_SERVER] = api_adapter @@ -557,7 +566,7 @@ async def test_inject_watch_notification_origin_session_id_wins(monkeypatch, tmp runner = _build_runner(monkeypatch, tmp_path, "all") api_adapter = SimpleNamespace( supports_async_delivery=False, - handle_message=AsyncMock(), + handle_message=AdmittingHandler(), _host="127.0.0.1", _port=8642, _api_key="k", _model_name="m", ) runner.adapters[Platform.API_SERVER] = api_adapter @@ -591,7 +600,7 @@ async def test_async_delegation_apiserver_persists_delivery_not_self_post( runner = _build_runner(monkeypatch, tmp_path, "all") api_adapter = SimpleNamespace( supports_async_delivery=False, - handle_message=AsyncMock(), + handle_message=AdmittingHandler(), _host="127.0.0.1", _port=8642, _api_key="k", _model_name="m", ) runner.adapters[Platform.API_SERVER] = api_adapter @@ -639,7 +648,7 @@ async def test_async_delegation_apiserver_persist_failure_is_retryable( runner = _build_runner(monkeypatch, tmp_path, "all") api_adapter = SimpleNamespace( supports_async_delivery=False, - handle_message=AsyncMock(), + handle_message=AdmittingHandler(), _host="127.0.0.1", _port=8642, _api_key="k", _model_name="m", ) runner.adapters[Platform.API_SERVER] = api_adapter diff --git a/tests/gateway/test_completion_admission.py b/tests/gateway/test_completion_admission.py new file mode 100644 index 0000000000..e42afe7fda --- /dev/null +++ b/tests/gateway/test_completion_admission.py @@ -0,0 +1,122 @@ +"""Real adapter admission is the completion acknowledgement boundary.""" +import asyncio +import logging +import time + +import pytest + +from gateway.config import GatewayConfig, Platform, PlatformConfig +from gateway.platforms.event import MessageEvent +from gateway.run import GatewayRunner +from gateway.session import SessionSource, build_session_key +from hermes_state import SessionDB +from plugins.platforms.discord.adapter import DiscordAdapter +from tools import async_delegation as delegation + + +def pending(key, name): + evt = {"type": "async_delegation", "session_key": key, "delegation_id": name, + "summary": name, "status": "completed", "dispatched_at": time.time()} + delegation._persist_dispatch(evt) + delegation._persist_completion(evt, {"status": "completed", "summary": name}) + return evt + + +async def drain(adapter): + while adapter._background_tasks: + await asyncio.gather(*list(adapter._background_tasks)) + + +@pytest.mark.asyncio +async def test_completion_ack_requires_admission_and_replay_never_repeats(tmp_path): + runner = GatewayRunner(GatewayConfig()) + adapter = DiscordAdapter(PlatformConfig(enabled=True, typing_indicator=False)) + runner.adapters = {Platform.DISCORD: adapter} + source = SessionSource(platform=Platform.DISCORD, chat_type="dm", chat_id="42", user_id="42") + key = build_session_key(source) + events = [pending(key, f"admission-{i}") for i in range(2)] + received = [] + release, started = asyncio.Event(), asyncio.Event() + + async def handler(event): + received.append(event.text) + started.set() + await release.wait() + if key not in adapter._pending_messages: + queued = runner._promote_queued_event(key, adapter, None) + if queued is not None: + adapter._pending_messages[key] = queued + + try: + # Missing handler must not acknowledge either durable sibling. + for _ in range(10): + assert await runner._deliver_async_delegation_group(events) is False + for event in events: + row = delegation.get_durable_delegation(event["delegation_id"]) + assert (row["delivery_state"], row["delivery_attempts"]) == ("pending", 0) + assert not runner._completion_deliveries_delivered + adapter.set_message_handler(handler) + await adapter.handle_message(MessageEvent(text="human-active", source=source)) + await asyncio.wait_for(started.wait(), 2) + await adapter.handle_message(MessageEvent(text="human-pending", source=source)) + adapter.set_busy_session_handler(runner._handle_active_session_busy_message) + runner._BUSY_QUEUE_MAX_PENDING = 1 + for _ in range(10): + assert await runner._deliver_async_delegation_group(events) is False + assert adapter._pending_messages[key].text == "human-pending" + assert not runner._completion_deliveries_delivered + for event in events: + row = delegation.get_durable_delegation(event["delegation_id"]) + assert (row["delivery_state"], row["delivery_attempts"]) == ("pending", 0) + # An explicitly mismatched adapter key must fail closed too. + wrong = dict(events[0], session_key="agent:main:discord:dm:other", + platform="discord", chat_type="dm", chat_id="42") + assert await runner._inject_watch_notification("wrong-route", wrong) is False + runner._BUSY_QUEUE_MAX_PENDING = 4 + assert await runner._deliver_async_delegation_group(events) is True + assert await runner._deliver_async_delegation_group(events) is None + release.set() + await drain(adapter) + assert received[:2] == ["human-active", "human-pending"] + assert len(received) == 3 and all(event["summary"] in received[-1] for event in events) + for event in events: + assert delegation.get_durable_delegation(event["delegation_id"])["delivery_state"] == "delivered" + idle = pending(key, "idle-admitted") + assert await runner._deliver_async_delegation_group([idle]) is True + await drain(adapter) + assert len(received) == 4 and "idle-admitted" in received[-1] + finally: + release.set() + await drain(adapter) + await runner._cancel_process_completion_batch_tasks() + + +@pytest.mark.asyncio +async def test_unavailable_raw_route_is_quiet_without_hiding_invalid_routes(tmp_path, caplog): + runner = GatewayRunner(GatewayConfig()) + runner.adapters = {} + evt = pending("opaque-client-session", "raw-admission") + caplog.set_level(logging.WARNING, logger="gateway.run") + for _ in range(3): + assert await runner._deliver_async_delegation_group([evt]) is False + assert not caplog.records + row = delegation.get_durable_delegation(evt["delegation_id"]) + assert (row["delivery_state"], row["delivery_attempts"]) == ("pending", 0) + assert await runner._inject_watch_notification("watch", {"type": "watch_match", "session_key": "agent:broken"}) is None + assert any("unresolvable" in record.message for record in caplog.records) + # API recovery writes only the delivery row, never starts a model turn. + from gateway.platforms.api_server import APIServerAdapter + api = APIServerAdapter(PlatformConfig()) + db = SessionDB(tmp_path / "api.db") + db.create_session(evt["session_key"], "api_server") + api._ensure_session_db = lambda: db + runner.adapters = {Platform.API_SERVER: api} + try: + caplog.clear() + assert await runner._deliver_async_delegation_group([evt]) is True + assert await runner._deliver_async_delegation_group([evt]) is None + rows = db.get_messages(evt["session_key"]) + assert len(rows) == 1 and rows[0]["display_kind"] == "async_delegation_complete" + assert not api._background_tasks and not caplog.records + finally: + db.close() diff --git a/tests/gateway/test_completion_delivery.py b/tests/gateway/test_completion_delivery.py index 107b309604..4cb8f32d48 100644 --- a/tests/gateway/test_completion_delivery.py +++ b/tests/gateway/test_completion_delivery.py @@ -21,6 +21,15 @@ from gateway.session import SessionSource from tools.process_registry import ProcessRegistry, ProcessSession +class AdmittingHandler(AsyncMock): + """Fake transport whose successful insertion issues the production receipt.""" + + async def _execute_mock_call(self, event, *args, **kwargs): + result = await super()._execute_mock_call(event, *args, **kwargs) + event._gateway_accepted = True + return result + + @pytest.fixture(autouse=True) def isolated_registry(tmp_path, monkeypatch): """Any current/future durable compatibility path must stay in tmp state.""" @@ -104,7 +113,7 @@ def test_duplicate_async_queue_replay_injects_once(monkeypatch, isolated_registr isolated.put(dict(_async_event())) isolated.put(dict(_async_event())) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=2) @@ -122,7 +131,7 @@ def test_unroutable_async_event_remains_retryable( event["session_key"] = "20260711_unparseable_ui_session" isolated.put(event) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=2) @@ -141,7 +150,7 @@ def test_concurrent_claims_share_the_same_narrow_delivery_seam(): entered.set() await release.wait() - adapter = SimpleNamespace(handle_message=AsyncMock(side_effect=_blocked_injection)) + adapter = SimpleNamespace(handle_message=AdmittingHandler(side_effect=_blocked_injection)) runner = _runner(adapter) event = _async_event() text = "completion" @@ -166,7 +175,7 @@ def test_failed_async_injection_is_retried_and_only_success_is_acked( isolated.put(_async_event()) adapter = SimpleNamespace( - handle_message=AsyncMock(side_effect=[RuntimeError("temporary"), None]) + handle_message=AdmittingHandler(side_effect=[RuntimeError("temporary"), None]) ) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=3) @@ -227,7 +236,7 @@ def test_explicit_kill_returns_output_before_consuming_notification(monkeypatch) assert result["output"] == "important terminal output\n" assert registry.is_completion_consumed(session.id) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) async def _instant_sleep(*_a, **_kw): @@ -297,7 +306,7 @@ def test_autonomous_completion_redacts_real_command_and_output_secrets(monkeypat monkeypatch.setattr(pr_module, "process_registry", registry) monkeypatch.setattr(redact_module, "_REDACT_ENABLED", True) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) async def _instant_sleep(*_a, **_kw): @@ -348,7 +357,7 @@ def test_concurrent_process_watchers_coalesce_one_session_completion_turn(monkey }) monkeypatch.setattr(pr_module, "process_registry", registry) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) async def _exercise(): @@ -379,7 +388,7 @@ def test_completion_arriving_during_batch_delivery_schedules_next_flush(): first_delivery_entered.set() await release_first_delivery.wait() - adapter = SimpleNamespace(handle_message=AsyncMock(side_effect=_deliver)) + adapter = SimpleNamespace(handle_message=AdmittingHandler(side_effect=_deliver)) runner = _runner(adapter) async def _exercise(): @@ -402,7 +411,7 @@ def test_completion_arriving_during_batch_delivery_schedules_next_flush(): def test_completion_batches_do_not_cross_conversation_routes(): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) first = _completion_event(started_at=1.0, session_id="proc_route_a") @@ -429,7 +438,7 @@ def test_failed_coalesced_delivery_retries_all_entries(): if attempts == 1: raise RuntimeError("temporary adapter failure") - adapter = SimpleNamespace(handle_message=AsyncMock(side_effect=_deliver)) + adapter = SimpleNamespace(handle_message=AdmittingHandler(side_effect=_deliver)) runner = _runner(adapter) events = [ _completion_event(started_at=float(index), session_id=f"proc_retry_{index}") @@ -451,7 +460,7 @@ def test_failed_coalesced_delivery_retries_all_entries(): def test_coalesced_success_records_every_completion_identity(): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) events = [ _completion_event(started_at=float(index), session_id=f"proc_ledger_{index}") @@ -543,7 +552,7 @@ def test_coalesced_format_redacts_before_truncating_output(monkeypatch): def test_duplicate_primary_does_not_discard_fresh_batch_sibling(): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) duplicate = _completion_event(started_at=1.0, session_id="proc_duplicate") fresh = _completion_event(started_at=2.0, session_id="proc_fresh") @@ -563,7 +572,7 @@ def test_duplicate_primary_does_not_discard_fresh_batch_sibling(): def test_batch_format_failure_resolves_waiters_for_retry(monkeypatch): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) monkeypatch.setattr( runner, @@ -587,7 +596,7 @@ def test_batch_format_failure_resolves_waiters_for_retry(monkeypatch): def test_shutdown_cancels_batch_during_window_and_settles_waiter_for_retry(): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) sleep_entered = asyncio.Event() release_sleep = asyncio.Event() @@ -629,7 +638,7 @@ def test_shutdown_cancels_blocked_batch_delivery_and_keeps_it_retryable(): delivery_entered.set() await asyncio.Event().wait() - adapter = SimpleNamespace(handle_message=AsyncMock(side_effect=_blocked_delivery)) + adapter = SimpleNamespace(handle_message=AdmittingHandler(side_effect=_blocked_delivery)) runner = _runner(adapter) runner._completion_notification_batch_window = 0 event = _completion_event(started_at=1.0, session_id="proc_cancel_delivery") @@ -655,7 +664,7 @@ def test_shutdown_cancels_blocked_batch_delivery_and_keeps_it_retryable(): def test_completion_enqueue_stays_retryable_after_shutdown_starts(): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) async def _exercise(): @@ -672,7 +681,7 @@ def test_completion_enqueue_stays_retryable_after_shutdown_starts(): def test_successful_batch_releases_all_lifecycle_task_references(): - adapter = SimpleNamespace(handle_message=AsyncMock(return_value=None)) + adapter = SimpleNamespace(handle_message=AdmittingHandler(return_value=None)) runner = _runner(adapter) runner._completion_notification_batch_window = 0 @@ -697,7 +706,7 @@ def test_shutdown_cancels_overlapping_flushes_for_same_route(): delivery_entered.set() await asyncio.Event().wait() - adapter = SimpleNamespace(handle_message=AsyncMock(side_effect=_blocked_delivery)) + adapter = SimpleNamespace(handle_message=AdmittingHandler(side_effect=_blocked_delivery)) runner = _runner(adapter) runner._completion_notification_batch_window = 0 first_event = _completion_event(started_at=1.0, session_id="proc_old_flush") @@ -764,7 +773,7 @@ def test_same_tick_async_batch_coalesces_into_one_turn_and_acks_all_rows( _persist_pending_completion(event) isolated.put(dict(event)) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=2) @@ -792,7 +801,7 @@ def test_same_tick_async_events_for_different_sessions_do_not_coalesce( "deleg_route_b", session_key="agent:main:telegram:dm:99999:678", )) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=2) @@ -813,7 +822,7 @@ def test_single_async_event_latency_and_text_are_unchanged( monkeypatch.setattr(isolated_registry, "completion_queue", isolated) isolated.put(_distinct_async_event("deleg_single")) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=2) @@ -839,7 +848,7 @@ def test_failed_coalesced_async_batch_releases_claims_and_retries( isolated.put(dict(event)) adapter = SimpleNamespace( - handle_message=AsyncMock(side_effect=[RuntimeError("temporary"), None]) + handle_message=AdmittingHandler(side_effect=[RuntimeError("temporary"), None]) ) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=3) @@ -872,7 +881,7 @@ def test_sibling_claimed_by_other_consumer_is_not_double_delivered( events[1]["delegation_id"], "other-consumer:claim", ) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) _stop_after_sleeps(monkeypatch, runner, count=2) @@ -902,7 +911,7 @@ def test_unavailable_delivery_preserves_budget_across_restarts(tmp_path, unavail event["parent_session_id"] = "parent-session" _persist_pending_completion(event) - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) api = SimpleNamespace(supports_async_delivery=False, _ensure_session_db=lambda: None) for _restart in range(10): runner = _runner(adapter) @@ -940,8 +949,8 @@ def test_unavailable_delivery_preserves_budget_across_restarts(tmp_path, unavail @pytest.mark.parametrize("available", [False, True]) def test_completion_profile_transport_never_falls_back(monkeypatch, isolated_registry, available): - primary = SimpleNamespace(handle_message=AsyncMock(), send=AsyncMock()) - secondary = SimpleNamespace(handle_message=AsyncMock(), send=AsyncMock()) + primary = SimpleNamespace(handle_message=AdmittingHandler(), send=AsyncMock()) + secondary = SimpleNamespace(handle_message=AdmittingHandler(), send=AsyncMock()) runner = _runner(primary) runner._profile_adapters = {"research": {Platform.TELEGRAM: secondary}} if available else {} evt = dict(_completion_event(started_at=1), session_key="agent:research:telegram:dm:12345") @@ -959,7 +968,7 @@ def test_completion_profile_transport_never_falls_back(monkeypatch, isolated_reg @pytest.mark.parametrize("event_type", ["watch_match", "watch_disabled"]) @pytest.mark.parametrize("mode", ["concise", "off"]) def test_idle_watch_drain_respects_notify_mode(monkeypatch, isolated_registry, event_type, mode): - adapter = SimpleNamespace(handle_message=AsyncMock()) + adapter = SimpleNamespace(handle_message=AdmittingHandler()) runner = _runner(adapter) runner._load_background_notifications_mode = lambda: mode evt = dict(_completion_event(started_at=1), type=event_type, @@ -972,7 +981,7 @@ def test_idle_watch_drain_respects_notify_mode(monkeypatch, isolated_registry, e def test_watch_drain_retries_transport_failure(monkeypatch, isolated_registry): - adapter = SimpleNamespace(handle_message=AsyncMock(side_effect=[RuntimeError("offline"), None])) + adapter = SimpleNamespace(handle_message=AdmittingHandler(side_effect=[RuntimeError("offline"), None])) runner = _runner(adapter) runner._load_background_notifications_mode = lambda: "concise" evt = dict(_completion_event(started_at=1), type="watch_match", pattern="READY", output="READY") diff --git a/tests/gateway/test_completion_session_boundary.py b/tests/gateway/test_completion_session_boundary.py index b01b0f6c1e..a7dffa99af 100644 --- a/tests/gateway/test_completion_session_boundary.py +++ b/tests/gateway/test_completion_session_boundary.py @@ -51,6 +51,9 @@ class _SessionDB: def _runner(adapter, *, session_db=...): + def admit(event): + event._gateway_accepted = True + adapter.handle_message.side_effect = admit runner = object.__new__(GatewayRunner) runner._running = True runner.adapters = {Platform.TELEGRAM: adapter} diff --git a/tests/gateway/test_relay_completion_injection_routing.py b/tests/gateway/test_relay_completion_injection_routing.py index 327e7377c1..96418bd04a 100644 --- a/tests/gateway/test_relay_completion_injection_routing.py +++ b/tests/gateway/test_relay_completion_injection_routing.py @@ -32,7 +32,11 @@ class _RelayAdapter: def __init__(self): self.handled = [] - self.handle_message = AsyncMock(side_effect=self.handled.append) + self.handle_message = AsyncMock(side_effect=self._admit) + + def _admit(self, event): + self.handled.append(event) + event._gateway_accepted = True def fronts_platform(self, platform): return platform == Platform.SLACK diff --git a/tests/gateway/test_relay_injection_egress_priming.py b/tests/gateway/test_relay_injection_egress_priming.py index 3200c0b3b5..2e33ccefae 100644 --- a/tests/gateway/test_relay_injection_egress_priming.py +++ b/tests/gateway/test_relay_injection_egress_priming.py @@ -86,6 +86,7 @@ async def test_injection_path_primes_before_handle_message(): async def handle_message(self, event): calls.append(("handle", getattr(event.source, "chat_id", None))) + event._gateway_accepted = True runner = object.__new__(GatewayRunner) runner._running = True diff --git a/tools/async_delegation.py b/tools/async_delegation.py index 7a46d0ca7c..27fd4523fd 100644 --- a/tools/async_delegation.py +++ b/tools/async_delegation.py @@ -402,6 +402,15 @@ def release_completion_delivery(delegation_id: str, claim_id: str) -> bool: return cur.rowcount == 1 +def defer_completion_delivery(delegation_id: str, claim_id: str) -> bool: + """Return an unadmitted completion to pending without spending a delivery attempt.""" + return _update_delivery("""UPDATE async_delegations SET delivery_claim=NULL, + delivery_claimed_at=NULL, delivery_attempts=MAX(0, delivery_attempts-1), + updated_at=? + WHERE delegation_id=? AND delivery_state='pending' AND delivery_claim=?""", + (time.time(), delegation_id, claim_id)) + + def drop_completion_delivery(delegation_id: str, claim_id: str) -> bool: """Terminally drop a claimed completion whose target is permanently gone (the spawning session ended at an explicit user boundary such as /new or reset). diff --git a/website/docs/user-guide/features/delegation.md b/website/docs/user-guide/features/delegation.md index 33c32ad69d..dbe5e21c08 100644 --- a/website/docs/user-guide/features/delegation.md +++ b/website/docs/user-guide/features/delegation.md @@ -10,6 +10,21 @@ The `delegate_task` tool spawns child AIAgent instances with isolated context, i Top-level model calls run in the background automatically. Hermes returns a handle immediately so the conversation can continue, then posts the result back as a new message. An orchestrator subagent waits for its own workers so it can synthesize their results before returning. +## Completion delivery + +Messaging gateways acknowledge background completions only after their adapter actually +schedules the event or inserts it into the session's queue. Missing handlers, mismatched +session routes, and full queues leave the completion pending for retry; these admission +refusals do not consume the durable delivery-attempt budget. A successful admission suppresses +repeat delivery within the running gateway, but is not proof that a model turn or outbound +reply completed. Crash/restart delivery remains at least once, subject to the existing replay +age limit; actual transport failures retain their bounded retry policy. + +An unavailable API-server route stays pending without repeated missing-route warnings. +Malformed messaging routes still produce diagnostics. On the API server, an async delegation +completion adds a durable timeline delivery row only: the client owns the next model turn. +Setting background process notifications to `off` still drains pattern-watch events silently. + ## Background process lifetime Background terminal processes belong to the agent that starts them. Closing a child during delegation teardown terminates its remaining processes, including work started in earlier turns, without stopping processes owned by the parent or sibling agents. Sharing a terminal environment does not transfer process ownership.