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>
107 lines
4.2 KiB
Python
107 lines
4.2 KiB
Python
"""The remote transcript mirror against a real (loopback) worker feed:
|
|
first sync, append, shrink, subagents, stamps, offline back-off."""
|
|
|
|
import os
|
|
import pathlib
|
|
import tempfile
|
|
import unittest
|
|
|
|
from workers import remotefeed
|
|
|
|
import workers
|
|
|
|
from _fakeworker import FakeWorker
|
|
|
|
|
|
class Mirror(unittest.TestCase):
|
|
def setUp(self):
|
|
self.tmp = pathlib.Path(tempfile.mkdtemp())
|
|
self.remote = self.tmp / "projects"
|
|
self.fw = FakeWorker(self.remote, self.tmp / "state")
|
|
self.store = workers.WorkerStore(str(self.tmp / "workers.json"))
|
|
self.w = workers.pair(self.store, self.fw.pairing_string("work"),
|
|
account=None, name=None, hub_name="t",
|
|
resolve_login=lambda _l: None,
|
|
valid_account=lambda a: a)
|
|
self.dest = self.tmp / "mirror"
|
|
self.landed: list[tuple[str, list[str]]] = []
|
|
self.m = remotefeed.FeedMirror(
|
|
self.store, self.dest,
|
|
on_files=lambda w, rels: self.landed.append((w["id"], rels)))
|
|
|
|
def tearDown(self):
|
|
self.fw.stop()
|
|
|
|
def _remote(self, rel, text, mode="w"):
|
|
p = self.remote / rel
|
|
p.parent.mkdir(parents=True, exist_ok=True)
|
|
with open(p, mode) as f:
|
|
f.write(text)
|
|
return p
|
|
|
|
def test_first_sync_mirrors_with_remote_mtime(self):
|
|
p = self._remote("-Users-g-repo/s1.jsonl", '{"a":1}\n')
|
|
os.utime(p, (1_700_000_000, 1_700_000_000))
|
|
self._remote("-Users-g-repo/s1/subagents/agent-x.jsonl", "{}\n")
|
|
self._remote("-Users-g-repo/s1/subagents/agent-x.meta.json", "{}")
|
|
self.assertEqual(self.m.tick(), 3)
|
|
local = self.dest / "-Users-g-repo/s1.jsonl"
|
|
self.assertEqual(local.read_text(), '{"a":1}\n')
|
|
self.assertAlmostEqual(local.stat().st_mtime, 1_700_000_000, places=2)
|
|
self.assertTrue((self.dest / "-Users-g-repo/s1/subagents/"
|
|
"agent-x.meta.json").is_file())
|
|
wid, rels = self.landed[0]
|
|
self.assertEqual(wid, self.w["id"])
|
|
self.assertEqual(sorted(rels), ["-Users-g-repo/s1.jsonl",
|
|
"-Users-g-repo/s1/subagents/agent-x.jsonl"])
|
|
self.assertEqual(self.m.tick(), 0) # nothing changed
|
|
|
|
def test_append_fetches_only_the_tail(self):
|
|
p = self._remote("p/s.jsonl", '{"a":1}\n')
|
|
self.m.tick()
|
|
self._remote("p/s.jsonl", '{"b":2}\n', mode="a")
|
|
self.assertEqual(self.m.tick(), 1)
|
|
self.assertEqual((self.dest / "p/s.jsonl").read_text(),
|
|
'{"a":1}\n{"b":2}\n')
|
|
self.assertEqual(p.stat().st_size,
|
|
(self.dest / "p/s.jsonl").stat().st_size)
|
|
|
|
def test_shrink_redownloads_whole(self):
|
|
self._remote("p/s.jsonl", "x" * 50 + "\n")
|
|
self.m.tick()
|
|
self._remote("p/s.jsonl", "short\n")
|
|
self.m.tick()
|
|
self.assertEqual((self.dest / "p/s.jsonl").read_text(), "short\n")
|
|
self.assertFalse((self.dest / "p/s.jsonl.part").exists())
|
|
|
|
def test_large_file_spans_chunks(self):
|
|
old = remotefeed.CHUNK
|
|
remotefeed.CHUNK = 7
|
|
try:
|
|
self._remote("p/big.jsonl", "".join(f"{i}\n" for i in range(100)))
|
|
self.m.tick()
|
|
finally:
|
|
remotefeed.CHUNK = old
|
|
self.assertEqual((self.dest / "p/big.jsonl").read_text(),
|
|
"".join(f"{i}\n" for i in range(100)))
|
|
|
|
def test_offline_worker_is_skipped_then_backed_off(self):
|
|
self._remote("p/s.jsonl", "{}\n")
|
|
self.store.set_status(self.w["id"], status="offline")
|
|
self.assertEqual(self.m.tick(), 0)
|
|
self.store.set_status(self.w["id"], status="online")
|
|
self.fw.stop() # online per monitor, but down
|
|
self.assertEqual(self.m.tick(), 0)
|
|
self.assertIn(self.w["id"], self.m._backoff)
|
|
self.assertIn("feed:", self.store.status(self.w["id"])["lastError"])
|
|
|
|
def test_safe_rel(self):
|
|
for bad in ("../x.jsonl", "/abs.jsonl", "x.jsonl", "p/../x.jsonl",
|
|
"p/x.txt", "p//x.jsonl"):
|
|
self.assertIsNone(remotefeed.safe_rel(bad), bad)
|
|
self.assertIsNotNone(remotefeed.safe_rel("p/s/subagents/a.meta.json"))
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|