_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>
99 lines
3.7 KiB
Python
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)
|