Files
ai-agent/sidecar/test_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

135 lines
5.0 KiB
Python

"""The outbox: a run queues, the hub collects + answers, the run's wait returns."""
import threading
import time
import pytest
from fastapi import FastAPI, HTTPException
from fastapi.testclient import TestClient
import outbox
HUB = "hub-t0ken"
RUN = "run-t0ken"
def _hub_auth(h):
if h != f"Bearer {HUB}":
raise HTTPException(401, "unauthorized")
@pytest.fixture()
def box(tmp_path):
return outbox.Outbox(tmp_path / "outbox.json")
@pytest.fixture()
def app(box):
a = FastAPI()
a.include_router(outbox.make_router(box, _hub_auth, lambda: RUN))
return a
@pytest.fixture()
def run(app):
return TestClient(app, headers={"Authorization": f"Bearer {RUN}"})
@pytest.fixture()
def hub(app):
return TestClient(app, headers={"Authorization": f"Bearer {HUB}"})
ASK = {"id": "ask_0123456789ab", "question": "Ship?", "options": ["Yes", "No"]}
def test_run_token_is_minted_once_and_private(tmp_path):
t1 = outbox.run_token(tmp_path)
t2 = outbox.run_token(tmp_path)
assert t1 == t2 and len(t1) >= 32
assert oct((tmp_path / "outbox-token").stat().st_mode & 0o777) == "0o600"
def test_put_needs_run_token_and_a_cli_id(run, app):
anon = TestClient(app)
assert anon.post("/outbox", json={"kind": "ask", "payload": ASK}).status_code == 401
r = run.post("/outbox", json={"kind": "ask", "payload": {"id": "nope"}})
assert r.status_code == 400
r = run.post("/outbox", json={"kind": "ask", "payload": ASK})
assert r.status_code == 202 and r.json() == {"queued": True, "id": ASK["id"],
"status": "queued"}
# the same id can't be queued twice
assert run.post("/outbox", json={"kind": "ask", "payload": ASK}).status_code == 400
def test_hub_collects_acks_and_answers(run, hub, box):
run.post("/outbox", json={"kind": "ask", "payload": ASK})
# the run can't list the queue, the hub can
assert run.get("/outbox").status_code == 401
items = hub.get("/outbox").json()["items"]
assert [i["id"] for i in items] == [ASK["id"]] and items[0]["status"] == "queued"
hub.post(f"/outbox/{ASK['id']}/ack", json={"hub": {"ok": True, "id": ASK["id"]}})
it = run.get(f"/outbox/{ASK['id']}").json()
assert it["status"] == "sent" and it["hub"]["ok"] is True
# still listed (open) so the hub keeps checking for the answer
assert hub.get("/outbox").json()["count"] == 1
r = hub.post(f"/outbox/{ASK['id']}/answer",
json={"status": "answered", "answer": "Yes", "answerIndex": 0})
assert r.json()["status"] == "answered"
it = run.get(f"/outbox/{ASK['id']}").json()
assert (it["status"], it["answer"], it["answerIndex"]) == ("answered", "Yes", 0)
assert hub.get("/outbox").json()["count"] == 0 # done — off the list
assert box.get(ASK["id"])["endedAt"]
def test_wait_returns_the_instant_the_answer_lands(box):
box.put("ask", ASK)
box.ack(ASK["id"], {"ok": True})
got = {}
def waiter():
got["it"] = box.wait(ASK["id"], 10)
th = threading.Thread(target=waiter)
t0 = time.monotonic()
th.start()
time.sleep(0.2)
box.answer(ASK["id"], {"status": "answered", "answer": "No", "answerIndex": 1})
th.join(5)
assert got["it"]["answer"] == "No" and time.monotonic() - t0 < 3
def test_failed_ack_ends_the_item(run, hub):
run.post("/outbox", json={"kind": "form", "payload": {"id": "frm_0123456789ab",
"title": "", "fields": []}})
hub.post("/outbox/frm_0123456789ab/ack", json={"hub": {"failed": True,
"error": "need a title"}})
it = run.get("/outbox/frm_0123456789ab").json()
assert it["status"] == "failed" and it["hub"]["error"] == "need a title"
assert hub.get("/outbox").json()["count"] == 0
def test_cancel_before_and_after_collection(run, hub):
a, b = dict(ASK, id="ask_aaaaaaaaaaaa"), dict(ASK, id="ask_bbbbbbbbbbbb")
run.post("/outbox", json={"kind": "ask", "payload": a})
run.post("/outbox", json={"kind": "ask", "payload": b})
# a: cancelled while still queued → simply gone from the hub's view
assert run.post(f"/outbox/{a['id']}/cancel").json()["status"] == "cancelled"
assert [i["id"] for i in hub.get("/outbox").json()["items"]] == [b["id"]]
# b: collected, then cancelled → the hub is told once, then it's off the list
hub.post(f"/outbox/{b['id']}/ack", json={"hub": {"ok": True}})
run.post(f"/outbox/{b['id']}/cancel")
items = hub.get("/outbox").json()["items"]
assert items and items[0]["id"] == b["id"] and items[0]["cancelRequested"] is True
hub.post(f"/outbox/{b['id']}/ack", json={"hub": {"cancelled": True}})
assert hub.get("/outbox").json()["count"] == 0
def test_persists_across_a_restart(tmp_path):
b1 = outbox.Outbox(tmp_path / "o.json")
b1.put("notify", {"id": "ntf_0123456789ab", "title": "t", "message": "m"})
b2 = outbox.Outbox(tmp_path / "o.json")
assert b2.get("ntf_0123456789ab")["status"] == "queued"
assert [i["id"] for i in b2.pending()] == ["ntf_0123456789ab"]