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>
219 lines
11 KiB
Python
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()
|