diff --git a/CHANGELOG.md b/CHANGELOG.md index 2b7e86a..175d29f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,11 @@ ## Unreleased +- `lich serve` queues frames per session on each connection instead of + per connection (#144). A long `prompt.submit` no longer holds up other + sessions or session-less frames such as `health` on the same socket. Frames + for one session keep their order; replies across sessions may interleave, + so clients match them by JSON-RPC `id` (Ossuary already does). - Gateway (#144): sender allowlists and the gateway toolset are one `GatewayPolicy`, built from config at startup (`src/gateway/access.ts`). Platform adapters declare `capabilities` (`{ kind: "text", max_reply_chars }`); diff --git a/docs/architecture/serve.md b/docs/architecture/serve.md index a1ca313..78056d9 100644 --- a/docs/architecture/serve.md +++ b/docs/architecture/serve.md @@ -84,8 +84,9 @@ about it from the next `session.clear` or `prompt.submit` failing with `not_found`. Evicted transcripts stay on disk and can be resumed again. `prompt.submit` runs the server Agent with that bag's history and -`AgentRunOptions.session` (one JSONL file per serve session). Runs are serialized -so AgentEvent fan-out stays correctly tagged with `session_id`. While a run is +`AgentRunOptions.session` (one JSONL file per serve session). Runs of one +session are serialized; different sessions run concurrently. Each run's events +arrive through its own `on_event` and are tagged with that run's `session_id`. While a run is in flight, `prompt.abort` aborts it via `AbortSignal`. The submit result mirrors `AgentRunResult` (`reply`, `usage`, `session_path`, `turns_used`, `stopped_reason`). @@ -113,8 +114,12 @@ same WebSocket that issued `prompt.submit` while the call is still in flight. do not assume `file://` / `app://` behavior here. - Frame size capped at ~1 MiB (`maxPayload`). - `health` returns `{ status: "ok", version }` (`LICH_VERSION` from `src/version.ts`). -- Frames are handled per-connection in arrival order: pipelined requests get - in-order replies even when earlier requests hit slower filesystem awaits. +- Frames are queued per session within a connection: frames whose + `params.session_id` match are handled in arrival order and get in-order + replies, while other sessions, and frames that name no session (such as + `health`, `session.list` or `session.create`), do not wait behind them. + Replies from different sessions can therefore interleave; match them by + JSON-RPC `id`. `prompt.abort` is dispatched immediately so it can cancel an in-flight `prompt.submit` on the same socket instead of waiting behind that run. A submit already accepted on that socket, but still waiting behind another diff --git a/src/serve/server.ts b/src/serve/server.ts index 082117d..a8f5ebd 100644 --- a/src/serve/server.ts +++ b/src/serve/server.ts @@ -233,11 +233,13 @@ function attach_client( prompts: ServePromptService | undefined, inflight: Set>, ): void { - // Serialize frames per connection: pipelined requests get in-order replies, - // and the tail catch keeps any rejection from becoming an unhandled one. - // prompt.abort is not on that tail. It must run while prompt.submit is still + // Serialize frames per session on this connection: frames for one session_id + // get in-order replies, while different sessions (and session-less frames + // such as health or session.list) no longer wait behind each other's runs. + // Replies across lanes can interleave; clients match them by JSON-RPC id. + // prompt.abort is on no lane. It must run while prompt.submit is still // awaiting the model; waiting would deadlock a hung run with its own cancel. - let tail: Promise = Promise.resolve(); + const tails = new Map>(); const track = (chain: Promise): void => { inflight.add(chain); void chain.finally(() => { @@ -256,12 +258,29 @@ function attach_client( // Arm the controller before this frame waits on `tail`, so prompt.abort // (which skips the tail) can cancel a submit that has not started yet. const prepared = arm_queued_submit(data, prompts); - const chain = tail.then(() => run_frame(data, prepared)); - tail = chain; + const lane = frame_lane(data); + const chain = (tails.get(lane) ?? Promise.resolve()).then(() => run_frame(data, prepared)); + tails.set(lane, chain); track(chain); + // Drop an idle lane so the map does not grow with every session ever used. + void chain.finally(() => { + if (tails.get(lane) === chain) { + tails.delete(lane); + } + }); }); } +/** Lane for a frame: its params.session_id, or "" for frames that name no session. */ +function frame_lane(data: RawData): string { + const params = parse_frame_object(data)?.params; + if (typeof params !== "object" || params === null || Array.isArray(params)) { + return ""; + } + const session_id = (params as { session_id?: unknown }).session_id; + return typeof session_id === "string" && session_id.length > 0 ? `session:${session_id}` : ""; +} + interface PreparedSubmit { session_id: string; controller: AbortController; diff --git a/test/serve_prompts.test.ts b/test/serve_prompts.test.ts index f1a4dc7..b384e30 100644 --- a/test/serve_prompts.test.ts +++ b/test/serve_prompts.test.ts @@ -1104,11 +1104,12 @@ describe("serve prompt over websocket", () => { return body; }); await started; + // Same session_id, so this frame shares the submit's lane and must wait behind it. const health_promise = request_rpc(ws, { jsonrpc: "2.0", id: 4, method: "health", - params: { note: "prompt.abort" }, + params: { session_id, note: "prompt.abort" }, }).then((body) => { order.push(4); return body; @@ -1142,6 +1143,67 @@ describe("serve prompt over websocket", () => { } }); + it("runs different sessions on one connection concurrently, and session-less frames do not wait", async () => { + const work_dir = await make_temp_dir("serve-ws-lanes"); + const session_dir = path.join(work_dir, "sessions"); + let release_hang!: () => void; + const hang = new Promise((resolve) => { + release_hang = resolve; + }); + let hang_started!: () => void; + const started = new Promise((resolve) => { + hang_started = resolve; + }); + const fetch_fn: typeof fetch = async (_url, init) => { + const body = JSON.parse(String(init?.body)) as { messages: Array<{ content?: string }> }; + if (body.messages.some((message) => message.content === "hang")) { + hang_started(); + await hang; + } + return new Response( + JSON.stringify({ + model: "mock-model", + choices: [{ message: { role: "assistant", content: "done" }, finish_reason: "stop" }], + usage: { prompt_tokens: 1, completion_tokens: 1, total_tokens: 2 }, + }), + { status: 200 }, + ); + }; + const agent = mock_agent(work_dir, fetch_fn); + const server = create_serve_server({ port: 0, boot_stdout: null, session_dir, agent, version: "9.9.9" }); + servers.push(server); + const boot = await server.start(); + const ws = await open_ws(`ws://127.0.0.1:${boot.port}/?token=${encodeURIComponent(boot.token)}`); + try { + const create = async (id: number): Promise => + ((await request_rpc(ws, { jsonrpc: "2.0", id, method: "session.create", params: { source: "test" } })).result as { + session_id: string; + }).session_id; + const a = await create(1); + const b = await create(2); + const order: string[] = []; + const slow = request_rpc(ws, { jsonrpc: "2.0", id: 3, method: "prompt.submit", params: { session_id: a, text: "hang" } }).then( + (body) => { + order.push("a"); + return body; + }, + ); + await started; + const fast = await request_rpc(ws, { jsonrpc: "2.0", id: 4, method: "prompt.submit", params: { session_id: b, text: "hi" } }); + order.push("b"); + const health = await request_rpc(ws, { jsonrpc: "2.0", id: 5, method: "health" }); + order.push("health"); + expect(fast.result).toMatchObject({ session_id: b, stopped_reason: "final" }); + expect(health.result).toBeDefined(); + release_hang(); + expect((await slow).result).toMatchObject({ session_id: a, stopped_reason: "final" }); + expect(order).toEqual(["b", "health", "a"]); + } finally { + release_hang(); + ws.close(); + } + }); + it("stop() aborts an in-flight prompt instead of draining the model call", async () => { const work_dir = await make_temp_dir("serve-stop-abort"); const session_dir = path.join(work_dir, "sessions");