Files
ai-agent/backend/workers/remotefeed.py
Gabriel Vidal 0ff9e40242 refactor(backend): split main.py into domain packages with app.state injection
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>
2026-10-06 23:55:47 +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