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>
177 lines
7.6 KiB
Python
177 lines
7.6 KiB
Python
"""The outbox mirror against a fake worker: queued items are created through
|
|
the hub's own functions and acked, answers are pushed back, cancels honoured.
|
|
``python -m unittest test_outbox_mirror`` from ``backend/``."""
|
|
|
|
import json
|
|
import threading
|
|
import unittest
|
|
from http.server import BaseHTTPRequestHandler, HTTPServer
|
|
|
|
from fastapi import HTTPException
|
|
|
|
import outbox_mirror
|
|
|
|
|
|
class FakeOutbox(BaseHTTPRequestHandler):
|
|
"""A worker's /outbox surface: an in-memory item list the test pokes."""
|
|
items: list[dict] = []
|
|
calls: list[tuple[str, str, dict]] = []
|
|
token = "w0rker"
|
|
|
|
def log_message(self, *a):
|
|
pass
|
|
|
|
def _body(self) -> dict:
|
|
n = int(self.headers.get("Content-Length") or 0)
|
|
return json.loads(self.rfile.read(n) or b"{}") if n else {}
|
|
|
|
def _send(self, code: int, obj: dict) -> None:
|
|
data = json.dumps(obj).encode()
|
|
self.send_response(code)
|
|
self.send_header("Content-Type", "application/json")
|
|
self.send_header("Content-Length", str(len(data)))
|
|
self.end_headers()
|
|
self.wfile.write(data)
|
|
|
|
def do_GET(self):
|
|
if self.headers.get("Authorization") != f"Bearer {self.token}":
|
|
return self._send(401, {"detail": "unauthorized"})
|
|
self._send(200, {"items": list(self.items), "count": len(self.items)})
|
|
|
|
def do_POST(self):
|
|
if self.headers.get("Authorization") != f"Bearer {self.token}":
|
|
return self._send(401, {"detail": "unauthorized"})
|
|
body = self._body()
|
|
FakeOutbox.calls.append(("POST", self.path, body))
|
|
oid = self.path.split("/")[2]
|
|
for it in self.items:
|
|
if it["id"] == oid:
|
|
if self.path.endswith("/ack"):
|
|
it["status"] = "failed" if body["hub"].get("failed") else "sent"
|
|
it["hub"] = body["hub"]
|
|
elif self.path.endswith("/answer"):
|
|
it.update(body)
|
|
return self._send(200, {"id": oid, "status": it["status"]})
|
|
self._send(404, {"detail": "unknown"})
|
|
|
|
|
|
class FakeStore:
|
|
def __init__(self, w):
|
|
self.w, self.errors = w, []
|
|
|
|
def list(self):
|
|
return [self.w]
|
|
|
|
def online(self, wid):
|
|
return True
|
|
|
|
def set_status(self, wid, **st):
|
|
self.errors.append(st)
|
|
|
|
|
|
class MirrorTest(unittest.TestCase):
|
|
@classmethod
|
|
def setUpClass(cls):
|
|
cls.srv = HTTPServer(("127.0.0.1", 0), FakeOutbox)
|
|
threading.Thread(target=cls.srv.serve_forever, daemon=True).start()
|
|
cls.url = f"http://127.0.0.1:{cls.srv.server_port}"
|
|
|
|
@classmethod
|
|
def tearDownClass(cls):
|
|
cls.srv.shutdown()
|
|
|
|
def setUp(self):
|
|
FakeOutbox.items, FakeOutbox.calls = [], []
|
|
self.records: dict[str, dict] = {}
|
|
self.created: list[dict] = []
|
|
|
|
def create(kind):
|
|
def fn(d):
|
|
self.created.append({"kind": kind, **d})
|
|
if kind == "form" and not d.get("title"):
|
|
raise HTTPException(400, "need a title")
|
|
self.records[d["id"]] = {"id": d["id"], "status": "pending"}
|
|
return {"id": d["id"], "ok": True, "url": "https://v/c/x"}
|
|
return fn
|
|
self.store = FakeStore({"id": "w1", "url": self.url, "token": "w0rker"})
|
|
self.m = outbox_mirror.OutboxMirror(
|
|
self.store,
|
|
create={k: create(k) for k in ("notify", "ask", "form")},
|
|
status={p: (lambda i: self.records.get(i)) for p in ("ntf", "ask", "frm")},
|
|
cancel={"frm": lambda i: self.records.get(i) and
|
|
self.records[i].update(status="cancelled") or self.records.get(i)},
|
|
publish=lambda ev: None)
|
|
|
|
def test_queued_item_is_created_once_and_acked(self):
|
|
FakeOutbox.items = [{"id": "ask_0123456789ab", "kind": "ask", "status": "queued",
|
|
"payload": {"id": "ask_0123456789ab", "question": "Q?",
|
|
"options": ["A", "B"]}}]
|
|
self.m.tick()
|
|
self.assertEqual(len(self.created), 1)
|
|
self.assertEqual(self.created[0]["worker"], "w1")
|
|
self.assertEqual(FakeOutbox.items[0]["status"], "sent")
|
|
self.assertEqual(FakeOutbox.items[0]["hub"]["url"], "https://v/c/x")
|
|
# next tick: it's "sent" and still pending hub-side → nothing pushed
|
|
self.m.tick()
|
|
self.assertEqual(len(self.created), 1)
|
|
self.assertEqual([c[1] for c in FakeOutbox.calls],
|
|
["/outbox/ask_0123456789ab/ack"])
|
|
|
|
def test_answer_is_pushed_once_the_record_ends(self):
|
|
FakeOutbox.items = [{"id": "ask_0123456789ab", "kind": "ask", "status": "sent"}]
|
|
self.records["ask_0123456789ab"] = {"id": "ask_0123456789ab", "status": "pending"}
|
|
self.m.tick()
|
|
self.assertEqual(FakeOutbox.calls, [])
|
|
self.records["ask_0123456789ab"].update(status="answered", answer="B",
|
|
answerIndex=1)
|
|
self.m.tick()
|
|
self.assertEqual(FakeOutbox.calls[-1][1:], ("/outbox/ask_0123456789ab/answer",
|
|
{"status": "answered", "answer": "B",
|
|
"answerIndex": 1}))
|
|
|
|
def test_form_submit_and_unknown_record(self):
|
|
FakeOutbox.items = [{"id": "frm_0123456789ab", "kind": "form", "status": "sent"},
|
|
{"id": "ask_ffffffffffff", "kind": "ask", "status": "sent"}]
|
|
self.records["frm_0123456789ab"] = {"status": "submitted", "answers": {"a": 1}}
|
|
self.m.tick()
|
|
paths = {c[1]: c[2] for c in FakeOutbox.calls}
|
|
self.assertEqual(paths["/outbox/frm_0123456789ab/answer"],
|
|
{"status": "submitted", "answers": {"a": 1}})
|
|
# a record the hub doesn't know (lost store) ends the wait as expired
|
|
self.assertEqual(paths["/outbox/ask_ffffffffffff/answer"], {"status": "expired"})
|
|
|
|
def test_create_failure_is_acked_as_failed(self):
|
|
FakeOutbox.items = [{"id": "frm_0123456789ab", "kind": "form", "status": "queued",
|
|
"payload": {"id": "frm_0123456789ab", "title": "", "fields": []}}]
|
|
self.m.tick()
|
|
self.assertEqual(FakeOutbox.items[0]["status"], "failed")
|
|
self.assertIn("need a title", FakeOutbox.items[0]["hub"]["error"])
|
|
|
|
def test_already_existing_record_is_not_recreated(self):
|
|
self.records["ntf_0123456789ab"] = {"id": "ntf_0123456789ab"}
|
|
FakeOutbox.items = [{"id": "ntf_0123456789ab", "kind": "notify", "status": "queued",
|
|
"payload": {"id": "ntf_0123456789ab", "title": "t", "message": "m"}}]
|
|
self.m.tick()
|
|
self.assertEqual(self.created, [])
|
|
self.assertEqual(FakeOutbox.items[0]["hub"].get("existing"), True)
|
|
|
|
def test_cancel_request_cancels_the_hub_record(self):
|
|
self.records["frm_0123456789ab"] = {"id": "frm_0123456789ab", "status": "pending"}
|
|
FakeOutbox.items = [{"id": "frm_0123456789ab", "kind": "form",
|
|
"status": "cancelled", "cancelRequested": True}]
|
|
self.m.tick()
|
|
self.assertEqual(self.records["frm_0123456789ab"]["status"], "cancelled")
|
|
self.assertEqual(FakeOutbox.calls[-1][1:], ("/outbox/frm_0123456789ab/ack",
|
|
{"hub": {"cancelled": True}}))
|
|
|
|
def test_worker_down_backs_off(self):
|
|
store = FakeStore({"id": "w2", "url": "http://127.0.0.1:1", "token": "x"})
|
|
m = outbox_mirror.OutboxMirror(store, create={}, status={}, cancel={})
|
|
m.tick()
|
|
self.assertTrue(store.errors and store.errors[-1]["lastError"].startswith("outbox:"))
|
|
self.assertIn("w2", m._backoff)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|