Files
ai-agent/backend/test_remotefeed.py
Gabriel Vidal 18f4bae5f1 feat(workers): pair remote workers, mirror their transcripts, route runs to them
- 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>
2026-09-28 13:57:24 +02:00

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()