- 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>
81 lines
2.8 KiB
Python
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)
|