- 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>
198 lines
8.2 KiB
Python
198 lines
8.2 KiB
Python
"""
|
||
The remote transcript mirror — a worker's transcripts, pulled into a local dir.
|
||
|
||
A worker (``workers.py``) writes its transcripts on *its* disk; the hub has no
|
||
mount for them. This thread pulls them over the worker's ``/feed`` into
|
||
``WORKERS_TRANSCRIPTS_DIR`` using the same ``<project-slug>/<sid>.jsonl``
|
||
layout Claude Code writes, and that dir is simply one more of the backend's
|
||
``SOURCE_DIRS`` — so the Watcher, the indexer, the resumable parser, costs and
|
||
SSE all treat a Mac session exactly like a local one. Nothing downstream knows
|
||
the bytes crossed the tailnet.
|
||
|
||
Per tick, per online worker: ``GET /feed?since=<last now − slack>`` lists what
|
||
changed; each file is compared with the mirror copy — a bigger remote file has
|
||
its tail appended from the local size, a smaller one (a rewritten file) is
|
||
re-downloaded whole, an equal one is skipped. The mirror copy then takes the
|
||
remote mtime, so a months-old terminal session backfilled on first pairing
|
||
doesn't look like it ran a second ago (the "running" recency window reads it).
|
||
|
||
Every pass reports the transcripts that landed to ``on_files(worker, rels)``
|
||
— main stamps ``meta.workerSource`` (which runner a terminal-started session lives
|
||
on, so resume/fork/interrupt route there) and ``meta.accountSource`` (the
|
||
worker's account). A worker that stops answering is backed off (2 s → 60 s)
|
||
until the monitor sees it again.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import os
|
||
import pathlib
|
||
import threading
|
||
import time
|
||
import urllib.error
|
||
import urllib.parse
|
||
import urllib.request
|
||
|
||
SLACK_S = 10.0 # re-list files touched this close to the last cursor
|
||
CHUNK = 8 * 1024 * 1024
|
||
_SUFFIXES = (".jsonl", ".meta.json")
|
||
|
||
|
||
def safe_rel(rel: str) -> pathlib.PurePosixPath | None:
|
||
"""A feed path we're willing to write under the mirror dir, else None —
|
||
the worker is trusted with runs, not with our filesystem layout."""
|
||
if not rel or rel.startswith("/") or "\\" in rel or "\0" in rel \
|
||
or not rel.endswith(_SUFFIXES):
|
||
return None
|
||
parts = rel.split("/")
|
||
if len(parts) < 2 or any(p in ("", ".", "..") for p in parts):
|
||
return None
|
||
return pathlib.PurePosixPath(rel)
|
||
|
||
|
||
class WorkerDown(Exception):
|
||
pass
|
||
|
||
|
||
class FeedMirror(threading.Thread):
|
||
def __init__(self, store, dest: pathlib.Path, on_files=None,
|
||
interval: float = 2.0):
|
||
super().__init__(daemon=True, name="feed-mirror")
|
||
self.store = store
|
||
self.dest = pathlib.Path(dest)
|
||
self.on_files = on_files
|
||
self.interval = interval
|
||
self._cursor: dict[str, float] = {}
|
||
self._backoff: dict[str, tuple[float, float]] = {} # wid → (until, step)
|
||
self._stop = threading.Event()
|
||
|
||
def stop(self) -> None:
|
||
self._stop.set()
|
||
|
||
# ── HTTP ──────────────────────────────────────────────────────────────
|
||
def _get(self, w: dict, path: str, params: dict, timeout: float = 20):
|
||
url = f"{w['url']}{path}?{urllib.parse.urlencode(params)}"
|
||
req = urllib.request.Request(
|
||
url, headers={"Authorization": f"Bearer {w['token']}"})
|
||
try:
|
||
r = urllib.request.urlopen(req, timeout=timeout)
|
||
except urllib.error.HTTPError as e:
|
||
if e.code == 404:
|
||
raise FileNotFoundError(path)
|
||
raise WorkerDown(f"HTTP {e.code}")
|
||
except (urllib.error.URLError, OSError) as e:
|
||
raise WorkerDown(str(e))
|
||
return r
|
||
|
||
def _chunk(self, w: dict, rel: str, offset: int) -> tuple[bytes, int, float]:
|
||
with self._get(w, "/feed/file", {"rel": rel, "offset": offset,
|
||
"limit": CHUNK}, timeout=60) as r:
|
||
data = r.read()
|
||
size = int(r.headers.get("X-Feed-Size") or 0)
|
||
mtime = float(r.headers.get("X-Feed-Mtime") or 0)
|
||
return data, size, mtime
|
||
|
||
# ── one worker ────────────────────────────────────────────────────────
|
||
def sync(self, w: dict) -> int:
|
||
"""One pass over a worker's feed; returns how many files changed."""
|
||
since = max(0.0, self._cursor.get(w["id"], 0.0) - SLACK_S)
|
||
with self._get(w, "/feed", {"since": since}) as r:
|
||
listing = json.loads(r.read().decode())
|
||
files = sorted(listing.get("files") or [],
|
||
key=lambda f: f.get("mtime") or 0)
|
||
n = 0
|
||
landed: list[str] = []
|
||
try:
|
||
for f in files:
|
||
rel = safe_rel(f.get("rel") or "")
|
||
if rel is None:
|
||
continue
|
||
try:
|
||
if self._sync_file(w, rel, int(f.get("size") or 0),
|
||
float(f.get("mtime") or 0)):
|
||
n += 1
|
||
if rel.suffix == ".jsonl":
|
||
landed.append(rel.as_posix())
|
||
except FileNotFoundError:
|
||
continue # vanished between listing and fetch
|
||
finally:
|
||
# Even when the worker drops mid-pass: what did land is on disk
|
||
# now and won't be re-reported (sizes match next time).
|
||
if landed and self.on_files:
|
||
self.on_files(w, landed)
|
||
self._cursor[w["id"]] = float(listing.get("now") or time.time())
|
||
return n
|
||
|
||
def _sync_file(self, w: dict, rel: pathlib.PurePosixPath, size: int,
|
||
mtime: float) -> bool:
|
||
local = self.dest / rel
|
||
try:
|
||
st = local.stat()
|
||
lsize, lmtime = st.st_size, st.st_mtime
|
||
except OSError:
|
||
lsize, lmtime = -1, 0.0
|
||
if lsize == size and abs(lmtime - mtime) < 1e-3:
|
||
return False
|
||
local.parent.mkdir(parents=True, exist_ok=True)
|
||
if 0 <= lsize <= size:
|
||
# grew (or only its mtime moved): append the tail
|
||
offset, mode = max(lsize, 0), "ab"
|
||
target = local
|
||
else:
|
||
# new, or shrank — the file was rewritten: fetch it whole
|
||
offset, mode = 0, "wb"
|
||
target = local.with_name(local.name + ".part")
|
||
with open(target, mode) as out:
|
||
while True:
|
||
data, rsize, rmtime = self._chunk(w, rel.as_posix(), offset)
|
||
if data:
|
||
out.write(data)
|
||
offset += len(data)
|
||
mtime = rmtime or mtime
|
||
if not data or offset >= rsize:
|
||
break
|
||
if target is not local:
|
||
os.replace(target, local)
|
||
if mtime:
|
||
os.utime(local, (mtime, mtime))
|
||
return True
|
||
|
||
# ── loop ──────────────────────────────────────────────────────────────
|
||
def tick(self) -> int:
|
||
n = 0
|
||
now = time.monotonic()
|
||
for w in self.store.list():
|
||
wid = w["id"]
|
||
if not self.store.online(wid):
|
||
continue
|
||
until, step = self._backoff.get(wid, (0.0, 0.0))
|
||
if now < until:
|
||
continue
|
||
try:
|
||
self.store.set_status(wid, syncing=True)
|
||
got = self.sync(w)
|
||
n += got
|
||
self._backoff.pop(wid, None)
|
||
if got:
|
||
st = self.store.status(wid)
|
||
self.store.set_status(
|
||
wid, mirrored=(st.get("mirrored") or 0) + got)
|
||
except WorkerDown as e:
|
||
step = min(60.0, max(2.0, step * 2))
|
||
self._backoff[wid] = (now + step, step)
|
||
self.store.set_status(wid, lastError=f"feed: {e}"[:200])
|
||
except Exception as e: # never let one bad file kill the thread
|
||
self.store.set_status(wid, lastError=f"feed: {e}"[:200])
|
||
finally:
|
||
self.store.set_status(wid, syncing=False)
|
||
return n
|
||
|
||
def run(self) -> None:
|
||
self.dest.mkdir(parents=True, exist_ok=True)
|
||
while not self._stop.wait(self.interval):
|
||
try:
|
||
self.tick()
|
||
except Exception:
|
||
pass
|