fix: retain completions until explicit adapter admission
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
122
tests/gateway/test_completion_admission.py
Normal file
122
tests/gateway/test_completion_admission.py
Normal file
@@ -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()
|
||||
@@ -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")
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user