Files
ai-agent/backend/test_worker_routing.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

219 lines
11 KiB
Python

"""main's run routing across runners: a new run picks the account's online
worker (or the composer's pick), a continue follows the session to its worker,
the terminal-live guard, interrupt routing and per-runner exit baselines.
Imports ``main`` against throwaway paths (no startup event → no threads)."""
import os
import pathlib
import tempfile
import time
import unittest
TMP = pathlib.Path(tempfile.mkdtemp())
for k, v in {"DB_PATH": TMP / "x.db", "TRANSCRIPTS_DIR": TMP / "t",
"WORKSPACE": TMP, "REPO_DIR": TMP, "ARCHIVE_DIR": TMP / "arch",
"META_PATH": TMP / "meta.json", "WORKERS_PATH": TMP / "w.json",
"WORKERS_TRANSCRIPTS_DIR": TMP / "mirror",
"CRON_PATH": TMP / "cron.json", "SIDECAR_TOKEN": "",
"PI_TRANSCRIPTS_DIR": ""}.items():
os.environ[k] = str(v)
from fastapi import HTTPException # noqa: E402
import accounts # noqa: E402
import main # noqa: E402
import workers # noqa: E402
from core.state import _record_worker_files # noqa: E402
from cron.routes import CronCreateBody, CronPatchBody, cron_create, cron_delete, cron_update # noqa: E402
from cron.service import _fire_cron_job # noqa: E402
from runs.routes import InterruptBody, interrupt_conversation # noqa: E402
from runs.service import _send_message # noqa: E402
from runs.sidecar_client import _sidecar_live_sids # noqa: E402
DEPS = main.app.state.ai
from _fakeworker import FakeWorker # noqa: E402
SID = "11111111-2222-3333-4444-555555555555"
class Routing(unittest.TestCase):
def setUp(self):
self.fw = FakeWorker(TMP / "remote", TMP / "state")
self.w = workers.pair(DEPS.worker_store, self.fw.pairing_string("work"),
account=None, name="mac", hub_name="t",
resolve_login=lambda _l: None,
valid_account=accounts.valid,
auto_route=True)
def tearDown(self):
self.fw.stop()
DEPS.worker_store.delete(self.w["id"])
def test_new_run_on_the_account_goes_to_its_worker(self):
out = _send_message(DEPS, "hi", projects=["orus-monorepo"])
path, body = self.fw.calls[-1]
self.assertEqual(path, "/spawn")
self.assertNotIn("account", body) # the worker's own login
self.assertIsNone(body["cwd"]) # the worker's own default
m = DEPS.meta_store.get(out["sessionId"])
self.assertEqual((m["worker"], m["account"]), (self.w["id"], "work"))
def test_without_auto_route_the_lab_keeps_the_account(self):
DEPS.worker_store.patch(self.w["id"], {"autoRoute": False})
# auto: the worker is online but opt-out ⇒ the host sidecar
# (unconfigured here ⇒ 503), never the worker
with self.assertRaises(HTTPException) as cm:
_send_message(DEPS, "hi", projects=["orus-monorepo"])
self.assertEqual(cm.exception.status_code, 503)
self.assertFalse([c for c in self.fw.calls if c[0] == "/spawn"])
# an explicit pick still lands there
out = _send_message(DEPS, "hi", worker=self.w["id"])
self.assertEqual(self.fw.calls[-1][0], "/spawn")
self.assertEqual(DEPS.meta_store.get(out["sessionId"])["worker"],
self.w["id"])
def test_local_pick_and_offline_fallback(self):
DEPS.worker_store.set_status(self.w["id"], status="offline")
# auto: offline worker ⇒ the host sidecar (unconfigured here ⇒ 503)
with self.assertRaises(HTTPException) as cm:
_send_message(DEPS, "hi", projects=["orus-monorepo"])
self.assertEqual(cm.exception.status_code, 503)
# an explicit pick of an offline worker is refused, not rerouted
with self.assertRaises(HTTPException) as cm:
_send_message(DEPS, "hi", worker=self.w["id"])
self.assertEqual(cm.exception.detail["kind"], "worker_offline")
self.assertEqual(self.fw.calls, [])
def test_explicit_worker_bills_its_account(self):
out = _send_message(DEPS, "hi", account="personal", worker=self.w["id"])
self.assertEqual(DEPS.meta_store.get(out["sessionId"])["account"],
"work")
def test_pi_never_goes_to_a_worker(self):
with self.assertRaises(HTTPException) as cm:
_send_message(DEPS, "hi", harness="pi", worker=self.w["id"])
self.assertEqual(cm.exception.detail["kind"], "worker_harness")
def test_continue_follows_the_session_and_guards_a_live_terminal(self):
DEPS.meta_store.update(SID, {"workerSource": self.w["id"],
"accountSource": "work"})
mirror = TMP / "mirror" / "-Users-g-repo" / f"{SID}.jsonl"
mirror.parent.mkdir(parents=True, exist_ok=True)
mirror.write_text("{}\n") # grew just now
with self.assertRaises(HTTPException) as cm:
_send_message(DEPS, "more", session_id=SID, cwd="/Users/g/repo")
self.assertEqual(cm.exception.detail["kind"], "terminal_live")
self.assertEqual(self.fw.calls, [])
# confirmed takeover goes through, to the worker, force-resume
_send_message(DEPS, "more", session_id=SID, cwd="/Users/g/repo",
takeover=True)
path, body = self.fw.calls[-1]
self.assertEqual((path, body["force"], body["cwd"]),
("/resume", True, "/Users/g/repo"))
self.assertEqual(DEPS.meta_store.get(SID)["worker"], self.w["id"])
# an old transcript isn't guarded
os.utime(mirror, (time.time() - 3600, time.time() - 3600))
_send_message(DEPS, "again", session_id=SID, cwd="/Users/g/repo")
# …nor is our own live run on the worker
mirror.write_text("{}\n{}\n")
self.fw.live = {SID}
_send_message(DEPS, "again", session_id=SID, cwd="/Users/g/repo")
self.assertEqual([c[0] for c in self.fw.calls],
["/resume", "/resume", "/resume"])
def test_growth_from_our_own_finished_run_is_not_a_terminal(self):
sid = "22222222-2222-3333-4444-555555555555"
mirror = TMP / "mirror" / "-Users-g-repo" / f"{sid}.jsonl"
mirror.parent.mkdir(parents=True, exist_ok=True)
mirror.write_text("{}\n")
now = time.time()
os.utime(mirror, (now - 10, now - 10))
stamp = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(now - 8))
DEPS.meta_store.update(sid, {"worker": self.w["id"], "state": "finished",
"exitedAt": stamp})
_send_message(DEPS, "next", session_id=sid, cwd="/Users/g/repo")
self.assertEqual(self.fw.calls[-1][0], "/resume")
# …but a write after our run exited is a terminal's
DEPS.meta_store.update(sid, {"state": "finished"})
os.utime(mirror, (now, now))
with self.assertRaises(HTTPException) as cm:
_send_message(DEPS, "next", session_id=sid, cwd="/Users/g/repo")
self.assertEqual(cm.exception.detail["kind"], "terminal_live")
def test_continue_into_a_live_run_goes_through_its_inbox(self):
DEPS.meta_store.update(SID, {"worker": self.w["id"], "state": "running"})
self.fw.live = {SID}
self.fw.inbox = {SID}
out = _send_message(DEPS, "Up?", session_id=SID, cwd="/Users/g/repo")
self.assertEqual(out["delivered"], "inbox")
self.assertEqual([c[0] for c in self.fw.calls], ["/message"])
self.assertEqual(self.fw.calls[-1][1]["prompt"], "Up?")
self.assertEqual(DEPS.meta_store.get(SID)["state"], "running")
# the live run can't take it (409) ⇒ the force-resume as before
self.fw.inbox = set()
_send_message(DEPS, "Up?", session_id=SID, cwd="/Users/g/repo")
path, body = self.fw.calls[-1]
self.assertEqual((path, body["force"]), ("/resume", True))
# nothing live (404) ⇒ a plain resume, still force (harmless)
self.fw.live = set()
_send_message(DEPS, "Up?", session_id=SID, cwd="/Users/g/repo")
self.assertEqual(self.fw.calls[-1][0], "/resume")
def test_interrupt_routes_to_the_worker(self):
DEPS.meta_store.update(SID, {"worker": self.w["id"]})
interrupt_conversation(DEPS, InterruptBody(sessionId=SID))
self.assertEqual(self.fw.calls[-1][0], "/interrupt")
def test_mirror_stamps_worker_and_account(self):
sid = "99999999-2222-3333-4444-555555555555"
_record_worker_files(DEPS, self.w, [f"-Users-g-repo/{sid}.jsonl",
f"-Users-g-repo/{sid}/subagents/a.jsonl"])
m = DEPS.meta_store.get(sid)
self.assertEqual((m["workerSource"], m["accountSource"]),
(self.w["id"], "work"))
self.assertEqual(accounts.resolve({}, m), "work")
def test_cron_job_runs_on_its_worker_and_never_auto_routes(self):
agents = TMP / ".claude" / "agents"
agents.mkdir(parents=True, exist_ok=True)
# account: work would auto-route to the Mac — an unpinned job stays home
(agents / "pr.md").write_text("---\naccount: work\n---\nwatch PRs\n")
job = cron_create(DEPS, CronCreateBody(
name="pr", schedule="0 9 * * *",
promptFile=".claude/agents/pr.md"))
self.assertIsNone(job["worker"])
with self.assertRaises(HTTPException) as cm: # lab sidecar: 503 here
_fire_cron_job(DEPS, DEPS.cron_store.get("pr"))
self.assertEqual(cm.exception.status_code, 503)
self.assertEqual(self.fw.calls, [])
# pinned to the worker: the prompt text goes there
job = cron_update(DEPS, "pr", CronPatchBody(worker=self.w["id"]))
self.assertEqual(job["worker"], self.w["id"])
out = _fire_cron_job(DEPS, DEPS.cron_store.get("pr"))
path, body = self.fw.calls[-1]
self.assertEqual((path, body["prompt"]), ("/spawn", "watch PRs"))
m = DEPS.meta_store.get(out["sessionId"])
self.assertEqual((m["worker"], m["cron"]["id"]), (self.w["id"], "pr"))
# offline ⇒ refused, not rerouted to the lab
DEPS.worker_store.set_status(self.w["id"], status="offline")
with self.assertRaises(HTTPException) as cm:
_fire_cron_job(DEPS, DEPS.cron_store.get("pr"))
self.assertEqual(cm.exception.detail["kind"], "worker_offline")
# unknown ids are rejected up front; "" moves it back to the lab
with self.assertRaises(HTTPException):
cron_update(DEPS, "pr", CronPatchBody(worker="nope"))
job = cron_update(DEPS, "pr", CronPatchBody(worker=""))
self.assertIsNone(job["worker"])
cron_delete(DEPS, "pr")
def test_live_sids_per_runner(self):
self.fw.live = {"a", "b"}
self.assertEqual(_sidecar_live_sids(DEPS, self.w["id"]), {"a", "b"})
DEPS.worker_store.set_status(self.w["id"], status="offline")
self.assertIsNone(_sidecar_live_sids(DEPS, self.w["id"]))
if __name__ == "__main__":
unittest.main()