Files
ai-agent/backend/workers/outbox_mirror.py
Gabriel Vidal 0ff9e40242 refactor(backend): split main.py into domain packages with app.state injection
main.py (5146 lines, 106 routes) becomes an assembly only: one package per
domain — core, files, settings, dashboards, conversations, diff,
notifications, forms, runs, accounts, workers, models, cron, agents,
projects, services, goals, memories, plans, templates — each exposing an
APIRouter; the flat domain modules move into their package behind a barrel
that keeps the old `import conversations` / `import projects` spellings.

The shared singletons (store, meta_store, hub, indexer, …) are built once by
core.state.build_state() and attached to app.state.ai; routes take them as
the `deps: State` dependency and helpers as an explicit `deps: AppState`.
conversations/pricing.py carries the per-model rates out of the parser.

Verified: route table and OpenAPI byte-identical; 90 read endpoints
golden-diffed against the monolith on a copy of the live data (identical);
write routes smoke-tested; 66 backend tests pass.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-06 23:55:47 +02:00

179 lines
7.8 KiB
Python

"""
The outbox mirror — the hub's side of a worker's outbox (sidecar/outbox.py).
A session on a remote worker can't reach the hub (the worker holds no hub URL
or credential — the hub always calls in). So its ``notify`` CLI queues the
push / ask / form in the worker's own sidecar, and this thread — one more
poller next to ``FeedMirror`` — collects it:
every tick, per online worker:
GET /outbox → queued items + open ones + cancels
queued → create it here (the same functions /api/notify, /api/ask and
/api/forms run), then POST /outbox/{id}/ack {hub: result}
sent → look the record up; once answered / submitted / acted /
expired / cancelled, POST /outbox/{id}/answer {…} so the
worker's long-polling `notify wait` returns
cancelled (cancel requested after collection) → cancel the hub record,
ack
Everything is idempotent against a restart on either side: a queued item is
created at most once (the ack flips it to ``sent`` on the worker; a create
that already exists answers the existing id), and an answer is pushed again
until the worker confirms it. A worker that stops answering is backed off
(2 s → 60 s) until the monitor sees it again.
"""
from __future__ import annotations
import json
import threading
import time
import urllib.error
import urllib.request
from typing import Callable
from fastapi import HTTPException
class WorkerDown(Exception):
pass
class OutboxMirror(threading.Thread):
def __init__(self, store, *, create: dict[str, Callable[[dict], dict]],
status: dict[str, Callable[[str], dict | None]],
cancel: dict[str, Callable[[str], dict | None]],
publish: Callable[[dict], None] | None = None,
interval: float = 2.0):
super().__init__(daemon=True, name="outbox-mirror")
self.store = store
self.create, self.status, self.cancel = create, status, cancel
self.publish = publish or (lambda ev: None)
self.interval = interval
self._backoff: dict[str, tuple[float, float]] = {}
self._stop = threading.Event()
def stop(self) -> None:
self._stop.set()
# ── HTTP ──────────────────────────────────────────────────────────────
def _call(self, w: dict, path: str, payload: dict | None = None, *,
method: str = "POST", timeout: float = 10) -> dict:
data = json.dumps(payload).encode() if payload is not None else None
req = urllib.request.Request(
f"{w['url']}{path}", data=data, method=method,
headers={"Content-Type": "application/json",
"Authorization": f"Bearer {w['token']}"})
try:
with urllib.request.urlopen(req, timeout=timeout) as r:
return json.loads(r.read().decode() or "{}")
except urllib.error.HTTPError as e:
if e.code == 404:
raise FileNotFoundError(path)
raise WorkerDown(f"HTTP {e.code} on {path}")
except (urllib.error.URLError, OSError, ValueError) as e:
raise WorkerDown(str(e))
# ── one item ──────────────────────────────────────────────────────────
@staticmethod
def terminal(kind: str, rec: dict | None) -> dict | None:
"""The answer payload to push for a record that has ended, else None."""
if not rec:
return None
st = rec.get("status") or ""
if kind == "ask" and st in ("answered", "expired"):
return {"status": st, "answer": rec.get("answer"),
"answerIndex": rec.get("answerIndex")}
if kind == "frm" and st in ("submitted", "cancelled"):
return {"status": st, "answers": rec.get("answers") or {}}
if kind == "ntf" and st == "acted":
return {"status": st, "actionIndex": rec.get("actionIndex"),
"actionTitle": rec.get("actionTitle")}
return None
def _handle(self, w: dict, item: dict) -> None:
oid = str(item.get("id") or "")
kind = str(item.get("kind") or "")
st = item.get("status")
if st == "queued":
fn = self.create.get(kind)
if not fn:
self._call(w, f"/outbox/{oid}/ack", {"hub": {"failed": True,
"error": f"unknown kind {kind!r}"}})
return
payload = dict(item.get("payload") or {})
payload["id"] = oid
payload["worker"] = w["id"]
try:
# The record may already exist (ack lost on a previous pass):
# a create then 400s on the duplicate id → treat as created.
existing = self.status[oid[:3]](oid)
res = existing and {"id": oid, "ok": True, "existing": True} \
or fn(payload)
except HTTPException as e:
res = {"failed": True, "error": str(e.detail)[:300]}
self._call(w, f"/outbox/{oid}/ack", {"hub": {
k: v for k, v in res.items() if k in ("id", "ok", "url", "focusUrl",
"delivered", "failed", "error",
"existing")}})
elif st == "sent":
fn = self.status.get(oid[:3])
rec = fn(oid) if fn else None
if rec is None:
self._call(w, f"/outbox/{oid}/answer",
{"status": "expired"})
return
done = self.terminal(oid[:3], rec)
if done:
self._call(w, f"/outbox/{oid}/answer", done)
elif st == "cancelled":
fn = self.cancel.get(oid[:3])
rec = fn(oid) if fn else None
if rec is not None:
self.publish({"type": "form", "id": oid,
"sessionId": rec.get("sessionId") or ""})
self._call(w, f"/outbox/{oid}/ack", {"hub": {"cancelled": True}})
# ── loop ──────────────────────────────────────────────────────────────
def sync(self, w: dict) -> int:
try:
listing = self._call(w, "/outbox", method="GET")
except FileNotFoundError:
return 0 # a worker that predates the outbox — nothing to collect
n = 0
for item in listing.get("items") or []:
try:
self._handle(w, item)
n += 1
except FileNotFoundError:
continue
return n
def tick(self) -> int:
n = 0
now = time.monotonic()
for w in self.store.list():
wid = w["id"]
if not self.store.online(wid):
continue
until, step = self._backoff.get(wid, (0.0, 0.0))
if now < until:
continue
try:
n += self.sync(w)
self._backoff.pop(wid, None)
except WorkerDown as e:
step = min(60.0, max(2.0, step * 2))
self._backoff[wid] = (now + step, step)
self.store.set_status(wid, lastError=f"outbox: {e}"[:200])
except Exception as e: # never let one bad item kill the thread
self.store.set_status(wid, lastError=f"outbox: {e}"[:200])
return n
def run(self) -> None:
while not self._stop.wait(self.interval):
try:
self.tick()
except Exception:
pass