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>
179 lines
7.8 KiB
Python
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
|