Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions cli-config.yaml.example
Original file line number Diff line number Diff line change
Expand Up @@ -1219,6 +1219,40 @@ platform_toolsets:
# # Route scripts default to a 30 second timeout. Scripts must live under
# # the active profile's scripts directory and receive webhook JSON on stdin.
# script_timeout_seconds: 30
# routes:
# # Orca completion bridge: a LOCAL Orca instance wakes Hermes the
# # moment a run reports in, instead of Hermes finding out on its next
# # poll. The body is data, never a prompt or a command, and no payload
# # can declare itself complete — Hermes re-queries Orca by run_id and
# # only Orca's own ledger completes a run. `Stop` and `StopFailure`
# # are hook events: recorded and ignored, never completions. And the
# # worker ledger has a veto — only workerState `succeeded` completes,
# # while `failed`/`stopped`/`abandoned`/`start_unknown`/`stop_unknown`
# # deliver nothing at all. A run with no worker ledger (no dispatches,
# # or context-only ones reported as `unsupervised`) is decided by the
# # Task ledger alone.
# #
# # Constraints (enforced at startup and per request):
# # * HMAC secret required; INSECURE_NO_AUTH is refused for this kind
# # * the request must AUTHENTICATE under a replay-protected scheme
# # (X-Webhook-Signature-V2 + X-Webhook-Timestamp, or Svix); the
# # GitHub/GitLab-token/V1 branches stay available to every OTHER
# # route but are not evaluated here at all, so a junk V2 header
# # alongside a valid GitHub signature does not get in
# # * loopback bind — omit `host` and it pins itself to 127.0.0.1;
# # a non-loopback `host` with an orca_bridge route fails to start
# # * off-machine peers are rejected, and so is any request carrying
# # a forwarding header (the bridge must not be proxied)
# # * config.yaml only — `hermes webhook subscribe` cannot mint one
# #
# # Register a run from the conversation that launched it so the report
# # comes back to that thread:
# # hermes webhook orca-register --run-id run_xxx --goal "..."
# # Then have the Orca hook notify:
# # hermes webhook orca-notify --run-id run_xxx --event worker_done
# orca:
# orca_bridge: true
# secret: "${ORCA_BRIDGE_SECRET}"
#
# Discord-specific settings (config.yaml top-level, not under platforms:):
#
Expand Down
315 changes: 287 additions & 28 deletions gateway/platforms/webhook.py

Large diffs are not rendered by default.

165 changes: 165 additions & 0 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@

import asyncio
import concurrent.futures
import contextlib
import dataclasses
import faulthandler
import inspect
Expand Down Expand Up @@ -6734,6 +6735,105 @@ def _approval_notify_sync(approval_data: dict) -> None:
# DB-backed commands and is how many suites construct a bare runner). A plain
# ``None`` cannot express both. Mirrors ``gateway.session._DB_UNPINNED``.
_SESSION_DB_UNPINNED = object()
class OrcaBridgeSupervisor:
"""Lifecycle for the Orca completion bridge and the listener it rides on.

The bridge has no socket of its own — it is reached through the webhook
adapter's loopback route — and that coupling is the whole reason these two
are supervised as a pair rather than as independent subsystems. A live
listener whose bridge failed to start would accept authenticated,
replay-protected completion events and then drop them on the floor, which
is strictly worse than not listening: the sender gets a 5xx it can retry
from, or a 200 for work nobody recorded.

``webhook`` is the already-connected adapter and ``bridge`` is the
``tools.orca_bridge`` module (injected so the ordering can be tested
without a socket or a database).
"""

def __init__(self, adapters, platform, webhook=None, bridge=None):
# The gateway's live platform map. Rollback removes the webhook entry
# from it so nothing downstream routes to a listener we just closed.
self.adapters = adapters
self.platform = platform
self.webhook = webhook
self.bridge = bridge
self._sweep_task = None

async def start(self) -> bool:
"""Bring the bridge up behind an already-started listener.

Returns False after rolling the listener back, rather than raising:
an Orca bridge that cannot open its durable state is a reason to run
the gateway without the bridge, not a reason to take every other
platform down with it.
"""
try:
self.bridge.start()
except Exception as exc: # noqa: BLE001
logger.error(
"Orca bridge failed to start (%s); shutting the webhook "
"listener back down so it cannot accept events the bridge "
"would drop", exc, exc_info=True,
)
await self.shutdown()
self.adapters.pop(self.platform, None)
return False

# A run that finished while the gateway was down never gets its POST
# redelivered, so recovery cannot depend on the sender — ask Orca.
self._sweep_task = asyncio.create_task(self._startup_sweep())
return True

async def _startup_sweep(self) -> None:
"""Reconcile runs that ended while this gateway was not running.

Runs once, off the event loop, and never propagates: a missing `orca`
binary or an unreachable runtime is a reason to keep serving webhooks,
not a reason to fail startup.

ponytail: the sweep runs on a worker thread, so cancelling this task
unblocks shutdown but does not interrupt an in-flight `orca` call —
the subprocess timeout (20s per query) is the real ceiling. Good
enough while the sweep is a bounded one-shot over registered runs; if
it ever becomes periodic, give it a cancellable subprocess instead.
"""
try:
published = await asyncio.to_thread(self.bridge.sweep)
if published:
logger.info(
"Orca startup sweep recovered %d completed run(s)",
published,
)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001
logger.warning("Orca startup sweep failed: %s", exc)

async def shutdown(self) -> None:
"""Close the listener first, then the bridge behind it.

Order is load-bearing. Stopping the bridge while the socket is still
accepting would leave a window where an authenticated event is taken
in and then refused by a bridge that has already closed its state —
the sender sees a failure for a request that was genuinely delivered.
Closing the listener first means no new event can arrive after the
bridge stops recording.
"""
task, self._sweep_task = self._sweep_task, None
if task is not None:
task.cancel()
# Cancellation is the expected outcome and anything else the sweep
# raised is already logged inside it. Neither may abort the rest of
# shutdown, so await it for the join and swallow both.
with contextlib.suppress(asyncio.CancelledError, Exception):
await task
if self.webhook:
await self.webhook.disconnect()
self.webhook = None
if self.bridge:
self.bridge.stop()
self.bridge = None


class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, GatewaySlashCommandsMixin):
Expand Down Expand Up @@ -6873,6 +6973,9 @@ def __init__(self, config: Optional[GatewayConfig] = None):
# user knows persistence is broken instead of discovering it later
# via a missing /resume or empty history (#88235).
self._session_db_init_error: Optional[str] = None
# Set by _start_orca_bridge() only when a webhook Orca bridge route
# is configured; stays None on every other deployment.
self._orca_supervisor: Optional["OrcaBridgeSupervisor"] = None
# Multi-profile multiplexing: adapters for NON-default profiles live
# here, keyed by profile name then Platform. self.adapters stays the
# default/active profile's map so the ~93 existing self.adapters[...]
Expand Down Expand Up @@ -12489,6 +12592,50 @@ def _start_loop_heartbeat_task(self) -> None:
self._loop_heartbeat_task.add_done_callback(_bg.discard)
except Exception:
logger.debug("Failed to start gateway loop heartbeat", exc_info=True)
async def _start_orca_bridge(self) -> bool:
"""Start the Orca completion bridge if the webhook serves one.

No-op — and no import of ``tools.orca_bridge`` at all — unless the
connected webhook adapter actually declares an ``orca_bridge`` route.
"""
adapter = self.adapters.get(Platform.WEBHOOK)
routes = getattr(adapter, "orca_bridge_routes", None)
if adapter is None or not callable(routes) or not routes():
return False

from tools import orca_bridge

self._orca_supervisor = OrcaBridgeSupervisor(
self.adapters, Platform.WEBHOOK, webhook=adapter, bridge=orca_bridge
)
started = await self._orca_supervisor.start()
if not started:
self._orca_supervisor = None
self._update_platform_runtime_status(
Platform.WEBHOOK.value,
platform_state="fatal",
error_code="orca_bridge_start_failed",
error_message="Orca completion bridge could not start",
)
return False
logger.info("✓ Orca completion bridge started")
return True

async def _stop_orca_bridge(self) -> None:
"""Close the listener and then the bridge. Never raises.

getattr-guard for the same reason ``stop()`` uses one: shutdown-path
tests build bare runners via ``object.__new__``, which never ran
``__init__`` and so have no ``_orca_supervisor`` attribute at all.
"""
supervisor = getattr(self, "_orca_supervisor", None)
self._orca_supervisor = None
if supervisor is None:
return
try:
await supervisor.shutdown()
except Exception as exc: # noqa: BLE001
logger.warning("Orca bridge shutdown error: %s", exc)

async def start(self) -> bool:
"""
Expand Down Expand Up @@ -13252,6 +13399,12 @@ async def _connect_one_startup(p, p_cfg, adp):
# Update delivery router with adapters
if await self._abort_startup_if_shutdown_requested():
return True
# Orca completion bridge — started only when the webhook platform is
# actually serving a bridge route. Must run BEFORE the router is
# published, so a rollback removes the listener before anything can
# route to it.
await self._start_orca_bridge()

self.delivery_router.adapters = self.adapters
self._wire_teams_pipeline_runtime()

Expand Down Expand Up @@ -15152,6 +15305,18 @@ def _phase_elapsed() -> float:
)
if cancel_completion_batches is not None:
await cancel_completion_batches()
# Orca bridge first: it closes the webhook listener and only then
# stops the bridge, so no authenticated completion event can be
# accepted into a bridge that has already stopped recording. The
# generic teardown below still runs for the webhook adapter —
# adapter.disconnect() is idempotent by contract.
#
# getattr/callable-guarded like the liveness guards and agent
# cache above: shutdown-path tests drive this body with partial
# runner doubles that implement only what they assert on.
_stop_bridge = getattr(self, "_stop_orca_bridge", None)
if callable(_stop_bridge):
await _stop_bridge()

for platform, adapter in list(self.adapters.items()):
await self._bounded_adapter_teardown(adapter, platform)
Expand Down
58 changes: 58 additions & 0 deletions hermes_cli/subcommands/webhook.py
Original file line number Diff line number Diff line change
Expand Up @@ -80,4 +80,62 @@ def build_webhook_parser(subparsers, *, cmd_webhook: Callable) -> None:
"--payload", default="", help="JSON payload to send (default: test payload)"
)

# ── Orca completion bridge ──────────────────────────────────────────
wh_orca_reg = webhook_subparsers.add_parser(
"orca-register",
help="Register an Orca run so its completion wakes this conversation",
)
wh_orca_reg.add_argument("--run-id", required=True, help="Orca run id")
wh_orca_reg.add_argument(
"--goal", default="", help="What the run was asked to do"
)
wh_orca_reg.add_argument(
"--session-key",
default="",
help="Gateway session key to report back to (default: the live "
"HERMES_SESSION_KEY, so the completion returns to the thread that "
"launched the run)",
)
wh_orca_reg.add_argument(
"--worktree", default="", help="Worktree path, for the report"
)
wh_orca_reg.add_argument(
"--terminal", default="", help="Orca terminal handle, for the report"
)

wh_orca_runs = webhook_subparsers.add_parser(
"orca-runs", help="List registered Orca runs"
)
wh_orca_runs.add_argument(
"--state", default="", help="Filter by state (open, completed)"
)

webhook_subparsers.add_parser(
"orca-sweep",
help="Re-query Orca for every open run and deliver any that finished",
)

wh_orca_notify = webhook_subparsers.add_parser(
"orca-notify",
help="Send a signed completion notification to the local bridge",
)
wh_orca_notify.add_argument("--run-id", required=True, help="Orca run id")
wh_orca_notify.add_argument(
"--event", default="worker_done",
help="Signal kind (worker_done, hermes-ready, exit, stop, ...)",
)
wh_orca_notify.add_argument(
"--route", default="orca", help="Webhook route name for the bridge"
)
wh_orca_notify.add_argument(
"--secret", default="", help="HMAC secret (default: from config.yaml)"
)
wh_orca_notify.add_argument(
"--event-id", default="", help="Sender-supplied id used for dedupe"
)
wh_orca_notify.add_argument(
"--sequence", type=int, default=-1,
help="Monotonic sequence number for out-of-order rejection",
)

webhook_parser.set_defaults(func=cmd_webhook)
Loading