Files
ai-agent/backend/_fakeworker.py
Gabriel Vidal 74393fa474 feat(hub): deliver a follow-up into a live run through its inbox
_send_message now POSTs the prompt to the sidecar's /message first: a
run that is still alive keeps its background work (a build, a subagent,
a Monitor on a form) and reads the message between tool calls, or at
once when idle. Only 404 (nothing live) or 409 (the run can't take it —
a pi run, an older CLI) falls back to the old stop-then-resume.

Claude Code records such a message as a peer's ("Another Claude session
sent a message … not typed by your user"); the sidecar prefixes the text
to say it is the user's, and conversations._unwrap_inbox shows the plain
user turn. SpawnResult gains `delivered: "inbox"`.

Docs: the headless background-task semantics in CLAUDE.md and
PROJECT_CLAUDE.md; FakeWorker learns /message for the routing tests.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-06 23:24:50 +02:00

99 lines
3.7 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()
# sessions whose live run takes an inbox message (/message delivers);
# a live one not listed here answers 409 inbox_unavailable, a session
# that isn't live at all 404 — the two fallbacks main.py handles.
self.inbox: 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.post("/message")
def message(body: dict,
authorization: str | None = Header(default=None)):
auth(authorization)
sid = body.get("sessionId")
if sid not in self.live:
raise HTTPException(404, {"message": "no running process",
"kind": "not_running"})
if sid not in self.inbox:
raise HTTPException(409, {"message": "no inbox socket",
"kind": "inbox_unavailable"})
self.calls.append(("/message", body))
return {"sessionId": sid, "pid": 4242, "delivered": "inbox"}
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)