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>
322 lines
13 KiB
Python
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
|