refactor(gateway): compact run_* docstrings/comments by hand; AST-neutral bracket packing to 110 cols
This commit is contained in:
+37
-70
@@ -38,11 +38,9 @@ class GatewayAdapterLifecycleMixin:
|
||||
|
||||
@staticmethod
|
||||
async def _wait_or_detach(task: "asyncio.Future", timeout: float) -> bool:
|
||||
"""Wait up to ``timeout`` for ``task``; on deadline (or our own cancellation) detach it.
|
||||
|
||||
Not ``asyncio.wait_for``: that WAITS for the cancelled child to exit, so a connect()/close()
|
||||
that swallows ``CancelledError`` blocks recovery forever. Returns True if it finished in time.
|
||||
"""
|
||||
"""Wait up to ``timeout`` for ``task``; on deadline (or our own cancellation) detach it. Not
|
||||
``asyncio.wait_for``: that WAITS for the cancelled child, so a connect()/close() swallowing
|
||||
``CancelledError`` blocks recovery forever. True if it finished in time."""
|
||||
try:
|
||||
done, _pending = await asyncio.wait({task}, timeout=timeout)
|
||||
except asyncio.CancelledError:
|
||||
@@ -67,11 +65,8 @@ class GatewayAdapterLifecycleMixin:
|
||||
return True
|
||||
|
||||
async def _safe_adapter_disconnect(self, adapter, platform) -> None:
|
||||
"""Call adapter.disconnect() defensively (bounded, never raises).
|
||||
|
||||
After a failed connect() partial resources (ClientSession, poll tasks, subprocesses) would
|
||||
otherwise leak; must tolerate partial-init state.
|
||||
"""
|
||||
"""Call adapter.disconnect() defensively (bounded, never raises, tolerates partial-init state):
|
||||
after a failed connect() partial resources (ClientSession, poll tasks, subprocesses) leak."""
|
||||
timeout = self._adapter_disconnect_timeout_secs()
|
||||
label = platform.value if platform is not None else "adapter"
|
||||
with _log_suppressed(logging.DEBUG, "Defensive %s disconnect after failed connect raised: %s", label):
|
||||
@@ -82,18 +77,14 @@ class GatewayAdapterLifecycleMixin:
|
||||
)
|
||||
|
||||
async def _bounded_adapter_teardown(self, adapter, platform, *, profile: Optional[str] = None) -> None:
|
||||
"""Tear down one adapter on the shutdown path with bounded awaits (never raises).
|
||||
|
||||
Unbounded, a half-dead transport stalls shutdown past systemd's ``TimeoutStopSec``; the
|
||||
SIGKILL skips ``atexit`` PID-file cleanup and the next start dies with "PID file race lost".
|
||||
"""
|
||||
"""Tear down one adapter on the shutdown path with bounded awaits (never raises). Unbounded,
|
||||
a half-dead transport stalls past systemd's ``TimeoutStopSec``; the SIGKILL skips ``atexit``
|
||||
PID-file cleanup and the next start dies with "PID file race lost"."""
|
||||
timeout = self._adapter_disconnect_timeout_secs()
|
||||
suffix = f" (profile: {profile})" if profile else ""
|
||||
started_at = time.monotonic()
|
||||
try:
|
||||
if not await self._await_adapter_cleanup_with_timeout(
|
||||
adapter.cancel_background_tasks(), timeout
|
||||
):
|
||||
if not await self._await_adapter_cleanup_with_timeout(adapter.cancel_background_tasks(), timeout):
|
||||
logger.warning(
|
||||
"✗ %s background-task cancel timed out after %.1fs - forcing continue%s",
|
||||
platform.value, timeout, suffix,
|
||||
@@ -177,11 +168,9 @@ class GatewayAdapterLifecycleMixin:
|
||||
await self.hooks.emit(event_name, ctx)
|
||||
|
||||
async def _handle_adapter_fatal_error(self, adapter: BasePlatformAdapter) -> None:
|
||||
"""React to an adapter failure after startup (retryable → background reconnect queue).
|
||||
|
||||
Detached task: the notification arrives on the failing adapter's own polling task, which
|
||||
the handler's disconnect can cancel mid-flight, stranding the platform half-handled.
|
||||
"""
|
||||
"""React to an adapter failure after startup (retryable → background reconnect queue). Runs
|
||||
detached: the notification arrives on the failing adapter's own polling task, which the
|
||||
handler's disconnect can cancel mid-flight, stranding the platform half-handled."""
|
||||
tasks = getattr(self, "_fatal_handler_tasks", None)
|
||||
if tasks is None:
|
||||
tasks = self._fatal_handler_tasks = set()
|
||||
@@ -363,16 +352,12 @@ class GatewayAdapterLifecycleMixin:
|
||||
def _spawn_supervised(
|
||||
self, coro_factory, name, *, restart=True, _attempt=0, on_spawn=None, on_give_up=None
|
||||
):
|
||||
"""Launch a long-lived background task with task-level supervision.
|
||||
|
||||
Catches outer-loop/pre-try exceptions a bare ``create_task`` drops silently. Restarts with
|
||||
capped backoff up to ``_MAX_SUPERVISED_RESTARTS`` rapid failures; the counter resets after
|
||||
``_SUPERVISED_HEALTHY_SECS`` of healthy running. Each spawn uses a fresh ``Context`` (an
|
||||
inherited delegated-child marker would make the Kanban dispatcher reject its own writes).
|
||||
``on_spawn`` fires on EVERY spawn incl. respawns — callers tracking the handle elsewhere
|
||||
MUST pass it or a respawn leaves a stale handle and a SECOND watcher. ``on_give_up(name)``
|
||||
fires when the restart budget is spent.
|
||||
"""
|
||||
"""Launch a long-lived supervised background task: exceptions a bare ``create_task`` drops
|
||||
are logged, and it respawns with capped backoff up to ``_MAX_SUPERVISED_RESTARTS`` rapid
|
||||
failures (counter resets after ``_SUPERVISED_HEALTHY_SECS`` healthy). Fresh ``Context`` per
|
||||
spawn (an inherited delegated-child marker would make the Kanban dispatcher reject its own
|
||||
writes). ``on_spawn`` fires on EVERY spawn incl. respawns — handle trackers MUST pass it or a
|
||||
respawn leaves a stale handle and a SECOND watcher; ``on_give_up(name)`` fires at budget end."""
|
||||
# Spawn timestamp lets ``_done`` tell a rapid crash-loop from a healthy-run-then-crash.
|
||||
_started = time.monotonic()
|
||||
|
||||
@@ -434,11 +419,9 @@ class GatewayAdapterLifecycleMixin:
|
||||
return task
|
||||
|
||||
async def _handoff_watcher(self, interval: float = 2.0, drain_timeout: float = 30.0) -> None:
|
||||
"""Process pending CLI→gateway session handoffs from ``state.db``.
|
||||
|
||||
Claim atomically (pending → running), re-bind the home channel to the CLI session_id,
|
||||
dispatch a synthetic ``MessageEvent``, mark ``completed``/``failed``.
|
||||
"""
|
||||
"""Process pending CLI→gateway session handoffs from ``state.db``: claim atomically (pending
|
||||
→ running), re-bind the home channel to the CLI session_id, dispatch a synthetic event, mark
|
||||
``completed``/``failed``."""
|
||||
from gateway.run import _async_profile_runtime_scope, _handoff_watch_scopes, _reclaim_stale
|
||||
await asyncio.sleep(5) # let platforms connect before dispatching through them
|
||||
|
||||
@@ -470,12 +453,9 @@ class GatewayAdapterLifecycleMixin:
|
||||
inflight.pop(session_id, None)
|
||||
|
||||
async def _tick(profile_name: Optional[str] = None) -> None:
|
||||
"""One poll of the CURRENTLY-SCOPED store; ``profile_name`` (None = root) routes
|
||||
delivery to that profile's OWN adapter.
|
||||
|
||||
A closure, not a method: tests bind ``_handoff_watcher`` onto a ``SimpleNamespace``
|
||||
with only ``_session_db``/``_running``/``_process_handoff``.
|
||||
"""
|
||||
"""One poll of the CURRENTLY-SCOPED store; ``profile_name`` (None = root) routes delivery
|
||||
to that profile's OWN adapter. A closure, not a method: tests bind ``_handoff_watcher`` onto
|
||||
a ``SimpleNamespace`` with only ``_session_db``/``_running``/``_process_handoff``."""
|
||||
session_db = getattr(self, "_session_db", None)
|
||||
if session_db is None:
|
||||
return
|
||||
@@ -637,8 +617,7 @@ class GatewayAdapterLifecycleMixin:
|
||||
queued_for / 3600.0, info.get("attempts", 0),
|
||||
)
|
||||
self._update_platform_runtime_status(
|
||||
platform.value, platform_state="retrying", needs_attention=True,
|
||||
retrying_since=retrying_since_iso,
|
||||
platform.value, platform_state="retrying", needs_attention=True, retrying_since=retrying_since_iso
|
||||
)
|
||||
|
||||
def _mark_platform_fatal(self, status_key: str, adapter) -> None:
|
||||
@@ -763,9 +742,7 @@ class GatewayAdapterLifecycleMixin:
|
||||
try:
|
||||
self._schedule_resume_pending_sessions(platform=platform)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"resume-pending reschedule after %s reconnect failed", platform.value, exc_info=True
|
||||
)
|
||||
logger.debug("resume-pending reschedule after %s reconnect failed", platform.value, exc_info=True)
|
||||
|
||||
async def _cancel_secondary_profile_reconnect_tasks(self) -> None:
|
||||
"""Cancel profile-scoped reconnects before tearing down their registry, so a reconnect
|
||||
@@ -794,10 +771,8 @@ class GatewayAdapterLifecycleMixin:
|
||||
|
||||
async def _start_secondary_profile_adapters(self) -> int:
|
||||
"""Bring up adapters for every non-active profile (multiplex only); returns connected count.
|
||||
|
||||
Each profile connects under its own HERMES_HOME + secret scope; credential/listener
|
||||
collisions are refused here — the only point seeing every profile's credentials together.
|
||||
"""
|
||||
Each profile connects under its own HERMES_HOME + secret scope; credential/listener collisions
|
||||
are refused here — the only point seeing every profile's credentials together."""
|
||||
from gateway.run import (
|
||||
MultiplexConfigError, SecondaryPortBindingConfigError, _multiplex_profile_homes
|
||||
)
|
||||
@@ -825,9 +800,7 @@ class GatewayAdapterLifecycleMixin:
|
||||
except MultiplexConfigError:
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
"Failed to start adapters for profile '%s': %s", profile_name, e, exc_info=True
|
||||
)
|
||||
logger.error("Failed to start adapters for profile '%s': %s", profile_name, e, exc_info=True)
|
||||
|
||||
self._record_served_profiles(active, profile_homes)
|
||||
return connected
|
||||
@@ -863,11 +836,9 @@ class GatewayAdapterLifecycleMixin:
|
||||
write_runtime_status(served_profiles=served)
|
||||
|
||||
async def _load_secondary_profile_config(self, profile_name: str, profile_home: "Path"):
|
||||
"""Hydrate + enter ``profile_home``'s scope once; return its gateway config.
|
||||
|
||||
Raises ``MultiplexConfigError`` (open dm/group policy) or ``SecondaryPortBindingConfigError``
|
||||
(the default profile owns the single shared HTTP listener).
|
||||
"""
|
||||
"""Hydrate + enter ``profile_home``'s scope once; return its gateway config. Raises
|
||||
``MultiplexConfigError`` (open dm/group policy) or ``SecondaryPortBindingConfigError`` (the
|
||||
default profile owns the single shared HTTP listener)."""
|
||||
from gateway.run import (
|
||||
MultiplexConfigError, SecondaryPortBindingConfigError, _load_gateway_runtime_config,
|
||||
_own_policy_open_startup_violation, _profile_runtime_scope,
|
||||
@@ -897,8 +868,7 @@ class GatewayAdapterLifecycleMixin:
|
||||
raise MultiplexConfigError(
|
||||
f"Profile '{profile_name}' enables {violation}. "
|
||||
"Enable GATEWAY_ALLOW_ALL_USERS or the platform allow-all flag "
|
||||
"for that profile, or change dm_policy/group_policy away from "
|
||||
"'open'."
|
||||
"for that profile, or change dm_policy/group_policy away from 'open'."
|
||||
)
|
||||
|
||||
port_binding_platforms = sorted(
|
||||
@@ -1074,11 +1044,9 @@ class GatewayAdapterLifecycleMixin:
|
||||
adapter._hermes_profile_name = profile_name
|
||||
|
||||
async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platform):
|
||||
"""One scoped attempt to rebuild+connect a secondary adapter → ``(adapter, success)``.
|
||||
|
||||
``(None, None)`` = give up for good (disabled, credential removed, adapter unavailable).
|
||||
Caller tears down a RETURNED adapter; one whose configure/connect raised is torn down here.
|
||||
"""
|
||||
"""One scoped attempt to rebuild+connect a secondary adapter → ``(adapter, success)``;
|
||||
``(None, None)`` = give up for good (disabled, credential removed, adapter unavailable). Caller
|
||||
tears down a RETURNED adapter; one whose configure/connect raised is torn down here."""
|
||||
from gateway.run import _platform_has_bot_credential, _profile_runtime_scope
|
||||
# Lazy + per-attempt: keeps test monkeypatches on these modules live.
|
||||
from hermes_cli.profiles import get_profile_dir
|
||||
@@ -1162,8 +1130,7 @@ class GatewayAdapterLifecycleMixin:
|
||||
attempts += 1
|
||||
backoff = _reconnect_backoff(attempts)
|
||||
logger.info(
|
||||
"Secondary %s reconnect retry in %ds (profile: %s)", platform.value, backoff,
|
||||
profile_name,
|
||||
"Secondary %s reconnect retry in %ds (profile: %s)", platform.value, backoff, profile_name
|
||||
)
|
||||
await asyncio.sleep(backoff)
|
||||
finally:
|
||||
|
||||
@@ -245,31 +245,24 @@ class GatewayNotificationsMixin:
|
||||
if switched is None:
|
||||
logger.warning(
|
||||
"Async-delegation completion could not bind routing key %s to "
|
||||
"owning session %s; dropping injection.",
|
||||
session_entry.session_key, target_session_id,
|
||||
"owning session %s; dropping injection.", session_entry.session_key, target_session_id,
|
||||
)
|
||||
return None
|
||||
|
||||
logger.info(
|
||||
"Pinned async-delegation completion to owning session %s "
|
||||
"(was %s) for routing key %s (#57498)",
|
||||
"Pinned async-delegation completion to owning session %s (was %s) for routing key %s (#57498)",
|
||||
target_session_id, prior_session_id, session_entry.session_key,
|
||||
)
|
||||
return switched
|
||||
|
||||
async def _deliver_media_from_response(
|
||||
self, response: str, event: MessageEvent, adapter,
|
||||
thread_metadata: Optional[Dict[str, Any]] = None,
|
||||
self, response: str, event: MessageEvent, adapter, thread_metadata: Optional[Dict[str, Any]] = None
|
||||
) -> None:
|
||||
"""Extract explicit MEDIA: tags from an already-streamed response and deliver them.
|
||||
|
||||
The text is already delivered; this only handles file attachments. Unlike the non-streaming
|
||||
path in ``gateway/platforms/base.py`` this rescan is EXPLICIT-ONLY: a bare local path in a
|
||||
streamed reply was either shown as text or is stale inspected content, and promoting it sent
|
||||
files the model never asked to deliver. MEDIA tags are NOT deduped against prior turns: a
|
||||
MEDIA: directive in the final reply is a deliberate attach (incl. user-requested resends);
|
||||
stale auto-appended tags are deduped upstream (_collect_auto_append_media_tags).
|
||||
"""
|
||||
"""Deliver explicit MEDIA: tags from an already-streamed response (text already delivered).
|
||||
EXPLICIT-ONLY, unlike the non-streaming path in ``gateway/platforms/base.py``: a bare local
|
||||
path in a streamed reply is shown text or stale inspected content, and promoting it sent
|
||||
files the model never asked for. MEDIA tags are NOT deduped against prior turns (a final-reply
|
||||
directive is a deliberate attach); stale auto-appended tags are deduped upstream."""
|
||||
from urllib.parse import quote as _quote
|
||||
with _log_suppressed(logging.WARNING, "Post-stream media extraction failed: %s"):
|
||||
# Capture [[as_document]] before extract_media strips it: image files then go through
|
||||
@@ -378,8 +371,7 @@ class GatewayNotificationsMixin:
|
||||
from gateway.run import _hermes_home
|
||||
return cls._UpdatePaths(
|
||||
pending=_hermes_home / ".update_pending.json",
|
||||
claimed=_hermes_home / ".update_pending.claimed.json",
|
||||
output=_hermes_home / ".update_output.txt",
|
||||
claimed=_hermes_home / ".update_pending.claimed.json", output=_hermes_home / ".update_output.txt",
|
||||
exit_code=_hermes_home / ".update_exit_code",
|
||||
prompt=_hermes_home / ".update_prompt.json", response=_hermes_home / ".update_response",
|
||||
)
|
||||
@@ -664,8 +656,7 @@ class GatewayNotificationsMixin:
|
||||
platform_cfg = self.config.platforms.get(platform)
|
||||
if platform_cfg is not None and not platform_cfg.gateway_restart_notification:
|
||||
logger.info(
|
||||
"Restart notification suppressed: %s has gateway_restart_notification=false",
|
||||
platform_str,
|
||||
"Restart notification suppressed: %s has gateway_restart_notification=false", platform_str
|
||||
)
|
||||
return None
|
||||
|
||||
@@ -833,8 +824,7 @@ class GatewayNotificationsMixin:
|
||||
|
||||
try:
|
||||
platform = Platform(platform_name)
|
||||
# Reject arbitrary strings (dynamic pseudo-members): built-ins are always valid, plugin
|
||||
# platforms must be registered in the platform registry.
|
||||
# Reject dynamic pseudo-members: plugin platforms must be registered.
|
||||
if platform.value not in _BUILTIN_PLATFORM_VALUES:
|
||||
try:
|
||||
from gateway.platform_registry import platform_registry
|
||||
@@ -851,15 +841,12 @@ class GatewayNotificationsMixin:
|
||||
|
||||
scope_id = _opt("scope_id")
|
||||
if scope_id is None and chat_type not in ("dm", "thread"):
|
||||
# Reconstructed (non-persisted) source for a scoped chat with no scope discriminator: a
|
||||
# relay connector's fail-closed tenant guard may decline the reply unless user_id resolves
|
||||
# it. Don't fail — DMs/author-bound chats still route and native adapters need no
|
||||
# scope_id — but warn so a post-restart egress decline isn't silent.
|
||||
# Reconstructed scoped-chat source without scope_id: a relay connector's tenant guard may
|
||||
# decline the reply. Warn, don't fail (native adapters need no scope_id).
|
||||
logger.warning(
|
||||
"Synthetic event source for %s chat=%s (%s) reconstructed "
|
||||
"without scope_id; scoped relay egress may be declined by "
|
||||
"the connector's tenant guard (user_id fallback only).",
|
||||
platform_name, chat_id, chat_type,
|
||||
"the connector's tenant guard (user_id fallback only).", platform_name, chat_id, chat_type,
|
||||
)
|
||||
return SessionSource(
|
||||
platform=platform, chat_id=chat_id, chat_type=chat_type, thread_id=_opt("thread_id"),
|
||||
@@ -916,15 +903,9 @@ class GatewayNotificationsMixin:
|
||||
return False
|
||||
|
||||
def _resolve_injection_adapter(self, platform_name: str):
|
||||
"""Adapter for a synthetic-event platform: transport resolver first, literal scan as fallback.
|
||||
|
||||
Alias-aware resolution (relay-plane): one adapter under Platform.RELAY fronts N logical
|
||||
platforms, so a literal ``p.value == platform_name`` scan misses "slack" and drops the
|
||||
completion as "no gateway route". Native adapter wins; relay is eligible only when it
|
||||
advertises fronting the logical platform. The legacy literal scan is still correct for
|
||||
native adapters and keeps minimal runner stubs (tests) / exotic platform strings working
|
||||
when the resolver can't run.
|
||||
"""
|
||||
"""Adapter for a synthetic-event platform: alias-aware transport resolver first (one
|
||||
Platform.RELAY adapter fronts N logical platforms; native wins), literal ``p.value`` scan as
|
||||
fallback for minimal runner stubs / exotic platform strings when the resolver can't run."""
|
||||
from gateway.run import resolve_delivery_transport
|
||||
try:
|
||||
_transport = resolve_delivery_transport(Platform(platform_name), self.config, self.adapters)
|
||||
@@ -1042,15 +1023,10 @@ class GatewayNotificationsMixin:
|
||||
return seen
|
||||
|
||||
async def _classify_completion_target(self, parent_session_id: str) -> str:
|
||||
"""Classify an async-completion delivery target before adapter acceptance.
|
||||
|
||||
- ``"deliver"``: spawning session live (or compression-rotated with a live continuation);
|
||||
proves deliverability only, the resolver still retargets.
|
||||
- ``"terminal"``: parent gone for good (unknown / explicit user boundary like /new); drop
|
||||
the durable row rather than falsely ack or replay forever.
|
||||
- ``"retry"``: transient uncertainty (DB unavailable, rotation mid-flight); release the
|
||||
claim for a later consumer; the attempt cap bounds churn.
|
||||
"""
|
||||
"""Classify an async-completion target before adapter acceptance: ``"deliver"`` (spawning
|
||||
session live or compression-rotated with a live continuation; the resolver still retargets),
|
||||
``"terminal"`` (parent gone for good — unknown / user boundary like /new; drop the durable row
|
||||
rather than falsely ack), ``"retry"`` (DB unavailable / rotation mid-flight; release the claim)."""
|
||||
from gateway.run import _USER_BOUNDARY_END_REASONS
|
||||
session_db = getattr(self, "_session_db", None)
|
||||
if session_db is None:
|
||||
@@ -1139,8 +1115,7 @@ class GatewayNotificationsMixin:
|
||||
logger.warning(
|
||||
"Background process %s completion targets "
|
||||
"permanently-gone session %s (user boundary such as "
|
||||
"/new); dropping notification (output remains "
|
||||
"available via process(action='log')).",
|
||||
"/new); dropping notification (output remains available via process(action='log')).",
|
||||
evt.get("session_id") or "<unknown>", parent_session_id,
|
||||
)
|
||||
claim.proceed = False
|
||||
@@ -1309,12 +1284,10 @@ class GatewayNotificationsMixin:
|
||||
# Some unit tests construct GatewayRunner with object.__new__. Keep the
|
||||
# batching seam lazy so those focused lifecycle tests remain valid.
|
||||
for attr, default in (
|
||||
("_completion_notification_batches", dict),
|
||||
("_completion_notification_batch_tasks", dict),
|
||||
("_completion_notification_batches", dict), ("_completion_notification_batch_tasks", dict),
|
||||
("_completion_notification_batch_flush_tasks", set),
|
||||
("_completion_notification_batch_window", lambda: 0.1),
|
||||
("_completion_notification_batches_stopping", lambda: False),
|
||||
("_background_tasks", set),
|
||||
("_completion_notification_batches_stopping", lambda: False), ("_background_tasks", set),
|
||||
):
|
||||
if not hasattr(self, attr):
|
||||
setattr(self, attr, default())
|
||||
@@ -1360,14 +1333,10 @@ class GatewayNotificationsMixin:
|
||||
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.
|
||||
|
||||
Single-event groups ride the per-event path. Multi-event groups deliver the primary via
|
||||
``_deliver_completion_notification`` with consolidated text of every sibling THIS runner
|
||||
claimed; sibling claims are acked only after adapter acceptance, and siblings claimed by
|
||||
another consumer are excluded (no double delivery). Returns True after acceptance, False
|
||||
to requeue the group, None when nothing is deliverable here (retry siblings requeued).
|
||||
"""
|
||||
"""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;
|
||||
their claims are acked only after adapter acceptance). True after acceptance, False to requeue
|
||||
the group, None when nothing is deliverable here (retry siblings requeued)."""
|
||||
from gateway.run import _format_gateway_process_notification
|
||||
from tools.process_registry import process_registry as _pr
|
||||
deliverable: list[tuple[dict, str]] = []
|
||||
@@ -1547,13 +1516,10 @@ class GatewayNotificationsMixin:
|
||||
)
|
||||
|
||||
async def _run_process_watcher(self, watcher: dict) -> None:
|
||||
"""Periodically check a background process and push updates to the user.
|
||||
|
||||
Runs as an asyncio task. Stays silent when nothing changed. Auto-removes when the process
|
||||
exits or is killed. Notification mode (``display.background_process_notifications``):
|
||||
concise (default, one-line; failures append output tail) / all (running updates + final
|
||||
raw output) / result (final raw only) / error (final raw only if exit != 0) / off.
|
||||
"""
|
||||
"""Poll a background process and push updates until it exits. Mode
|
||||
(``display.background_process_notifications``): concise (default one-liner; failures append
|
||||
the output tail) / all (running updates + final raw) / result (final raw) / error (final raw
|
||||
if exit != 0) / off."""
|
||||
from tools.process_registry import format_process_notification, process_registry
|
||||
session_id = watcher["session_id"]
|
||||
interval = watcher["check_interval"]
|
||||
|
||||
+28
-50
@@ -385,13 +385,10 @@ class GatewayShutdownMixin:
|
||||
return self.adapters.get(Platform.RELAY)
|
||||
|
||||
async def _scale_to_zero_watcher(self, interval: float = 30.0) -> None:
|
||||
"""Watch for idle, drive the relay dormant, then self-suspend the machine.
|
||||
|
||||
On sustained idle: status `draining` (NOT _running=False), relay go_dormant() (socket close,
|
||||
NOT disconnect()), no mark_resume_pending (suspend preserves RAM), THEN suspend via the flaps
|
||||
socket. The gateway owns the suspend because Fly autostop sees only INBOUND connections and
|
||||
would freeze mid-job or before the relay flip. Off-Fly the watcher does not quiesce at all.
|
||||
"""
|
||||
"""Watch for idle, drive the relay dormant, then self-suspend. On sustained idle: status
|
||||
`draining` (NOT _running=False), relay go_dormant() (socket close, NOT disconnect()), no
|
||||
mark_resume_pending (suspend preserves RAM), THEN suspend via the flaps socket — Fly autostop
|
||||
sees only INBOUND connections and would freeze mid-job. Off-Fly: no quiesce at all."""
|
||||
await asyncio.sleep(min(interval, 30.0)) # let startup settle
|
||||
while self._running:
|
||||
try:
|
||||
@@ -473,8 +470,7 @@ class GatewayShutdownMixin:
|
||||
self._external_drain_active = True
|
||||
logger.info(
|
||||
"External drain ENGAGED (.drain_request.json present) — refusing "
|
||||
"new turns; %d in-flight turn(s) will finish. Process stays up.",
|
||||
self._active_work_count(),
|
||||
"new turns; %d in-flight turn(s) will finish. Process stays up.", self._active_work_count(),
|
||||
)
|
||||
# Flip persisted lifecycle state so /api/status.gateway_busy / gateway_drainable track the
|
||||
# drain; active_agents is preserved (read-merge keeps the live count), only state changes.
|
||||
@@ -502,9 +498,8 @@ class GatewayShutdownMixin:
|
||||
from gateway.drain_control import drain_requested
|
||||
while self._running:
|
||||
try:
|
||||
# drain_requested() does a synchronous read_text() on the marker file: at 1s cadence
|
||||
# that is a blocking disk read on the event loop ~86k times/day, and under host I/O
|
||||
# pressure one read can stall 30s+ and take every platform heartbeat down. Off-thread it.
|
||||
# Off-thread: a synchronous marker read at 1s cadence can stall 30s+ under host I/O
|
||||
# pressure and take every platform heartbeat down.
|
||||
if await asyncio.to_thread(drain_requested):
|
||||
self._enter_external_drain()
|
||||
# API and cron work live outside messaging's _running_agents map; refresh the
|
||||
@@ -1085,9 +1080,8 @@ class GatewayShutdownMixin:
|
||||
)
|
||||
watcher_env = GatewayShutdownMixin._restart_watcher_env()
|
||||
project_root = Path(__file__).resolve().parent.parent
|
||||
# Console python under CREATE_NO_WINDOW owns one hidden console inherited by the restart
|
||||
# child, so nothing flashes. Do NOT swap in pythonw.exe — a console-less watcher forces
|
||||
# every console-subsystem descendant to allocate a visible conhost.
|
||||
# Console python under CREATE_NO_WINDOW: nothing flashes. NOT pythonw.exe — a console-less
|
||||
# watcher makes every console-subsystem descendant allocate a visible conhost.
|
||||
watcher_python = sys.executable
|
||||
venv_dir = Path(watcher_env.get("VIRTUAL_ENV") or project_root / "venv")
|
||||
site_packages = venv_dir / "Lib" / "site-packages"
|
||||
@@ -1102,10 +1096,8 @@ class GatewayShutdownMixin:
|
||||
str(current_pid), str(restart_after_s), *hermes_cmd, "gateway", "restart",
|
||||
]
|
||||
popen_kwargs = dict(stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, env=watcher_env)
|
||||
# The watcher must break away from the parent CLI's job object (Desktop wrappers, Windows
|
||||
# Terminal, schtasks) or it is reaped when the CLI exits. A job without
|
||||
# JOB_OBJECT_LIMIT_BREAKAWAY_OK rejects CREATE_BREAKAWAY_FROM_JOB (OSError); retry once
|
||||
# without the bit, preserving argv and the scrubbed env.
|
||||
# Break away from the parent CLI's job object or be reaped when the CLI exits; a job without
|
||||
# BREAKAWAY_OK rejects CREATE_BREAKAWAY_FROM_JOB (OSError) — retry once without the bit.
|
||||
try:
|
||||
subprocess.Popen(watcher_argv, **popen_kwargs, **windows_detach_popen_kwargs())
|
||||
except OSError:
|
||||
@@ -1406,8 +1398,7 @@ class GatewayShutdownMixin:
|
||||
if _cron_at_start and _cron_timeout > timeout:
|
||||
logger.info(
|
||||
"Shutdown drain: %d in-flight cron job(s) — waiting up to "
|
||||
"%.0fs for them (cron_drain_timeout=%.0fs, "
|
||||
"restart_drain_timeout=%.0fs)",
|
||||
"%.0fs for them (cron_drain_timeout=%.0fs, restart_drain_timeout=%.0fs)",
|
||||
_cron_at_start, _cron_timeout, _cron_drain_cfg, timeout,
|
||||
)
|
||||
_drain_started_at = time.monotonic()
|
||||
@@ -1415,8 +1406,7 @@ class GatewayShutdownMixin:
|
||||
ctx.drain_elapsed = time.monotonic() - _drain_started_at
|
||||
logger.info(
|
||||
"Shutdown phase: drain done at +%.2fs (drain took %.2fs, "
|
||||
"timed_out=%s, active_at_start=%d, active_now=%d, "
|
||||
"cron_at_start=%d, cron_now=%d, "
|
||||
"timed_out=%s, active_at_start=%d, active_now=%d, cron_at_start=%d, cron_now=%d, "
|
||||
"api_at_start=%d, api_now=%d, "
|
||||
"deferred_at_start=%d, deferred_now=%d)", ctx.elapsed(), ctx.drain_elapsed,
|
||||
ctx.timed_out, len(ctx.active_agents), self._running_agent_count(), _cron_at_start,
|
||||
@@ -1439,15 +1429,12 @@ class GatewayShutdownMixin:
|
||||
from gateway.run import GatewayRunner
|
||||
logger.warning(
|
||||
"Gateway drain timed out after %.1fs with %d active agent(s), "
|
||||
"%d in-flight cron job(s), %d api_server run(s), and "
|
||||
"%d deferred agent worker(s); "
|
||||
"%d in-flight cron job(s), %d api_server run(s), and %d deferred agent worker(s); "
|
||||
"interrupting remaining work.", ctx.drain_elapsed, self._running_agent_count(),
|
||||
self._active_cron_job_count(), self._active_api_run_count(), ctx.deferred_count(),
|
||||
)
|
||||
# Mark forcibly-interrupted sessions resume_pending BEFORE interrupting so the next message
|
||||
# auto-resumes instead of becoming a fresh session via suspend_recently_active(); stuck
|
||||
# sessions still escalate via ``.restart_failure_counts`` (suspended=True wins). Uses the
|
||||
# CURRENT _running_agents, not the drain-start snapshot, so cleanly finished sessions get no note.
|
||||
# Mark resume_pending BEFORE interrupting so the next message auto-resumes (stuck sessions
|
||||
# still escalate via .restart_failure_counts). CURRENT _running_agents, not the drain snapshot.
|
||||
await GatewayRunner._mark_running_sessions_resume_pending(self, "mark_resume_pending")
|
||||
reason = GatewayRunner._shutdown_interrupt_reason(self)
|
||||
self._interrupt_running_agents(reason)
|
||||
@@ -1568,11 +1555,9 @@ class GatewayShutdownMixin:
|
||||
def _stop_quiesce_and_close_session_dbs(self, timeout: float, ctx: "GatewayShutdownMixin._StopContext") -> None:
|
||||
"""Quiesce the executor, then close SessionDB handles only if no worker is still live."""
|
||||
from gateway.run import GatewayRunner, _EXECUTOR_QUIESCE_TIMEOUT
|
||||
# Quiesce the thread pool BEFORE closing session DBs: otherwise a coroutine can mint a fresh
|
||||
# pool (``_executor_closing`` still False) or an already-started run_in_executor future keeps
|
||||
# writing after ``SessionDB.close()`` checkpointed the WAL — the late write reopens the handle
|
||||
# and mints a split WAL generation (close-time corruption). Bounded and clamped to what is
|
||||
# left of the watchdog leash minus a second for the close itself.
|
||||
# Quiesce the thread pool BEFORE closing session DBs: a late executor write after
|
||||
# SessionDB.close() checkpointed the WAL reopens the handle and splits the WAL generation
|
||||
# (close-time corruption). Clamped to the remaining watchdog leash minus 1s for the close.
|
||||
_exec_quiesce_budget = max(
|
||||
0.0, min(_EXECUTOR_QUIESCE_TIMEOUT, resolve_shutdown_watchdog_delay(timeout) - ctx.elapsed() - 1.0),
|
||||
)
|
||||
@@ -1637,8 +1622,7 @@ class GatewayShutdownMixin:
|
||||
else:
|
||||
logger.info(
|
||||
"Skipping .clean_shutdown marker — drain timed out with "
|
||||
"interrupted agents; next startup will suspend recently "
|
||||
"active sessions."
|
||||
"interrupted agents; next startup will suspend recently active sessions."
|
||||
)
|
||||
# Stuck-loop detection: the counter increments for sessions active at each restart; at
|
||||
# the threshold (3 consecutive) the next startup auto-suspends the session.
|
||||
@@ -1661,14 +1645,12 @@ class GatewayShutdownMixin:
|
||||
self._exit_code = GATEWAY_SERVICE_RESTART_EXIT_CODE
|
||||
self._exit_reason = self._exit_reason or "Gateway restart requested"
|
||||
self._draining = False
|
||||
# Terminal gateway_state: "stopped" by default, "running" on an UNEXPECTED signal (s6 SIGTERM
|
||||
# on docker restart, OOM-kill) — container_boot.py only auto-starts gateways last seen
|
||||
# "running". Operator stops write a planned-stop marker BEFORE signalling; restarts persist "stopped".
|
||||
# Terminal gateway_state: "stopped", or "running" on an UNEXPECTED signal (docker restart,
|
||||
# OOM) — container_boot.py only auto-starts gateways last seen "running".
|
||||
if getattr(self, "_signal_initiated_shutdown", False) and not self._restart_requested:
|
||||
logger.info(
|
||||
"Gateway stopped by an unexpected signal — persisting "
|
||||
"gateway_state=running so container_boot auto-starts on "
|
||||
"the next boot (issue #42675)"
|
||||
"gateway_state=running so container_boot auto-starts on the next boot (issue #42675)"
|
||||
)
|
||||
self._update_runtime_status("running", self._exit_reason)
|
||||
else:
|
||||
@@ -1694,9 +1676,8 @@ class GatewayShutdownMixin:
|
||||
async def _stop_impl(self) -> None:
|
||||
"""Run every ``_stop_*`` phase under the thread-based shutdown watchdog."""
|
||||
from gateway.run import GatewayRunner
|
||||
# Thread-based watchdog (asyncio timeouts cannot recover a frozen loop): if teardown never
|
||||
# finishes within drain+grace it dumps faulthandler stacks and os._exit so KeepAlive/systemd
|
||||
# can revive. Skipped under pytest (no delayed hard-exit in the worker).
|
||||
# Thread-based watchdog (asyncio timeouts cannot recover a frozen loop): dumps stacks and
|
||||
# os._exit past drain+grace so the service manager revives us. Skipped under pytest.
|
||||
_watchdog_done = threading.Event()
|
||||
self._shutdown_watchdog_done = _watchdog_done
|
||||
# Shutdown-path tests and third-party runner doubles may only implement the older
|
||||
@@ -1706,10 +1687,8 @@ class GatewayShutdownMixin:
|
||||
)
|
||||
if not os.environ.get("PYTEST_CURRENT_TEST"):
|
||||
arm_shutdown_watchdog(
|
||||
resolve_shutdown_watchdog_delay(self._restart_drain_timeout),
|
||||
done_event=_watchdog_done,
|
||||
snapshot_fn=lambda: GatewayRunner._shutdown_watchdog_snapshot(self, ctx),
|
||||
exit_code=1,
|
||||
resolve_shutdown_watchdog_delay(self._restart_drain_timeout), done_event=_watchdog_done,
|
||||
snapshot_fn=lambda: GatewayRunner._shutdown_watchdog_snapshot(self, ctx), exit_code=1,
|
||||
)
|
||||
try:
|
||||
await GatewayRunner._stop_begin_teardown(self, ctx)
|
||||
@@ -1725,8 +1704,7 @@ class GatewayShutdownMixin:
|
||||
_watchdog_done.set()
|
||||
|
||||
async def stop(
|
||||
self, *, restart: bool = False, detached_restart: bool = False,
|
||||
service_restart: bool = False,
|
||||
self, *, restart: bool = False, detached_restart: bool = False, service_restart: bool = False
|
||||
) -> None:
|
||||
"""Stop the gateway and disconnect all adapters."""
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
+73
-146
@@ -130,8 +130,7 @@ class GatewayStartupMixin:
|
||||
await self._wait_bounded_or_release(
|
||||
{task}, timeout,
|
||||
"Turn-machinery warm-up still running after %.0fs; opening inbound gate anyway — the "
|
||||
"first turn may see lazily initialized machinery (#99373). Warm-up continues in the "
|
||||
"background.",
|
||||
"first turn may see lazily initialized machinery (#99373). Warm-up continues in the background.",
|
||||
"boot turn-machinery warm-up failed after gate release", level=logging.DEBUG,
|
||||
)
|
||||
|
||||
@@ -158,13 +157,10 @@ class GatewayStartupMixin:
|
||||
return done
|
||||
|
||||
async def _finish_startup_restore(self) -> None:
|
||||
"""Wait (BOUNDED) for startup auto-resume, then release + drain inbound.
|
||||
|
||||
Bounded by ``_startup_restore_drain_timeout_secs`` so one pathological boot-resume turn
|
||||
cannot hold the gate shut for every channel; on timeout the gate opens and resume turns
|
||||
finish in the background (NOT cancelled). Safe because ``_schedule_resume_pending_sessions``
|
||||
claims each ``_running_agents`` slot SYNCHRONOUSLY first, so drained inbound queues behind.
|
||||
"""
|
||||
"""Wait (BOUNDED by ``_startup_restore_drain_timeout_secs``) for startup auto-resume, then
|
||||
release + drain inbound. On timeout the gate opens and resume turns finish in the background
|
||||
(NOT cancelled) — safe because ``_schedule_resume_pending_sessions`` claims each
|
||||
``_running_agents`` slot SYNCHRONOUSLY first, so drained inbound queues behind."""
|
||||
from gateway.run import _startup_restore_drain_timeout_secs
|
||||
tasks = list(getattr(self, "_startup_restore_tasks", []) or [])
|
||||
if tasks:
|
||||
@@ -203,15 +199,11 @@ class GatewayStartupMixin:
|
||||
return _report
|
||||
|
||||
async def _await_startup_boot_sends(self, *, planned_restart_notification_pending: bool) -> None:
|
||||
"""Run boot-path sends without letting them pin the inbound restore gate.
|
||||
|
||||
Awaiting the sends inline before the gate releases lets one Telegram flood-control sleep
|
||||
freeze inbound on every platform. Same bounded ``asyncio.wait`` as the resume gate: on
|
||||
timeout return and let the sends finish in the background (not cancelled). The ledger claim
|
||||
+ ``resume_pending`` clear run INLINE before the send task exists (bounded DB work): deferring
|
||||
it let a hung notification expire the gate with zero rows claimed, so answered turns were
|
||||
replayed AND redelivered.
|
||||
"""
|
||||
"""Run boot-path sends without letting them pin the inbound restore gate (one Telegram
|
||||
flood-control sleep must not freeze inbound on every platform): same bounded wait as the
|
||||
resume gate, sends finish in the background on timeout. The ledger claim + ``resume_pending``
|
||||
clear run INLINE before the send task exists — deferring it let a hung notification expire the
|
||||
gate with zero rows claimed, so answered turns were replayed AND redelivered."""
|
||||
from gateway.run import _clear_planned_restart_notification, _startup_restore_drain_timeout_secs
|
||||
claimed = await self._claim_pending_obligations()
|
||||
|
||||
@@ -259,15 +251,11 @@ class GatewayStartupMixin:
|
||||
return sendable
|
||||
|
||||
async def _claim_pending_obligations(self) -> list:
|
||||
"""Claim recoverable delivery-ledger rows and clear their ``resume_pending`` flags.
|
||||
|
||||
Pure DB work, no sends. Must run INLINE at startup BEFORE ``_schedule_resume_pending_sessions``
|
||||
and before the abandonable boot-send task exists: these sessions already produced their
|
||||
answer, so clearing ``resume_pending`` here stops the resume path from re-running (and
|
||||
re-paying for) the turn however long the sends take. Rows that were mid-send or previously
|
||||
rejected carry a visible recovered-reply marker so a possible duplicate is labeled, never
|
||||
silent (gateway/delivery_ledger.py). Returns the claimed rows for redelivery.
|
||||
"""
|
||||
"""Claim recoverable delivery-ledger rows and clear their ``resume_pending`` flags (pure DB
|
||||
work, no sends). Must run INLINE BEFORE ``_schedule_resume_pending_sessions`` and the
|
||||
abandonable boot-send task: these sessions already produced their answer, so the resume path
|
||||
must not re-run (re-pay for) the turn however long the sends take. Mid-send / rejected rows
|
||||
carry a visible recovered-reply marker (gateway/delivery_ledger.py). Returns the rows."""
|
||||
try:
|
||||
from gateway.delivery_ledger import ledger_enabled, sweep_recoverable
|
||||
if not await asyncio.to_thread(ledger_enabled):
|
||||
@@ -333,9 +321,7 @@ class GatewayStartupMixin:
|
||||
try:
|
||||
result = await adapter.send(chat_id=row["chat_id"], content=content, metadata=metadata)
|
||||
except Exception as send_err:
|
||||
logger.warning(
|
||||
"obligation %s: redelivery send raised: %s", row["obligation_id"], send_err
|
||||
)
|
||||
logger.warning("obligation %s: redelivery send raised: %s", row["obligation_id"], send_err)
|
||||
result = None
|
||||
with _log_suppressed(logging.DEBUG, "delivery ledger update failed", exc_info=True):
|
||||
if result is not None and getattr(result, "success", False):
|
||||
@@ -348,8 +334,7 @@ class GatewayStartupMixin:
|
||||
)
|
||||
else:
|
||||
await asyncio.to_thread(
|
||||
mark_failed, row["obligation_id"],
|
||||
str(getattr(result, "error", "") or "send failed"),
|
||||
mark_failed, row["obligation_id"], str(getattr(result, "error", "") or "send failed")
|
||||
)
|
||||
return redelivered
|
||||
|
||||
@@ -358,9 +343,7 @@ class GatewayStartupMixin:
|
||||
try:
|
||||
platform = Platform(row["platform"])
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"obligation %s: unknown platform %r", row["obligation_id"], row.get("platform"),
|
||||
)
|
||||
logger.debug("obligation %s: unknown platform %r", row["obligation_id"], row.get("platform"))
|
||||
return None
|
||||
if "profile" in row:
|
||||
adapter = self._authorization_adapter(platform, row.get("profile"))
|
||||
@@ -384,12 +367,9 @@ class GatewayStartupMixin:
|
||||
async def _redeliver_failed_obligations_for_platform(
|
||||
self, platform: Platform, *, profile: Optional[str] = None,
|
||||
) -> int:
|
||||
"""Replay one adapter identity's transient failures after reconnect.
|
||||
|
||||
The startup sweep cannot claim live-owner rows, and an adapter can reconnect without the
|
||||
process exiting, so ``send_path_degraded`` responses would otherwise stay failed until the
|
||||
next restart. Claim/clear/send are best-effort and reuse the startup redelivery contract.
|
||||
"""
|
||||
"""Replay one adapter identity's transient failures after reconnect: the startup sweep cannot
|
||||
claim live-owner rows, so ``send_path_degraded`` responses would otherwise stay failed until
|
||||
the next restart. Best-effort; reuses the startup redelivery contract."""
|
||||
try:
|
||||
from gateway.delivery_ledger import ledger_enabled, sweep_failed_for_runtime
|
||||
if not await asyncio.to_thread(ledger_enabled):
|
||||
@@ -405,9 +385,7 @@ class GatewayStartupMixin:
|
||||
|
||||
# Clear before any send so the reconnect path cannot both redeliver an
|
||||
# already-produced answer and schedule the same agent turn for resume.
|
||||
sendable = await self._clear_resume_pending_for_claimed_obligations(
|
||||
claimed, require_success=True
|
||||
)
|
||||
sendable = await self._clear_resume_pending_for_claimed_obligations(claimed, require_success=True)
|
||||
sendable_ids = {row["obligation_id"] for row in sendable}
|
||||
for row in claimed:
|
||||
if row["obligation_id"] not in sendable_ids:
|
||||
@@ -434,10 +412,8 @@ class GatewayStartupMixin:
|
||||
logger.warning("Failed to enumerate resume-pending sessions: %s", exc)
|
||||
return None
|
||||
|
||||
# Defense-3: break the SIGTERM-respawn loop. Only count this boot when there are restart-
|
||||
# interrupted sessions to resume — a clean boot must not accrue toward the breaker. If too
|
||||
# many such boots hit the window, skip auto-resume for THIS boot only: the gateway still
|
||||
# serves inbound; the session stays resume_pending so a real user message can continue it.
|
||||
# Restart-loop breaker: only boots WITH restart-interrupted sessions count; when tripped, skip
|
||||
# auto-resume for THIS boot only (inbound still served; sessions stay resume_pending).
|
||||
if candidates:
|
||||
try:
|
||||
from gateway import restart_loop_guard as _rlg
|
||||
@@ -460,20 +436,14 @@ class GatewayStartupMixin:
|
||||
"longer authorized under the current allowlist", session_key,
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.warning(
|
||||
"Skipping auto-resume for %s: authorization check failed: %s", session_key, exc,
|
||||
)
|
||||
logger.warning("Skipping auto-resume for %s: authorization check failed: %s", session_key, exc)
|
||||
return False
|
||||
|
||||
def _schedule_resume_pending_sessions(self, platform=None) -> int:
|
||||
"""Auto-continue fresh restart-interrupted sessions after startup.
|
||||
|
||||
Synthesizes the next turn once adapters are back online; the event text is empty so the
|
||||
existing ``_is_resume_pending`` injection path owns the recovery wording. Sessions whose
|
||||
adapter is not in ``self.adapters`` stay ``resume_pending`` for the reconnect watcher, which
|
||||
re-calls this scoped to that ``platform`` (a reconnecting platform never touches another's
|
||||
recoveries); sessions with a running agent are skipped so none is resumed twice.
|
||||
"""
|
||||
"""Auto-continue fresh restart-interrupted sessions: synthesize an empty-text turn (the
|
||||
``_is_resume_pending`` injection path owns the wording). Sessions whose adapter is offline stay
|
||||
``resume_pending`` for the reconnect watcher, which re-calls this scoped to that ``platform``;
|
||||
sessions with a running agent are skipped so none is resumed twice."""
|
||||
from gateway.run import _AGENT_PENDING_SENTINEL, _auto_continue_freshness_window
|
||||
window = _auto_continue_freshness_window()
|
||||
candidates = self._resume_pending_candidates(platform)
|
||||
@@ -678,24 +648,19 @@ class GatewayStartupMixin:
|
||||
def _open_faulthandler_log(self):
|
||||
"""Open (append) ``<log_dir>/gateway_faulthandler.log``, creating the directory."""
|
||||
from gateway.run import get_hermes_home
|
||||
log_dir = getattr(self.config, "log_dir", None) or os.path.join(
|
||||
str(get_hermes_home()), "logs",
|
||||
)
|
||||
log_dir = getattr(self.config, "log_dir", None) or os.path.join(str(get_hermes_home()), "logs")
|
||||
os.makedirs(log_dir, exist_ok=True)
|
||||
return open(os.path.join(log_dir, "gateway_faulthandler.log"), "a", encoding="utf-8")
|
||||
|
||||
def _start_install_faulthandler(self) -> None:
|
||||
"""Enable faulthandler (stderr or a log file) plus the SIGUSR2 stack-dump hook."""
|
||||
# Falls back to a log file when sys.stderr is None (Windows VBS / pythonw / detached
|
||||
# service) — otherwise the gateway would die here and take every adapter offline.
|
||||
# sys.stderr may be None (Windows VBS / pythonw / detached service): fall back to a log file.
|
||||
try:
|
||||
faulthandler.enable()
|
||||
except (RuntimeError, ValueError, OSError):
|
||||
with _log_suppressed(logging.DEBUG, "faulthandler.enable() unavailable", exc_info=True):
|
||||
faulthandler.enable(file=self._open_faulthandler_log(), all_threads=True)
|
||||
# Also dump stacks to a file on SIGUSR2 for off-line analysis under a service manager that
|
||||
# doesn't capture stderr. faulthandler.register()/SIGUSR2 are POSIX-only: skip on Windows
|
||||
# (faulthandler.enable() above still covers fatal errors).
|
||||
# SIGUSR2 stack dump to file for service managers that drop stderr; POSIX-only.
|
||||
_sigusr2 = getattr(signal, "SIGUSR2", None)
|
||||
if _sigusr2 is not None and hasattr(faulthandler, "register"):
|
||||
with _log_suppressed(logging.DEBUG, "Could not set up faulthandler file logging", exc_info=True):
|
||||
@@ -711,10 +676,8 @@ class GatewayStartupMixin:
|
||||
self._gateway_loop = None
|
||||
if self._gateway_loop is not None:
|
||||
self._start_loop_liveness_guards(self._gateway_loop)
|
||||
# Loop confirmed live: the startup-liveness watchdog is done and the loop-liveness
|
||||
# watchdog (armed above) takes over. Disarm even when loop guards are config-disabled —
|
||||
# the startup watchdog covers only the pre-loop window. Deliberately inside this branch:
|
||||
# if the loop isn't live, startup has NOT reached the milestone and it must stay armed.
|
||||
# Loop live: the loop-liveness watchdog takes over from the startup watchdog. Disarm even
|
||||
# when loop guards are config-disabled; only inside this branch (no live loop = stay armed).
|
||||
with _log_suppressed(logging.DEBUG, "Startup watchdog disarm failed", exc_info=True):
|
||||
from gateway.startup_watchdog import disarm_startup_watchdog
|
||||
disarm_startup_watchdog()
|
||||
@@ -752,9 +715,7 @@ class GatewayStartupMixin:
|
||||
logger.info("Active profile: %s", _profile)
|
||||
with suppress(Exception):
|
||||
from gateway.status import write_runtime_status
|
||||
write_runtime_status(
|
||||
gateway_state="starting", exit_reason=None, clear_profile_platforms=True
|
||||
)
|
||||
write_runtime_status(gateway_state="starting", exit_reason=None, clear_profile_platforms=True)
|
||||
with _log_suppressed(logging.DEBUG, "gateway health OTLP export startup failed", exc_info=True):
|
||||
from hermes_cli.config import load_config
|
||||
from agent.monitoring.gateway_health_export import start_gateway_health_export
|
||||
@@ -910,9 +871,8 @@ class GatewayStartupMixin:
|
||||
if recovered:
|
||||
logger.info("Recovered %s background process(es) from previous run", recovered)
|
||||
|
||||
# Recover sessions active when the gateway last exited. Exact durable turn markers cover
|
||||
# long-running work; the 120s recency heuristic remains as a fallback for turns from older
|
||||
# versions without markers. SKIP after a clean shutdown — the previous process already drained.
|
||||
# Recover sessions active at last exit (exact turn markers + 120s recency fallback for
|
||||
# marker-less older turns). SKIP after a clean exit — the previous process already drained.
|
||||
_clean_marker = _hermes_home / ".clean_shutdown"
|
||||
if _clean_marker.exists():
|
||||
logger.info("Previous gateway exited cleanly — skipping session suspension")
|
||||
@@ -954,10 +914,8 @@ class GatewayStartupMixin:
|
||||
return True, enabled_platform_count, _multiplex_skipped_platforms, _pending_connects
|
||||
if not platform_config.enabled:
|
||||
continue
|
||||
# Under multiplexing, a platform may be enabled on the default profile's config.yaml
|
||||
# while its bot token lives only in a secondary profile's .env. Starting the primary with
|
||||
# an empty token fails at once and queues a reconnect loop that can never heal; the
|
||||
# secondary starts its own adapter with the real token, so skip the empty primary.
|
||||
# Multiplex: a platform enabled in the shared config.yaml may hold its token only in a
|
||||
# secondary profile's .env; an empty primary would queue a reconnect loop that never heals.
|
||||
if _multiplex_on and not _platform_has_bot_credential(platform, platform_config):
|
||||
logger.info(
|
||||
"Skipping %s on default profile: no bot credential in this "
|
||||
@@ -976,8 +934,7 @@ class GatewayStartupMixin:
|
||||
else:
|
||||
logger.warning(
|
||||
"No adapter for '%s' -- is the plugin installed? "
|
||||
"(platform is enabled in config.yaml but no plugin registered it)",
|
||||
platform.value,
|
||||
"(platform is enabled in config.yaml but no plugin registered it)", platform.value,
|
||||
)
|
||||
continue
|
||||
# Under multiplexing the default profile needs the same whole-handler runtime scope as
|
||||
@@ -1010,8 +967,7 @@ class GatewayStartupMixin:
|
||||
# restart/shutdown requested mid-connect must cancel still-pending connects, clean up the
|
||||
# ones already completed, and abort startup.
|
||||
_task_map = {
|
||||
asyncio.ensure_future(_connect_one_startup(p, c, a)): (p, c, a)
|
||||
for (p, c, a) in _pending_connects
|
||||
asyncio.ensure_future(_connect_one_startup(p, c, a)): (p, c, a) for (p, c, a) in _pending_connects
|
||||
}
|
||||
_pending_tasks = set(_task_map)
|
||||
while _pending_tasks:
|
||||
@@ -1064,9 +1020,8 @@ class GatewayStartupMixin:
|
||||
continue
|
||||
if outcome == "exception":
|
||||
logger.error("\u2717 %s error: %s", platform.value, exc)
|
||||
# Same defensive cleanup path for exceptions -- an adapter that raised mid-connect
|
||||
# may still have a live aiohttp.ClientSession or child subprocess. Unexpected
|
||||
# exceptions are typically transient -- queue for retry.
|
||||
# An adapter that raised mid-connect may still hold a ClientSession/subprocess; treat
|
||||
# unexpected exceptions as transient and queue for retry.
|
||||
await self._safe_adapter_disconnect(adapter, platform)
|
||||
self._startup_queue_transient_failure(
|
||||
platform, adapter, platform_config, str(exc), startup_retryable_errors
|
||||
@@ -1091,10 +1046,8 @@ class GatewayStartupMixin:
|
||||
platform, adapter, platform_config, "failed to connect", startup_retryable_errors
|
||||
)
|
||||
continue
|
||||
# A live foreign holder of this bot token is a single-writer ownership conflict, not a
|
||||
# blip — ``_acquire_platform_lock`` emits it retryable only so a MID-RUN reconnect can
|
||||
# recover. At startup route it non-retryable: with nothing connected the gateway exits
|
||||
# 78 instead of sitting alive and deaf in the retry queue.
|
||||
# A live foreign token holder is an ownership conflict, not a blip: retryable only for
|
||||
# MID-RUN reconnects; at startup route it non-retryable so the gateway exits 78, not deaf.
|
||||
_retryable = adapter.fatal_error_retryable and not is_global_startup_conflict(adapter.fatal_error_code)
|
||||
self._update_platform_runtime_status(
|
||||
platform.value, platform_state="retrying" if _retryable else "fatal",
|
||||
@@ -1164,10 +1117,8 @@ class GatewayStartupMixin:
|
||||
self._startup_fail_fatal_config(reason)
|
||||
return True
|
||||
if startup_nonretryable_errors:
|
||||
# Mixed failure mode (some fatal, some transient). Exiting 78 here would let exit-78
|
||||
# supervisors take the gateway PERMANENTLY down over a network blip and deny the
|
||||
# retryable ones their retry. Log the fatal side loudly, then fall through to the
|
||||
# degraded/retry path: the watcher recovers the retryable; the rest stay parked.
|
||||
# Mixed (some fatal, some transient): exiting 78 would take the gateway PERMANENTLY down
|
||||
# over a blip. Log the fatal side loudly and fall through to the degraded/retry path.
|
||||
logger.error(
|
||||
"%d platform(s) fatally misconfigured and parked: %s. "
|
||||
"Staying alive so retryable platforms can recover.",
|
||||
@@ -1238,27 +1189,22 @@ class GatewayStartupMixin:
|
||||
from gateway.run import _planned_restart_notification_pending, _restart_notification_pending
|
||||
await self._start_post_connect_services(connected_count)
|
||||
|
||||
# Give freshly connected adapters a brief moment to settle before sending restart/startup
|
||||
# lifecycle messages; in practice this helps Discord thread deliveries after reconnect.
|
||||
# Let fresh adapters settle before lifecycle sends (helps Discord thread deliveries).
|
||||
if connected_count > 0:
|
||||
await asyncio.sleep(1.0)
|
||||
|
||||
# Capture, before _send_restart_notification() unlinks the marker, whether this process
|
||||
# booted from a chat-originated /restart. One-shot signal for the /restart redelivery
|
||||
# guard (_is_stale_restart_redelivery): a missing dedup marker only suppresses a /restart
|
||||
# when we KNOW we just came out of a restart cycle.
|
||||
# Before _send_restart_notification() unlinks the marker: did we boot from a chat /restart?
|
||||
# One-shot signal for _is_stale_restart_redelivery.
|
||||
if _restart_notification_pending():
|
||||
self._booted_from_restart = True
|
||||
# Restart notification, home-channel startup notice, and obligation redelivery all call
|
||||
# adapter.send(). Those sends must not pin the inbound restore gate — a Telegram flood-
|
||||
# control sleep on this path froze every platform for the full penalty.
|
||||
# Boot-path adapter.send() calls must not pin the inbound restore gate (a Telegram flood-
|
||||
# control sleep here once froze every platform).
|
||||
await self._await_startup_boot_sends(
|
||||
planned_restart_notification_pending=_planned_restart_notification_pending(),
|
||||
)
|
||||
|
||||
# Auto-continue fresh sessions interrupted by the previous restart/shutdown. resume_pending
|
||||
# is cleared by the normal successful-turn path, so a failed auto-resume stays visible on the
|
||||
# next user message. _await_startup_boot_sends already cleared sessions answered in the ledger.
|
||||
# Auto-resume restart-interrupted sessions (ledger-answered ones were cleared above); a failed
|
||||
# auto-resume stays visible on the next user message.
|
||||
self._schedule_resume_pending_sessions()
|
||||
await self._finish_startup_restore()
|
||||
|
||||
@@ -1266,10 +1212,8 @@ class GatewayStartupMixin:
|
||||
# is broken before losing data.
|
||||
await self._send_session_db_warning_notifications()
|
||||
|
||||
# Drain recovered process watchers (crash-recovery checkpoint). Detach the batch atomically:
|
||||
# reassigning a fresh list takes ownership of exactly the watchers present now, so one
|
||||
# appended concurrently during the yield below isn't dropped by a clear() on the shared
|
||||
# list. Yield every 100 to avoid O(n^2) loop blocking when recovering thousands.
|
||||
# Resume recovered process watchers. Detach the batch atomically (fresh list, not clear(): a
|
||||
# concurrent append during the yield must not be lost); yield every 100 to keep the loop live.
|
||||
with _log_suppressed(logging.ERROR, "Recovered watcher setup error: %s"):
|
||||
from tools.process_registry import process_registry
|
||||
watchers = process_registry.pending_watchers
|
||||
@@ -1382,8 +1326,7 @@ class GatewayStartupMixin:
|
||||
if _aborted:
|
||||
return True
|
||||
if self._start_handle_no_connections(
|
||||
connected_count, enabled_platform_count, startup_retryable_errors,
|
||||
startup_nonretryable_errors,
|
||||
connected_count, enabled_platform_count, startup_retryable_errors, startup_nonretryable_errors
|
||||
):
|
||||
return True
|
||||
|
||||
@@ -1447,9 +1390,8 @@ class GatewayStartupMixin:
|
||||
|
||||
handoff_config, handoff_adapters = self._handoff_resolve_scope(profile_name)
|
||||
|
||||
# Adapter must be live. A relay-fronted gateway registers ONE adapter under Platform.RELAY
|
||||
# fronting N logical platforms, so a literal adapters.get(discord) misses a deliverable
|
||||
# platform; resolve_delivery_transport is the alias-aware resolver (native adapter wins).
|
||||
# Alias-aware transport: a relay-fronted gateway registers ONE Platform.RELAY adapter fronting
|
||||
# N logical platforms, so a literal adapters.get() would miss a deliverable one.
|
||||
transport = resolve_delivery_transport(platform, handoff_config, handoff_adapters)
|
||||
if not transport:
|
||||
raise RuntimeError(f"platform '{platform_name}' is not active in this gateway")
|
||||
@@ -1468,16 +1410,12 @@ class GatewayStartupMixin:
|
||||
home_chat_id, f"Hermes — {cli_title}",
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.debug(
|
||||
"Handoff: create_handoff_thread raised on %s: %s", platform_name, exc, exc_info=True
|
||||
)
|
||||
logger.debug("Handoff: create_handoff_thread raised on %s: %s", platform_name, exc, exc_info=True)
|
||||
new_thread_id = None
|
||||
effective_thread_id = new_thread_id or (str(home.thread_id) if home.thread_id else None)
|
||||
|
||||
# Telegram private-chat DM topics are shaped differently from group/forum threads by the
|
||||
# inbound adapter: a handoff-created topic in a positive chat_id must use the DM-topic source
|
||||
# shape (real user id == chat_id, so topic-mode checks and binding persistence match later
|
||||
# inbound turns), or the synthetic turn binds a `thread` key while replies arrive on `dm`.
|
||||
# Telegram private-chat DM topics use the DM-topic source shape (user_id == chat_id) so the
|
||||
# synthetic turn binds the same key later inbound turns arrive on (`dm`, not `thread`).
|
||||
is_telegram_private_chat = (
|
||||
platform == Platform.TELEGRAM and looks_like_telegram_private_chat_id(home_chat_id)
|
||||
)
|
||||
@@ -1501,14 +1439,11 @@ class GatewayStartupMixin:
|
||||
)
|
||||
|
||||
def _handoff_session_key(self, dest, profile_name: Optional[str]) -> str:
|
||||
"""Build the destination session_key with the adapters' own rules.
|
||||
|
||||
Thread keys omit user_id (thread_sessions_per_user default) so the next message shares it.
|
||||
The key is namespaced to the queuing profile: a multiplexed gateway would otherwise build
|
||||
``agent:main:...`` while the profile's adapter routes inbound on ``agent:<profile>:...``.
|
||||
The store resolver is only the root fallback (None when multiplexing is off; old key
|
||||
unchanged). The isinstance check is load-bearing: a Mock store returns a truthy MagicMock.
|
||||
"""
|
||||
"""Destination session_key by the adapters' own rules. Thread keys omit user_id so the next
|
||||
message shares it. Namespaced to the queuing profile (else a multiplexed gateway builds
|
||||
``agent:main:...`` while the profile's adapter routes on ``agent:<profile>:...``); the store
|
||||
resolver is only the root fallback. The isinstance check is load-bearing: a Mock store returns
|
||||
a truthy MagicMock."""
|
||||
platform_cfg = dest.handoff_config.platforms.get(dest.platform)
|
||||
extra = platform_cfg.extra if platform_cfg else {}
|
||||
handoff_profile = profile_name if (profile_name and profile_name != "default") else None
|
||||
@@ -1524,20 +1459,13 @@ class GatewayStartupMixin:
|
||||
logger.debug("Handoff: could not resolve profile namespace", exc_info=True)
|
||||
return build_session_key(
|
||||
dest.source, group_sessions_per_user=extra.get("group_sessions_per_user", True),
|
||||
thread_sessions_per_user=extra.get("thread_sessions_per_user", False),
|
||||
profile=handoff_profile,
|
||||
thread_sessions_per_user=extra.get("thread_sessions_per_user", False), profile=handoff_profile,
|
||||
)
|
||||
|
||||
async def _process_handoff(
|
||||
self, row: Dict[str, Any], profile_name: Optional[str] = None,
|
||||
) -> None:
|
||||
"""Execute one handoff row. Raises on failure (caller marks failed).
|
||||
|
||||
``profile_name`` (``None`` = root) is the profile whose store queued this handoff. Under
|
||||
multiplex it is load-bearing: ``self.adapters``/``self.config`` are the primary's (secondaries
|
||||
live in ``_profile_adapters``), and the session key must be namespaced ``agent:<profile>:...``
|
||||
or it binds a key nobody reads. Passing the name beats re-deriving it from the contextvar.
|
||||
"""
|
||||
async def _process_handoff(self, row: Dict[str, Any], profile_name: Optional[str] = None) -> None:
|
||||
"""Execute one handoff row; raises on failure (caller marks failed). ``profile_name`` (None =
|
||||
root) is the profile whose store queued it — load-bearing under multiplex: secondaries live in
|
||||
``_profile_adapters`` and the key must be namespaced ``agent:<profile>:...`` or nobody reads it."""
|
||||
cli_session_id = row["id"]
|
||||
dest = await self._handoff_resolve_destination(row, profile_name)
|
||||
session_key = self._handoff_session_key(dest, profile_name)
|
||||
@@ -1570,8 +1498,7 @@ class GatewayStartupMixin:
|
||||
logger.info(
|
||||
"Handoff: dispatching synthetic turn for CLI session %s → %s "
|
||||
"(home=%s, thread=%s, session_key=%s)",
|
||||
cli_session_id, dest.platform_name, dest.home.chat_id, dest.effective_thread_id,
|
||||
session_key,
|
||||
cli_session_id, dest.platform_name, dest.home.chat_id, dest.effective_thread_id, session_key,
|
||||
)
|
||||
# Dispatch through the runner directly: adapter.handle_message would spawn a background task
|
||||
# and lose error visibility; inline _handle_message keeps success/failure observable.
|
||||
|
||||
Reference in New Issue
Block a user