Files
ai-agent/backend/test_outbox_mirror.py
Gabriel Vidal 0ff9e40242 refactor(backend): split main.py into domain packages with app.state injection
main.py (5146 lines, 106 routes) becomes an assembly only: one package per
domain — core, files, settings, dashboards, conversations, diff,
notifications, forms, runs, accounts, workers, models, cron, agents,
projects, services, goals, memories, plans, templates — each exposing an
APIRouter; the flat domain modules move into their package behind a barrel
that keeps the old `import conversations` / `import projects` spellings.

The shared singletons (store, meta_store, hub, indexer, …) are built once by
core.state.build_state() and attached to app.state.ai; routes take them as
the `deps: State` dependency and helpers as an explicit `deps: AppState`.
conversations/pricing.py carries the per-model rates out of the parser.

Verified: route table and OpenAPI byte-identical; 90 read endpoints
golden-diffed against the monolith on a copy of the live data (identical);
write routes smoke-tested; 66 backend tests pass.

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

178 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
from workers 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()