A worker session resumed right after a hub-launched run exited looked 'active in a terminal' (its mirror copy was seconds old). Growth up to the exit watcher's exitedAt stamp, or while the hub's run still reads running, is ours; only writes after it are a terminal's. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
143 lines
6.7 KiB
Python
143 lines
6.7 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 main # noqa: E402
|
|
import workers # noqa: E402
|
|
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(main.worker_store, self.fw.pairing_string("work"),
|
|
account=None, name="mac", hub_name="t",
|
|
resolve_login=lambda _l: None,
|
|
valid_account=main.accounts_mod.valid)
|
|
|
|
def tearDown(self):
|
|
self.fw.stop()
|
|
main.worker_store.delete(self.w["id"])
|
|
|
|
def test_new_run_on_the_account_goes_to_its_worker(self):
|
|
out = main._send_message("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 = main.meta_store.get(out["sessionId"])
|
|
self.assertEqual((m["worker"], m["account"]), (self.w["id"], "work"))
|
|
|
|
def test_local_pick_and_offline_fallback(self):
|
|
main.worker_store.set_status(self.w["id"], status="offline")
|
|
# auto: offline worker ⇒ the host sidecar (unconfigured here ⇒ 503)
|
|
with self.assertRaises(HTTPException) as cm:
|
|
main._send_message("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:
|
|
main._send_message("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 = main._send_message("hi", account="personal", worker=self.w["id"])
|
|
self.assertEqual(main.meta_store.get(out["sessionId"])["account"],
|
|
"work")
|
|
|
|
def test_pi_never_goes_to_a_worker(self):
|
|
with self.assertRaises(HTTPException) as cm:
|
|
main._send_message("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):
|
|
main.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:
|
|
main._send_message("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
|
|
main._send_message("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(main.meta_store.get(SID)["worker"], self.w["id"])
|
|
# an old transcript isn't guarded
|
|
os.utime(mirror, (time.time() - 3600, time.time() - 3600))
|
|
main._send_message("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}
|
|
main._send_message("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))
|
|
main.meta_store.update(sid, {"worker": self.w["id"], "state": "finished",
|
|
"exitedAt": stamp})
|
|
main._send_message("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
|
|
main.meta_store.update(sid, {"state": "finished"})
|
|
os.utime(mirror, (now, now))
|
|
with self.assertRaises(HTTPException) as cm:
|
|
main._send_message("next", session_id=sid, cwd="/Users/g/repo")
|
|
self.assertEqual(cm.exception.detail["kind"], "terminal_live")
|
|
|
|
def test_interrupt_routes_to_the_worker(self):
|
|
main.meta_store.update(SID, {"worker": self.w["id"]})
|
|
main.interrupt_conversation(main.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"
|
|
main._record_worker_files(self.w, [f"-Users-g-repo/{sid}.jsonl",
|
|
f"-Users-g-repo/{sid}/subagents/a.jsonl"])
|
|
m = main.meta_store.get(sid)
|
|
self.assertEqual((m["workerSource"], m["accountSource"]),
|
|
(self.w["id"], "work"))
|
|
self.assertEqual(main.accounts_mod.resolve({}, m), "work")
|
|
|
|
def test_live_sids_per_runner(self):
|
|
self.fw.live = {"a", "b"}
|
|
self.assertEqual(main._sidecar_live_sids(self.w["id"]), {"a", "b"})
|
|
main.worker_store.set_status(self.w["id"], status="offline")
|
|
self.assertIsNone(main._sidecar_live_sids(self.w["id"]))
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|