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>
135 lines
5.0 KiB
Python
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"]
|