- workers.py: WorkerStore (/data/workers.json, 0600), pairing-string parse, pair/unpair, health+token probe, WorkerMonitor (workers SSE) - remotefeed.py: FeedMirror pulls each online worker's /feed into WORKERS_TRANSCRIPTS_DIR (a new SOURCE_DIR), append/shrink/backoff, stamps meta.workerSource + accountSource in one save per pass - main: _sidecar_call(worker=), new runs pick chip → account's online worker → local; continue/fork/interrupt follow meta.worker/workerSource; terminal_live 409 unless takeover; exit watcher keeps one baseline per runner; /api/workers CRUD + probe; ConvMeta.worker - tests against a real loopback worker (feed + pairing routers) Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
105 lines
4.2 KiB
Python
105 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
|
|
|
|
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()
|