Files
ai-agent/sidecar/outbox.py
Gabriel Vidal 2fc19c2470 feat(notify): one notify CLI (send/ask/form/wait) that also runs on workers via a sidecar outbox
cli/notify.py merges the lab's notify-done, notify-ask and ask-form scripts
into one stdlib-only CLI shipped with this repo (the aliases stay as shims),
so a session on a remote worker has it too. Two transports: the hub directly
on the lab, or — on a worker — the sidecar's new /outbox, which the hub's
OutboxMirror polls every 2 s, creating the record through the same functions
/api/notify, /api/ask and /api/forms run and pushing the answer back. The
worker still never calls the hub.

The hub now owns what the CLI used to compute: the default conversation URL
(worker transcripts included), the `notified` meta stamp, and a `--final`
push marks the session finished. A push's --action-cmd button answers on
/api/notify/{id}/action like an ask (HA integration updated). Ids are minted
by the CLI (ntf_/ask_/frm_) so viewer, phone and worker agree. Dropped: the
direct-HA fallback, the DONE button, the Forge cost line, the dashboard
mirror. The viewer's widgets also parse the `notify send|ask|form|wait`
spelling; the worker installer links the skill into ~/.claude/skills.

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

322 lines
13 KiB
Python

"""
The outbox — how a session on a *remote worker* reaches Gabriel without the
worker ever calling the hub.
A run on the Mac has no Home Assistant, no hub URL and no hub credential (that
is the worker's whole security shape — see pairing.py). So the ``notify`` CLI
it runs (``cli/notify.py``, transport ``outbox``) drops its notification / ask
/ form into **this sidecar**, and the hub, which already polls every paired
worker for transcripts, collects it on the same cadence:
run ──POST /outbox {kind, payload}──▶ sidecar (queued)
hub ──GET /outbox?pending=1────────▶ sidecar (bearer) every ~2 s
hub ──POST /outbox/{id}/ack {…}─────▶ sidecar (bearer) recorded + forwarded
hub ──POST /outbox/{id}/answer {…}──▶ sidecar (bearer) when Gabriel answers
run ──GET /outbox/{id}?waitSecs=55─▶ sidecar long-poll for that answer
The run authenticates with ``AI_AGENT_OUTBOX_TOKEN``, a secret this sidecar
mints once (``WORKER_STATE_DIR/outbox-token``) and exports to every run it
launches — so nothing else on the tailnet can push notifications through the
worker, and a run that outlives a sidecar restart keeps working. The hub uses
its pairing bearer, as for every other endpoint.
Ids are minted by the CLI in the hub's own format (``ntf_`` / ``ask_`` /
``frm_`` + 12 hex) and the hub records the item **under that id**, so an
``ASK_ID=ask_…`` printed on the worker is the id the viewer and the phone
know. State: ``WORKER_STATE_DIR/outbox.json`` (atomic rewrite, capped), the
same discipline as the hub's own notify store. Waiting is in-memory (one
``threading.Event`` per item); a restart just drops the wait and the CLI's
poll loop re-polls.
"""
from __future__ import annotations
import json
import os
import pathlib
import re
import secrets
import threading
import time
from typing import Callable
from fastapi import APIRouter, Header, HTTPException, Query
from pydantic import BaseModel
KINDS = ("notify", "ask", "form")
_ID_RE = re.compile(r"^(ntf|ask|frm)_[0-9a-f]{12}$")
_MAX_ITEMS = 500
_MAX_WAIT_S = 55.0
# A queued item the hub never collected (worker unpaired, hub down for days)
# is dropped after this long so the file can't grow forever.
STALE_S = 7 * 24 * 3600
TERMINAL = ("answered", "submitted", "acted", "expired", "cancelled", "failed")
def _now_iso() -> str:
return time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
class Outbox:
def __init__(self, path: pathlib.Path):
self.path = pathlib.Path(path)
self._lock = threading.Lock()
self._waiters: dict[str, threading.Event] = {}
self._items: list[dict] = []
self._load()
# ── persistence ────────────────────────────────────────────────────────
def _load(self) -> None:
try:
data = json.loads(self.path.read_text(encoding="utf-8"))
self._items = list(data.get("items") or [])
except (OSError, ValueError):
self._items = []
def _save(self) -> None:
self.path.parent.mkdir(parents=True, exist_ok=True)
tmp = self.path.with_name(self.path.name + ".tmp")
fd = os.open(tmp, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w", encoding="utf-8") as f:
json.dump({"items": self._items[-_MAX_ITEMS:]}, f)
os.replace(tmp, self.path)
def _find(self, oid: str) -> dict | None:
for it in self._items:
if it.get("id") == oid:
return it
return None
# ── the run's side ─────────────────────────────────────────────────────
def put(self, kind: str, payload: dict) -> dict:
oid = str(payload.get("id") or "")
if kind not in KINDS or not _ID_RE.match(oid):
raise ValueError("kind must be notify|ask|form and payload.id a "
"ntf_/ask_/frm_ id")
item = {"id": oid, "kind": kind, "payload": payload, "status": "queued",
"at": _now_iso(), "collectedAt": None, "hub": None}
with self._lock:
if self._find(oid):
raise ValueError(f"duplicate id {oid}")
self._items.append(item)
self._waiters[oid] = threading.Event()
self._gc()
self._save()
return json.loads(json.dumps(item))
def get(self, oid: str) -> dict | None:
with self._lock:
it = self._find(oid)
return json.loads(json.dumps(it)) if it else None
def wait(self, oid: str, wait_s: float) -> dict | None:
with self._lock:
ev = self._waiters.get(oid)
it = self._find(oid)
if it and it.get("status") not in TERMINAL and ev and wait_s > 0:
ev.wait(min(wait_s, _MAX_WAIT_S))
return self.get(oid)
def cancel(self, oid: str) -> dict | None:
"""The run gave up (``notify cancel``): a still-queued item is dropped
before the hub sees it; a collected one is flagged so the hub cancels
its own record on the next pull."""
with self._lock:
it = self._find(oid)
if not it:
return None
if it.get("status") in TERMINAL:
return json.loads(json.dumps(it))
it["status"] = "cancelled"
it["cancelRequested"] = it.get("collectedAt") is not None
self._save()
ev = self._waiters.pop(oid, None)
if ev:
ev.set()
return self.get(oid)
# ── the hub's side ─────────────────────────────────────────────────────
def pending(self) -> list[dict]:
"""What the hub must act on: queued items (to create), collected ones
still open (to check for an answer) and cancel requests."""
with self._lock:
return [json.loads(json.dumps(it)) for it in self._items
if it.get("status") == "queued"
or (it.get("status") == "sent")
or (it.get("status") == "cancelled" and it.get("cancelRequested"))]
def ack(self, oid: str, hub: dict) -> dict | None:
"""The hub recorded the item: remember what it answered (delivery,
the conversation url it chose…). A ``failed`` ack ends the item."""
with self._lock:
it = self._find(oid)
if not it:
return None
if it.get("status") == "cancelled":
it["cancelRequested"] = False # the hub saw the cancel
elif it.get("status") == "queued":
it["status"] = "failed" if hub.get("failed") else "sent"
it["collectedAt"] = _now_iso()
it["hub"] = hub
self._save()
ev = self._waiters.get(oid)
terminal = it["status"] in TERMINAL
if terminal:
self._waiters.pop(oid, None)
if terminal and ev:
ev.set()
return self.get(oid)
def answer(self, oid: str, result: dict) -> dict | None:
"""The hub pushes the terminal state (answered / submitted / acted /
expired / cancelled) with the answer fields the CLI prints."""
status = str(result.get("status") or "")
if status not in TERMINAL:
raise ValueError(f"status must be one of {', '.join(TERMINAL)}")
with self._lock:
it = self._find(oid)
if not it:
return None
it["status"] = status
it["cancelRequested"] = False
for k in ("answer", "answerIndex", "answers", "actionIndex",
"actionTitle"):
if k in result:
it[k] = result[k]
it["endedAt"] = _now_iso()
self._save()
ev = self._waiters.pop(oid, None)
if ev:
ev.set()
return self.get(oid)
def _gc(self) -> None:
cutoff = time.time() - STALE_S
keep = []
for it in self._items:
try:
at = time.mktime(time.strptime(it.get("at") or "", "%Y-%m-%dT%H:%M:%SZ")) \
- time.timezone
except (ValueError, OverflowError):
at = time.time()
if it.get("status") in TERMINAL or at >= cutoff or it.get("status") == "sent":
keep.append(it)
self._items = keep[-_MAX_ITEMS:]
# ── the run token ─────────────────────────────────────────────────────────────
def run_token(state_dir: pathlib.Path) -> str:
"""The secret a run presents to POST into the outbox — minted once per
worker, 0600, exported to every launched run as AI_AGENT_OUTBOX_TOKEN."""
p = pathlib.Path(state_dir) / "outbox-token"
try:
tok = p.read_text().strip()
if tok:
return tok
except OSError:
pass
tok = secrets.token_hex(24)
p.parent.mkdir(parents=True, exist_ok=True)
fd = os.open(p.with_name("outbox-token.tmp"), os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, "w") as f:
f.write(tok)
os.replace(p.with_name("outbox-token.tmp"), p)
return tok
# ── router ────────────────────────────────────────────────────────────────────
class PutBody(BaseModel):
kind: str
payload: dict
class AckBody(BaseModel):
hub: dict
class AnswerBody(BaseModel):
status: str
answer: str | None = None
answerIndex: int | None = None
answers: dict | None = None
actionIndex: int | None = None
actionTitle: str | None = None
def make_router(box: Outbox, hub_auth: Callable[[str | None], None],
run_token_value: Callable[[], str]) -> APIRouter:
import hmac
r = APIRouter()
def _run_auth(authorization: str | None) -> None:
tok = run_token_value()
if not tok or not hmac.compare_digest(authorization or "", f"Bearer {tok}"):
raise HTTPException(401, "unauthorized (run token)")
def _either(authorization: str | None) -> None:
try:
hub_auth(authorization)
except HTTPException:
_run_auth(authorization)
@r.post("/outbox", status_code=202)
def put(body: PutBody, authorization: str | None = Header(default=None)) -> dict:
_run_auth(authorization)
try:
it = box.put(body.kind, body.payload)
except ValueError as e:
raise HTTPException(400, str(e))
return {"queued": True, "id": it["id"], "status": it["status"]}
@r.get("/outbox")
def pending(authorization: str | None = Header(default=None)) -> dict:
hub_auth(authorization)
items = box.pending()
return {"items": items, "count": len(items), "now": time.time()}
@r.get("/outbox/{oid}")
def get(oid: str, waitSecs: float = Query(0.0),
authorization: str | None = Header(default=None)) -> dict:
_either(authorization)
it = box.wait(oid, waitSecs) if waitSecs > 0 else box.get(oid)
if not it:
raise HTTPException(404, "unknown outbox item")
# Flatten what the CLI's wait loop reads, next to the raw item.
out = {k: it.get(k) for k in ("id", "kind", "status", "at", "hub",
"answer", "answerIndex", "answers",
"actionIndex", "actionTitle")}
return out
@r.post("/outbox/{oid}/cancel")
def cancel(oid: str, authorization: str | None = Header(default=None)) -> dict:
_either(authorization)
it = box.cancel(oid)
if not it:
raise HTTPException(404, "unknown outbox item")
return {"id": oid, "status": it["status"]}
@r.post("/outbox/{oid}/ack")
def ack(oid: str, body: AckBody,
authorization: str | None = Header(default=None)) -> dict:
hub_auth(authorization)
it = box.ack(oid, body.hub)
if not it:
raise HTTPException(404, "unknown outbox item")
return {"id": oid, "status": it["status"]}
@r.post("/outbox/{oid}/answer")
def answer(oid: str, body: AnswerBody,
authorization: str | None = Header(default=None)) -> dict:
hub_auth(authorization)
try:
it = box.answer(oid, {k: v for k, v in body.model_dump().items()
if v is not None})
except ValueError as e:
raise HTTPException(400, str(e))
if not it:
raise HTTPException(404, "unknown outbox item")
return {"id": oid, "status": it["status"]}
return r