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

81 lines
2.8 KiB
Python

"""A real HTTP worker for the worker/mirror tests: the sidecar's own feed +
pairing routers (``../sidecar``) plus a stub ``/sessions`` + ``/health``,
served by uvicorn on a free loopback port in a thread."""
import os
import pathlib
import socket
import sys
import threading
import time
import uvicorn
from fastapi import FastAPI, Header, HTTPException
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parent.parent
/ "sidecar"))
import feed # noqa: E402
import pairing # noqa: E402
class FakeWorker:
def __init__(self, root: pathlib.Path, state: pathlib.Path):
os.environ["WORKER_STATE_DIR"] = str(state)
self.root = root
self.live: set[str] = set()
self.calls: list[tuple[str, dict]] = []
app = FastAPI()
def auth(h):
tok = pairing.token()
if not tok or h != f"Bearer {tok}":
raise HTTPException(401, "unauthorized")
@app.get("/health")
def health():
return {"ok": True, **pairing.identity(), "cwd": "/Users/g/repo",
"permissionMode": "bypassPermissions"}
@app.get("/sessions")
def sessions(authorization: str | None = Header(default=None)):
auth(authorization)
return {"sessions": [{"sessionId": s} for s in self.live],
"count": len(self.live)}
def _record(path):
def handler(body: dict,
authorization: str | None = Header(default=None)):
auth(authorization)
self.calls.append((path, body))
return {"sessionId": body.get("sessionId"), "pid": 4242}
app.post(path)(handler)
for path in ("/spawn", "/resume", "/fork", "/interrupt"):
_record(path)
app.include_router(feed.make_router(auth, lambda: self.root))
app.include_router(pairing.make_router(
auth, lambda: {"email": "g@orus.example", "orgName": "Orus",
"loggedIn": True},
lambda: {"cwd": "/Users/g/repo"}))
s = socket.socket()
s.bind(("127.0.0.1", 0))
self.port = s.getsockname()[1]
s.close()
self.server = uvicorn.Server(uvicorn.Config(
app, host="127.0.0.1", port=self.port, log_level="warning"))
self.thread = threading.Thread(target=self.server.run, daemon=True)
self.thread.start()
for _ in range(100):
if self.server.started:
break
time.sleep(0.05)
def pairing_string(self, login: str | None = "g@orus.example") -> str:
code = pairing.new_code()
return f"127.0.0.1:{self.port}/{code}" + (f"/{login}" if login else "")
def stop(self):
self.server.should_exit = True
self.thread.join(timeout=5)