Files
ai-agent/backend/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

198 lines
8.2 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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