From eeada33fa1d83a734d836d020832f7673f4e4fa3 Mon Sep 17 00:00:00 2001 From: Phoebie Builder Date: Fri, 18 Sep 2026 05:20:43 +0000 Subject: [PATCH] review: retry the cap, adopt the compression tip, scope the in-process route to multiplex Independent review findings addressed: - adopt the live continuation tip (api_server_runs._resolve_live_session_id) before the in-process wake, exactly as the HTTP self-post does: a rotated (compressed) origin must be woken on the transcript that is actually live, not the retired parent slice (CompressionSessionClosedError otherwise); - retry a saturated concurrent-run cap with the HTTP path's own backoff (gateway.wake._RETRY_DELAYS_SECONDS) instead of failing on the first check, so a busy listener no longer burns one notifier failure per tick toward the 12-strike drop of a durable subscription; - keep the historical HTTP self-post on a standalone (non-multiplex) gateway: the in-process route is now chosen by a served-profile helper that requires multiplex_profiles, and the adapter is told the authorized profile explicitly instead of re-deriving it from HERMES_HOME; - skip the ownership store read for the default profile (the fallthrough already authorizes it) and run the delivery-time recheck off the event loop; - tests: standalone named-profile gateway, missing store, failed-connect boundary, adapter without in-process delivery, platform-wide api_server route vs the default profile's destinations, cap retry/exhaustion, compression tip, and the profile handed to the adapter (8 -> 15 cases). --- gateway/kanban_watchers_notifier.py | 32 +++- gateway/platforms/api_server.py | 53 ++++-- gateway/wake.py | 3 +- ...t_kanban_notifier_served_apiserver_wake.py | 173 +++++++++++++++++- 4 files changed, 228 insertions(+), 33 deletions(-) diff --git a/gateway/kanban_watchers_notifier.py b/gateway/kanban_watchers_notifier.py index 8d4faf6e90..512257e707 100644 --- a/gateway/kanban_watchers_notifier.py +++ b/gateway/kanban_watchers_notifier.py @@ -7,6 +7,7 @@ per-subscription delivery (``_KanbanNotification``) live here. from __future__ import annotations +import asyncio import contextlib import re from functools import partial @@ -210,8 +211,9 @@ def _adapter_for_subscription(runner: Any, platform: Any, sub: dict, owner_profi # profile_routes entry can anchor it — and a platform-wide api_server route would deny the # default profile's own api_server destinations. The shared listener mirrors /p// for # every served profile, so the owner's own session store is the proof: authorize exactly the - # session that lives in the served profile's state.db, never the platform. - if getattr(platform, "value", platform) == "api_server" \ + # session that lives in the served profile's state.db, never the platform. The default profile + # keeps the historical fallthrough below (no store read). + if profile != primary_profile and getattr(platform, "value", platform) == "api_server" \ and _session_owned_by_profile(config, profile, chat): return primary return primary if profile == primary_profile else None @@ -599,14 +601,27 @@ class _KanbanNotification: logger.info("kanban notifier: woke agent for %s on %s/%s profile=%s events=%s", self.task_id, self.platform_str, self.sub["chat_id"], self.sub_profile or "default", self.wake_kinds) + def _served_wake_profile(self) -> Optional[str]: + """The subscription's profile when THIS gateway is a multiplexer serving it, else ``None``. + + ``None`` keeps the historical path: a standalone ``hermes -p `` gateway owns its own + listener and key, so its api_server wakes keep using the HTTP self-post. + """ + if not self.sub_profile: + return None + if not getattr(getattr(self.runner, "config", None), "multiplex_profiles", False): + return None + return self.sub_profile + def _owner_scope(self): """Runtime scope of the subscription's profile under multiplex, else a no-op context.""" runner = self.runner - if not (self.sub_profile and getattr(getattr(runner, "config", None), "multiplex_profiles", False)): + served_profile = self._served_wake_profile() + if not served_profile: return contextlib.nullcontext() from gateway.run import _async_profile_runtime_scope from gateway.session import SessionSource - source = SessionSource(platform=self.plat, chat_id=self.sub["chat_id"], profile=self.sub_profile) + source = SessionSource(platform=self.plat, chat_id=self.sub["chat_id"], profile=served_profile) return _async_profile_runtime_scope(runner._resolve_profile_home_for_source(source)) async def wake(self) -> None: @@ -620,7 +635,7 @@ class _KanbanNotification: # unprefixed self-post would resume the session in the DEFAULT profile's store. async with self._owner_scope(): await deliver_wake(self.adapter, text=self.synth, session_id=self.session_key, - profile=self.sub_profile or None, + profile=self._served_wake_profile(), notification_category="diagnostic" if self.wake_diagnostic else "result") self._log_woke() return @@ -732,8 +747,11 @@ class _KanbanNotification: except ValueError: await self.advance() return - # Recheck the exact route after claiming: config/adapters can change between ticks. - adapter = _adapter_for_subscription(self.runner, self.plat, self.sub, self.sub_profile or None) + # Recheck the exact route after claiming: config/adapters can change between ticks. The + # recheck reads the served profile's session store for a stateless destination, so it runs + # off the event loop (the claim path already collects in a worker thread). + adapter = await asyncio.to_thread( + _adapter_for_subscription, self.runner, self.plat, self.sub, self.sub_profile or None) if adapter is None: logger.debug("kanban notifier: adapter %s disconnected before delivery for %s; rewinding claim", self.platform_str, self.task_id) diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index bb4f2877e5..d749119894 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -2775,7 +2775,8 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): return [] async def run_internal_session_turn(self, *, session_id: str, text: str, - notification_category: str = "result") -> None: + notification_category: str = "result", + profile: str = "") -> None: """Run one background wake turn against a raw session id IN-PROCESS (no HTTP, no API key). The HTTP wake self-post cannot serve a multiplexed *served* profile: ``/p//`` on @@ -2783,13 +2784,12 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): route-only profile legitimately does not have — while an unprefixed self-post would resume the session in the DEFAULT profile's store. ``gateway.wake`` therefore runs the turn here, inside the owner profile's runtime scope (the session DB, model resolution and tool policy - all follow the ambient scope). Raises on failure so the caller can rewind its cursor; the - concurrent-run cap defers the wake instead of bypassing it. + all follow the ambient scope), with ``profile`` naming the profile the caller proved owns + the session. Raises on failure so the caller can rewind its cursor; a saturated + concurrent-run cap is retried with the same backoff the HTTP self-post uses. """ - limited = self._concurrency_limited_response() - if limited is not None: - raise RuntimeError("internal wake deferred: the API server is at its concurrent-run cap") - profile = _api_request_profile.get() + from gateway.wake import _RETRY_DELAYS_SECONDS + profile = (profile or "").strip() or (_api_request_profile.get() or "") if not profile: # The HTTP paths bind the profile from the /p// prefix; in-process we bind it # ourselves so _run_agent's scope and the session DB resolve to the caller's profile. @@ -2797,17 +2797,36 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): active = (get_active_profile_name() or "").strip() profile = active if active and active != "default" else "" token = _api_request_profile.set(profile) if profile else None + attempts = 1 + len(_RETRY_DELAYS_SECONDS) + last_err: Optional[BaseException] = None try: - session, err = await self._get_existing_session_or_404(session_id) - if err is not None or not session: - raise RuntimeError( - f"internal wake target session {session_id!r} is not in the active profile store") - history = await self._conversation_history_for_session(session_id) - await self._run_agent( - user_message=text, conversation_history=history, session_id=session_id, - gateway_session_key=None, requested_runtime={}, route_source="global", - session_history_delivery="1", notification_category=notification_category, - ) + for attempt in range(attempts): + if attempt: + await asyncio.sleep(_RETRY_DELAYS_SECONDS[attempt - 1]) + # Transient: the cap clears on its own, exactly as the HTTP self-post's 429 does. + if self._concurrency_limited_response() is not None: + last_err = RuntimeError( + "internal wake deferred: the API server is at its concurrent-run cap") + logger.warning("%s; attempt %d/%d", last_err, attempt + 1, attempts) + continue + # #98619/#13437: adopt the live continuation tip first, the same canonical + # resolution the HTTP self-post consumes — a compressed origin must be woken on the + # transcript that is actually live, never the retired parent slice. + from gateway.platforms.api_server_runs import _resolve_live_session_id + resolved = await _resolve_live_session_id(self, session_id) + session, err = await self._get_existing_session_or_404(resolved) + if err is not None or not session: + raise RuntimeError( + f"internal wake target session {resolved!r} is not in the active profile store") + history = await self._conversation_history_for_session(resolved) + await self._run_agent( + user_message=text, conversation_history=history, session_id=resolved, + gateway_session_key=None, requested_runtime={}, route_source="global", + session_history_delivery="1", notification_category=notification_category, + ) + return + raise RuntimeError( + f"internal wake gave up for session {session_id} after {attempts} attempts: {last_err}") finally: if token is not None: _api_request_profile.reset(token) diff --git a/gateway/wake.py b/gateway/wake.py index 968db3758a..726a3ea2ac 100644 --- a/gateway/wake.py +++ b/gateway/wake.py @@ -166,7 +166,8 @@ async def _self_post_chat_completion(adapter: Any, *, text: str, session_id: str raise RuntimeError( f"wake self-post for served profile {profile!r} requires in-process session " "delivery; refusing to self-post as the default profile") - await in_process(session_id=session_id, text=text, notification_category=notification_category) + await in_process(session_id=session_id, text=text, profile=str(profile), + notification_category=notification_category) return import aiohttp host = str(getattr(adapter, "_host", "") or "127.0.0.1") diff --git a/tests/gateway/test_kanban_notifier_served_apiserver_wake.py b/tests/gateway/test_kanban_notifier_served_apiserver_wake.py index 8d135e1b0f..173e93756f 100644 --- a/tests/gateway/test_kanban_notifier_served_apiserver_wake.py +++ b/tests/gateway/test_kanban_notifier_served_apiserver_wake.py @@ -43,10 +43,13 @@ class RecordingApiServerAdapter: def __init__(self, *, fail_first: bool = False): self.turns = [] self.homes = [] + self.profiles = [] self._fail_first = fail_first - async def run_internal_session_turn(self, *, session_id, text, notification_category="result"): + async def run_internal_session_turn(self, *, session_id, text, notification_category="result", + profile=""): self.homes.append(str(get_hermes_home())) + self.profiles.append(profile) if self._fail_first: self._fail_first = False raise RuntimeError("simulated wake failure") @@ -189,8 +192,10 @@ def test_served_profile_api_server_subscription_wakes_in_process(served, monkeyp assert [turn["session_id"] for turn in adapter.turns] == [SESSION] assert WORKER_SESSION not in [turn["session_id"] for turn in adapter.turns] assert task in adapter.turns[0]["text"] and "done once" in adapter.turns[0]["text"] - # The wake ran under the OWNING profile's runtime scope, not the launch profile's. + # The wake ran under the OWNING profile's runtime scope, not the launch profile's, and the + # adapter was told which profile the authorization proved (no re-derivation from HERMES_HOME). assert adapter.homes == [str(served.builder)] + assert adapter.profiles == ["builder"] # No HTTP self-post: no shared-listener request, so no secondary API_SERVER_KEY is involved. assert _FakeHttpSession.calls == [] assert _unseen(task) == [] @@ -315,7 +320,7 @@ def test_internal_session_turn_binds_profile_and_targets_the_session(served, mon monkeypatch.setattr(adapter, "_run_agent", fake_run_agent) with _profile_runtime_scope(served.builder): - asyncio.run(adapter.run_internal_session_turn(session_id=SESSION, text="wake")) + asyncio.run(adapter.run_internal_session_turn(session_id=SESSION, text="wake", profile="builder")) assert seen["session_id"] == SESSION assert seen["user_message"] == "wake" @@ -328,11 +333,163 @@ def test_internal_session_turn_binds_profile_and_targets_the_session(served, mon # A session the active profile's store does not own fails closed. with _profile_runtime_scope(served.builder): with pytest.raises(RuntimeError): - asyncio.run(adapter.run_internal_session_turn(session_id="not-a-session", text="wake")) + asyncio.run(adapter.run_internal_session_turn(session_id="not-a-session", text="wake", + profile="builder")) - # The concurrent-run cap defers the wake instead of bypassing it (caller rewinds the cursor). - adapter._max_concurrent_runs = 1 - adapter._inflight_agent_runs = 1 + +def test_internal_session_turn_adopts_the_compression_tip(served, monkeypatch): + """A rotated (compressed) origin wakes on the live continuation, not the retired parent.""" + from gateway.platforms.api_server import APIServerAdapter + + parent, tip = "20260918_010000_parent", "20260918_020000_tip" + db = SessionDB(served.builder / "state.db") + try: + db.create_session(parent, source="webui", profile_name="builder") + db.append_message(parent, "user", "old turn") + db.end_session(parent, "compression") + db.create_session(tip, source="webui", profile_name="builder", parent_session_id=parent) + db.append_message(tip, "user", "live turn") + assert db.resolve_resume_session_id(parent) == tip, "fixture must produce a live tip" + finally: + db.close() + + adapter = APIServerAdapter(PlatformConfig(enabled=True)) + seen = {} + + async def fake_run_agent(**kwargs): + seen.update(kwargs) + return {}, {} + + monkeypatch.setattr(adapter, "_run_agent", fake_run_agent) + with _profile_runtime_scope(served.builder): + asyncio.run(adapter.run_internal_session_turn(session_id=parent, text="wake", profile="builder")) + + assert seen["session_id"] == tip + assert "live turn" in str(seen["conversation_history"]) + + +def test_internal_session_turn_retries_a_saturated_cap(served, monkeypatch): + """The concurrent-run cap is transient: back off and retry (as the HTTP 429 path does).""" + from gateway.platforms.api_server import APIServerAdapter + + _own_session(served.builder, SESSION, "builder") + adapter = APIServerAdapter(PlatformConfig(enabled=True)) + seen, sleeps = {}, [] + + async def fake_run_agent(**kwargs): + seen.update(kwargs) + return {}, {} + + calls = {"n": 0} + + def limited_once(): + calls["n"] += 1 + return object() if calls["n"] == 1 else None + + async def fake_sleep(delay): + sleeps.append(delay) + + monkeypatch.setattr(adapter, "_run_agent", fake_run_agent) + monkeypatch.setattr(adapter, "_concurrency_limited_response", limited_once) + monkeypatch.setattr(asyncio, "sleep", fake_sleep) + + with _profile_runtime_scope(served.builder): + asyncio.run(adapter.run_internal_session_turn(session_id=SESSION, text="wake", profile="builder")) + assert seen["session_id"] == SESSION + assert sleeps == [2.0] # first backoff step, then the retry succeeds + + # Exhausting the attempts raises, so the caller rewinds instead of silently dropping the event. + sleeps.clear() + monkeypatch.setattr(adapter, "_concurrency_limited_response", lambda: object()) with _profile_runtime_scope(served.builder): with pytest.raises(RuntimeError): - asyncio.run(adapter.run_internal_session_turn(session_id=SESSION, text="wake")) + asyncio.run(adapter.run_internal_session_turn(session_id=SESSION, text="wake", profile="builder")) + assert sleeps == [2.0, 5.0, 10.0] + + +def test_standalone_named_profile_gateway_keeps_the_http_self_post(served, monkeypatch): + """Without a multiplexer the profile owns its own listener/key: the self-post is unchanged.""" + import aiohttp + + _FakeHttpSession.calls = [] + monkeypatch.setattr(aiohttp, "ClientSession", _FakeHttpSession) + kb.init_db() + adapter = RecordingApiServerAdapter() + adapter._api_key, adapter._host, adapter._port, adapter._model_name = "k" * 20, "127.0.0.1", 8642, "hermes" + runner = _make_runner(served, adapter=adapter) + runner.config = SimpleNamespace(multiplex_profiles=False, profile_routes=[]) + # A standalone ``hermes -p builder`` gateway: no multiplexer, no secondary adapter map, and the + # profile IS the primary profile of this process (so it owns its own listener and key). + runner._primary_profile_name = "builder" + runner._kanban_notifier_profile = "builder" + runner._profile_adapters = {} + task = _subscription() # notifier_profile="builder" on a non-multiplex gateway + + asyncio.run(_run_one_notifier_tick(monkeypatch, runner)) + + assert adapter.turns == [] + assert len(_FakeHttpSession.calls) == 1 + assert _FakeHttpSession.calls[0]["url"].endswith("/v1/chat/completions") + assert _unseen(task) == [] + + +def test_missing_session_store_fails_closed(served): + """No state.db at all for the served profile is not ownership.""" + runner = _make_runner(served) + assert not (served.builder / "state.db").exists() + assert _adapter_for_subscription(runner, Platform.API_SERVER, _api_sub(), "builder") is None + + +def test_failed_connect_profile_does_not_borrow_but_still_wakes_in_process(served): + """A profile whose own bot failed to connect still owns its session wake (no adapter is used). + + The in-process turn reads no credential and sends through no transport, so the transport + boundary documented in ``authz_mixin._is_shared_bot_satellite`` (a failed/reconnecting bot is + still that profile's credential) does not gate it. A profile that DID connect its own adapter + keeps the hard boundary: it never falls back to the primary listener. + """ + _own_session(served.builder, SESSION, "builder") + failed = _make_runner(served) + failed._profile_failed_platforms = {"builder": {Platform.DISCORD}} + assert _adapter_for_subscription(failed, Platform.API_SERVER, _api_sub(), "builder") \ + is failed.adapters[Platform.API_SERVER] + + connected = _make_runner(served, builder_adapters={Platform.DISCORD: object()}) + assert _adapter_for_subscription(connected, Platform.API_SERVER, _api_sub(), "builder") is None + + +def test_adapter_without_in_process_delivery_fails_closed(served, monkeypatch): + """A non-push adapter that cannot run in-process must never self-post as the default profile.""" + import aiohttp + + class InProcesslessAdapter: + supports_async_delivery = False + + async def send(self, chat_id, text, metadata=None): + from gateway.platforms.base import SendResult + return SendResult(success=False, error="stateless") + + _own_session(served.builder, SESSION, "builder") + _FakeHttpSession.calls = [] + monkeypatch.setattr(aiohttp, "ClientSession", _FakeHttpSession) + kb.init_db() + runner = _make_runner(served, adapter=InProcesslessAdapter()) + task = _subscription() + + asyncio.run(_run_one_notifier_tick(monkeypatch, runner)) + + assert _FakeHttpSession.calls == [] + assert _unseen(task) # denied and retryable, never delivered as the wrong profile + + +def test_platform_wide_api_server_route_denies_the_default_profiles_destinations(served): + """Why the ownership rule exists: a platform-wide api_server route is not a narrow fix.""" + _own_session(served.builder, SESSION, "builder") + _own_session(served.root, "default-session", "default") + runner = _make_runner(served, routes=[ + {"name": "api-builder", "platform": "api_server", "profile": "builder"}]) + assert _adapter_for_subscription(runner, Platform.API_SERVER, _api_sub(), "builder") \ + is runner.adapters[Platform.API_SERVER] + # The same route matches the default profile's own api_server destination and denies it. + assert _adapter_for_subscription(runner, Platform.API_SERVER, + _api_sub(chat_id="default-session"), "default") is None