diff --git a/docs/collab.md b/docs/collab.md index edcc9127f..5f199625a 100644 --- a/docs/collab.md +++ b/docs/collab.md @@ -91,32 +91,18 @@ Known v1 limit for guests: a turn already streaming when you join becomes visibl |---|---|---| | `collab.relayUrl` | `wss://my.omp.sh` | Relay used by `/collab` when no relay is passed inline | | `collab.displayName` | OS username | Name shown to other participants | -| `share.serverUrl` | `https://my.omp.sh/s` | Share viewer/upload base used by `/share` (same Go service; links are `/#`) | +| `share.serverUrl` | `https://my.omp.sh/s` | Share viewer/upload base used by `/share` (links are `/#`) | | `share.redactSecrets` | `true` | Run the secret obfuscator over `/share` snapshots before upload | ## Self-hosting the relay -The relay is a small content-blind Go service (`omp-collab-relay`, in the pi-www repo under `relay/`). It keeps no state beyond live connections and exposes: +The relay is a small content-blind Go service. It keeps no state beyond live connections and exposes: - `GET /` — the static collab-web guest client (target of the `/collab` deep link), - `GET /r/?role=host|guest` — WebSocket upgrade, -- `POST /s` / `GET /s/` / `GET /s//raw` — `/share` blob upload, viewer page, and blob fetch (see the relay README), +- `POST /s` / `GET /s/` / `GET /s//raw` — `/share` blob upload, viewer page, and blob fetch, - `GET /healthz` — liveness. -Run it: - -```sh -go build -o omp-collab-relay . -RELAY_BIND=0.0.0.0:7475 ./omp-collab-relay -``` - -`RELAY_BIND` accepts `host:port`, a bare port (binds localhost), or a unix socket path (front it with a TLS-terminating reverse proxy — guests other than localhost require `wss://`). Then: - -``` -/collab my-relay.example.com -``` - -or set `collab.relayUrl` in `/settings`. ## Architecture notes diff --git a/docs/tools/lsp.md b/docs/tools/lsp.md index ca4be1f86..9365a1bce 100644 --- a/docs/tools/lsp.md +++ b/docs/tools/lsp.md @@ -71,7 +71,7 @@ - Optional: `timeout`. **Execution** -- `file: "*"`: `runWorkspaceDiagnostics()` detects project type from root markers and runs one subprocess command: Rust `cargo check --message-format=short`, TypeScript `npx tsc --noEmit`, Go `go build ./...`, Python `pyright`. +- `file: "*"`: `runWorkspaceDiagnostics()` detects project type from root markers and runs one subprocess command: Rust `cargo check --message-format=short`, TypeScript `npx tsc --noEmit`, Python `pyright`. - Concrete file or glob: `resolveDiagnosticTargets()` treats non-globs as one target, otherwise expands a `Bun.Glob` up to `MAX_GLOB_DIAGNOSTIC_TARGETS`. - Per file, every matching server runs: custom clients call `lint(file)`; real LSP servers optionally wait for project load, capture `diagnosticsVersion`, `refreshFile()`, then `waitForDiagnostics()` for fresh `publishDiagnostics` (settles on the latest publish; exact-version match accepted immediately). - Results are deduplicated by range+message and severity-sorted. diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index ba1884405..fea1506dc 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -815,11 +815,15 @@ async function runLoopBody( // Agent would stop here. Drain non-interrupting asides + follow-up messages. await config.onBeforeYield?.(); // Skip queue drains when externally aborted (same stranding hazard as above). + // Re-poll steering too: a steer can land between the stop-boundary dequeue + // above and this yield point (e.g. queued while onBeforeYield ran). Without + // this poll it would strand in the queue until the next manual prompt. + const lateSteering = signal?.aborted ? [] : (await config.getSteeringMessages?.()) || []; const asideMessages = signal?.aborted ? [] : resolveAsides(await config.getAsideMessages?.()); const followUpMessages = signal?.aborted ? [] : (await config.getFollowUpMessages?.()) || []; - if (asideMessages.length > 0 || followUpMessages.length > 0) { + if (lateSteering.length > 0 || asideMessages.length > 0 || followUpMessages.length > 0) { // Set as pending so the inner loop processes them before stopping. - pendingMessages = [...asideMessages, ...followUpMessages]; + pendingMessages = [...lateSteering, ...asideMessages, ...followUpMessages]; continue; } diff --git a/packages/agent/test/agent.test.ts b/packages/agent/test/agent.test.ts index 2bbaa3914..0e18bae30 100644 --- a/packages/agent/test/agent.test.ts +++ b/packages/agent/test/agent.test.ts @@ -80,6 +80,38 @@ describe("Agent", () => { expect(mock.calls.length).toBe(2); }); + it("delivers a steer that lands at the yield boundary instead of stranding it", async () => { + // Regression: a steering message queued after the stop-boundary dequeue + // (e.g. while onBeforeYield runs) was silently stranded in the queue until + // the next manual prompt. The outer yield drain must re-poll steering. + const mock = createMockModel({ responses: [{ content: ["First answer"] }, { content: ["Steer answer"] }] }); + const agent = new Agent({ streamFn: mock.stream }); + let injected = false; + agent.setOnBeforeYield(() => { + if (injected) return; + injected = true; + agent.steer({ + role: "user", + content: [{ type: "text", text: "Late steer" }], + steering: true, + timestamp: Date.now(), + }); + }); + + await agent.prompt("Initial"); + + expect(mock.calls.length).toBe(2); + expect(agent.hasQueuedMessages()).toBe(false); + const steerDelivered = agent.state.messages.some( + message => + message.role === "user" && + Array.isArray(message.content) && + message.content.some(part => part.type === "text" && part.text === "Late steer"), + ); + expect(steerDelivered).toBe(true); + expect(agent.state.messages[agent.state.messages.length - 1].role).toBe("assistant"); + }); + it("prompt() emits assistant error lifecycle for Anthropic output-blocked stream errors before assistant start", async () => { const mock = createMockModel({ responses: [] }); const errorText = "Output blocked by content filtering policy"; diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 147feaa9a..0b1a4fd6e 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -56,7 +56,7 @@ - Updated collab link handling to accept compact `roomId#key` links and relay hosts without explicit scheme when joining or starting sessions - Added `/collab stop` and `/collab status` options to control and inspect active shared sessions -- Added `/collab`, `/join`, and `/leave` for live session sharing: the host shares an end-to-end encrypted link (AES-256-GCM key only in the link fragment; the relay sees opaque bytes) and guests render the session natively in their own TUI — streaming text, tool cards, footer state, ctrl+o expansion, `/dump` — and can prompt or interrupt the host's agent. Guest prompts render with an author badge; session-mutating commands stay host-only. Defaults to the public `my.omp.sh` relay (`collab.relayUrl`); a self-hostable Go relay lives in the pi-www repo as `omp-collab-relay` +- Added `/collab`, `/join`, and `/leave` for live session sharing: the host shares an end-to-end encrypted link (AES-256-GCM key only in the link fragment; the relay sees opaque bytes) and guests render the session natively in their own TUI — streaming text, tool cards, footer state, ctrl+o expansion, `/dump` — and can prompt or interrupt the host's agent. Guest prompts render with an author badge; session-mutating commands stay host-only. Defaults to the public `my.omp.sh` relay (`collab.relayUrl`) - Added `collab.relayUrl` and `collab.displayName` settings plus a `collab` status-line segment showing the participant count (host) or guest role. Guests mirror the host's real model/thinking state into their replica agent and the host's subagent ecosystem end-to-end: the live subagent HUD, the Agent Hub table with live progress, hub chat/kill/revive (routed to the host), and on-demand subagent transcript viewing - Added `omp join ` subcommand that launches the interactive TUI and immediately runs `/join ` diff --git a/packages/coding-agent/src/collab/host.ts b/packages/coding-agent/src/collab/host.ts index b3b8fdeb1..47b1bd6d8 100644 --- a/packages/coding-agent/src/collab/host.ts +++ b/packages/coding-agent/src/collab/host.ts @@ -29,6 +29,7 @@ import { COLLAB_PROTO, type CollabFrame, type CollabParticipant, + type CollabPromptDetails, type CollabSessionState, formatCollabLink, formatCollabWebLink, @@ -364,13 +365,24 @@ export class CollabHost { const name = peer.name; const content: string | (TextContent | ImageContent)[] = images && images.length > 0 ? [{ type: "text", text }, ...images] : text; + const details: CollabPromptDetails & { __pendingDisplayTag?: string } = { from: name }; + if (this.#ctx.session.isStreaming) { + // Mid-turn guest prompts are steered: register the pending-display twin + // so queuedMessageCount reflects the queued steer (host pending bar + + // guests' "queued ×N" badge). The tag dequeues the entry when the agent + // consumes the message (mirrors the skill-prompt path). + details.__pendingDisplayTag = this.#ctx.session.enqueueCustomMessageDisplay(text, "steer"); + this.#ctx.updatePendingMessagesDisplay(); + this.#ctx.ui.requestRender(); + this.#scheduleStateBroadcast(); + } this.#ctx.session .promptCustomMessage( { customType: COLLAB_PROMPT_MESSAGE_TYPE, content, display: true, - details: { from: name }, + details, attribution: "user", }, { streamingBehavior: "steer" }, diff --git a/packages/coding-agent/src/config/settings-schema.ts b/packages/coding-agent/src/config/settings-schema.ts index 3fb3db81b..35262eeec 100644 --- a/packages/coding-agent/src/config/settings-schema.ts +++ b/packages/coding-agent/src/config/settings-schema.ts @@ -1339,7 +1339,7 @@ export const SETTINGS_SCHEMA = { tab: "interaction", group: "Collab", label: "Relay URL", - description: "Relay used by /collab (wss://host[:port]; self-host with the omp-collab-relay service)", + description: "Relay used by /collab (wss://host[:port])", }, }, @@ -3948,9 +3948,6 @@ export const SETTINGS_SCHEMA = { "dev.autoqaPush.endpoint": { type: "string", - // Bundled QA collector — runs `/work/pi-www/autoqa` behind qa.omp.sh. - // Override via `PI_AUTO_QA_PUSH_URL` or `dev.autoqaPush.endpoint` - // in `config.yml` to point at a self-hosted instance. default: "https://qa.omp.sh/v1/grievances" as const, ui: { tab: "tools", diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index db7bf48f9..bc320d9a9 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -1139,9 +1139,24 @@ export class AgentSession { if (this.#promptInFlightCount === 0) { this.#releasePowerAssertion(); this.#flushPendingAgentEnd(); + this.#drainStrandedQueuedMessages(); } } + /** A steer/follow-up can land after the agent loop's final queue poll but + * before the prompt unwinds: #promptInFlightCount keeps isStreaming true + * through post-prompt recovery, so senders (collab guests, skills) still + * queue via agent.steer()/followUp() instead of starting a fresh prompt. + * Without a drain those messages strand invisibly until the next manual + * prompt. Runs when the session settles; the guard makes it a no-op when + * the queue was consumed normally or a new turn already started. */ + #drainStrandedQueuedMessages(): void { + if (!this.agent.hasQueuedMessages()) return; + this.#scheduleAgentContinue({ + shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(), + }); + } + #resetInFlight(): void { this.#promptInFlightCount = 0; this.#releasePowerAssertion(); @@ -5255,6 +5270,14 @@ export class AgentSession { } else { this.agent.steer(normalizedAppMessage); } + // The isStreaming check above can be stale: image normalization is + // awaited, so the turn may have ended in between, leaving the message + // queued on an idle agent. Mirror #queueSteer's idle-path delivery. + if (this.#canAutoContinueForFollowUp()) { + this.#scheduleAgentContinue({ + shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(), + }); + } return; } diff --git a/packages/coding-agent/test/agent-session-queued-steer-delivery.test.ts b/packages/coding-agent/test/agent-session-queued-steer-delivery.test.ts new file mode 100644 index 000000000..e6e07aa19 --- /dev/null +++ b/packages/coding-agent/test/agent-session-queued-steer-delivery.test.ts @@ -0,0 +1,159 @@ +/** + * Contract: a custom message steered into a streaming session (the collab-host + * and skill-prompt path: `promptCustomMessage(..., { streamingBehavior: "steer" })`) + * is always delivered — never silently stranded in the agent's steering queue. + * + * Two regression seams, both observed as "guest messages just disappear" in + * collab sessions: + * 1. A steer landing at the run's yield boundary (after the stop-boundary + * dequeue) must force another turn instead of stranding. + * 2. A steer landing while the prompt unwinds (isStreaming stays true through + * post-prompt recovery, but the loop is already done) must be drained when + * the session settles. + */ +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { Agent } from "@oh-my-pi/pi-agent-core"; +import { createMockModel, type MockModel } from "@oh-my-pi/pi-ai/providers/mock"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session"; +import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage"; +import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; +import { Snowflake } from "@oh-my-pi/pi-utils"; + +const COLLAB_PROMPT_TYPE = "collab-prompt"; + +interface SteerHarness { + session: AgentSession; + sessionManager: SessionManager; + mock: MockModel; +} + +describe("AgentSession queued steer delivery", () => { + let tempDir: string; + let session: AgentSession; + const authStorages: AuthStorage[] = []; + + beforeEach(() => { + tempDir = path.join(os.tmpdir(), `pi-steer-strand-${Snowflake.next()}`); + fs.mkdirSync(tempDir, { recursive: true }); + }); + + afterEach(async () => { + await session?.dispose(); + for (const authStorage of authStorages.splice(0)) { + authStorage.close(); + } + fs.rmSync(tempDir, { recursive: true, force: true }); + }); + + async function createSession(responses: { content: string[] }[]): Promise { + const model = getBundledModel("anthropic", "claude-sonnet-4-5")!; + const mock = createMockModel({ responses }); + const agent = new Agent({ + getApiKey: () => "test-key", + initialState: { model, systemPrompt: ["Test"], tools: [] }, + streamFn: mock.stream, + }); + const sessionManager = SessionManager.inMemory(); + const settings = Settings.isolated({ "compaction.enabled": false }); + const authStorage = await AuthStorage.create(path.join(tempDir, `auth-${Snowflake.next()}.db`)); + authStorages.push(authStorage); + authStorage.setRuntimeApiKey("anthropic", "test-key"); + const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml")); + session = new AgentSession({ agent, sessionManager, settings, modelRegistry }); + return { session, sessionManager, mock }; + } + + function steerCollabPrompt(target: AgentSession, text: string): Promise { + return target.promptCustomMessage( + { + customType: COLLAB_PROMPT_TYPE, + content: text, + display: true, + details: { from: "guest" }, + attribution: "user", + }, + { streamingBehavior: "steer" }, + ); + } + + /** Resolves with the entry text when a collab-prompt entry is persisted. */ + function nextCollabEntry(sessionManager: SessionManager): Promise { + const { promise, resolve } = Promise.withResolvers(); + sessionManager.onEntryAppended = entry => { + if (entry.type === "custom_message" && entry.customType === COLLAB_PROMPT_TYPE) { + resolve(typeof entry.content === "string" ? entry.content : JSON.stringify(entry.content)); + } + }; + return promise; + } + + it("delivers a collab steer that lands at the run's yield boundary", async () => { + const { session, sessionManager, mock } = await createSession([ + { content: ["host answer"] }, + { content: ["ack guest"] }, + ]); + const entryAppended = nextCollabEntry(sessionManager); + + let streamingAtInject: boolean | undefined; + let injected = false; + session.agent.setOnBeforeYield(async () => { + if (injected) return; + injected = true; + // The session is still mid-prompt here, so this takes the steer path. + streamingAtInject = session.isStreaming; + await steerCollabPrompt(session, "guest steer at yield"); + }); + + await session.prompt("hello"); + + expect(streamingAtInject).toBe(true); + expect(await entryAppended).toBe("guest steer at yield"); + expect(mock.calls.length).toBe(2); + expect(session.agent.hasQueuedMessages()).toBe(false); + }); + + it("drains a steer stranded in the agent queue when the session settles", async () => { + const { session, sessionManager, mock } = await createSession([ + { content: ["host answer"] }, + { content: ["ack guest"] }, + ]); + const entryAppended = nextCollabEntry(sessionManager); + + // Inject from the wire agent_end subscriber: it fires synchronously while + // the session settles (#promptInFlightCount just hit 0), after the agent + // loop's final queue poll — a message queued here is invisible to the run + // and must be picked up by the settle-time drain. + const secondRunDone = Promise.withResolvers(); + let agentEnds = 0; + session.subscribe(event => { + if (event.type !== "agent_end") return; + agentEnds++; + if (agentEnds === 1) { + session.agent.steer({ + role: "custom", + customType: COLLAB_PROMPT_TYPE, + content: "guest steer at settle", + display: true, + details: { from: "guest" }, + attribution: "user", + timestamp: Date.now(), + }); + } else if (agentEnds === 2) { + secondRunDone.resolve(); + } + }); + + await session.prompt("hello"); + expect(await entryAppended).toBe("guest steer at settle"); + await secondRunDone.promise; + + expect(mock.calls.length).toBe(2); + expect(session.agent.hasQueuedMessages()).toBe(false); + }); +}); diff --git a/packages/coding-agent/test/collab/read-only.test.ts b/packages/coding-agent/test/collab/read-only.test.ts index 30563ea20..40b82c448 100644 --- a/packages/coding-agent/test/collab/read-only.test.ts +++ b/packages/coding-agent/test/collab/read-only.test.ts @@ -25,7 +25,7 @@ interface RelayData { type RelaySocket = Bun.ServerWebSocket; -/** Single-room relay mirroring the omp-collab-relay forwarding contract. */ +/** Single-room test relay mirroring the production forwarding contract. */ function startTestRelay(): { url: string; stop(): void } { let host: RelaySocket | null = null; const guests = new Map(); diff --git a/packages/coding-agent/test/collab/steer-queue.test.ts b/packages/coding-agent/test/collab/steer-queue.test.ts new file mode 100644 index 000000000..6f540c3d6 --- /dev/null +++ b/packages/coding-agent/test/collab/steer-queue.test.ts @@ -0,0 +1,214 @@ +/** + * Contract: a guest prompt that arrives while the host agent is streaming is + * steered AND becomes visible as a queued message — the host registers the + * pending-display twin (so `queuedMessageCount` covers it, feeding the web + * composer's "queued ×N" badge and the TUI pending bar) and tags the custom + * message so the entry is dequeued when the agent consumes it. + */ +import { afterEach, describe, expect, it } from "bun:test"; +import { importRoomKey } from "@oh-my-pi/pi-coding-agent/collab/crypto"; +import { CollabHost } from "@oh-my-pi/pi-coding-agent/collab/host"; +import { + COLLAB_PROTO, + type CollabFrame, + parseCollabLink, + rewriteEnvelopePeer, + unpackEnvelope, +} from "@oh-my-pi/pi-coding-agent/collab/protocol"; +import { CollabSocket } from "@oh-my-pi/pi-coding-agent/collab/relay-client"; +import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types"; + +interface RelayData { + role: "host" | "guest"; + peerId: number; +} + +type RelaySocket = Bun.ServerWebSocket; + +/** Single-room relay mirroring the omp-collab-relay forwarding contract. */ +function startTestRelay(): { url: string; stop(): void } { + let host: RelaySocket | null = null; + const guests = new Map(); + let nextPeerId = 1; + const server = Bun.serve({ + port: 0, + fetch(req, srv): Response | undefined { + const role = new URL(req.url).searchParams.get("role") === "host" ? "host" : "guest"; + const data: RelayData = { role, peerId: 0 }; + if (srv.upgrade(req, { data })) return undefined; + return new Response("upgrade failed", { status: 400 }); + }, + websocket: { + open(ws: RelaySocket): void { + if (ws.data.role === "host") { + host = ws; + return; + } + ws.data.peerId = nextPeerId++; + guests.set(ws.data.peerId, ws); + host?.send(JSON.stringify({ t: "peer-joined", peer: ws.data.peerId })); + }, + message(ws: RelaySocket, message: string | Buffer): void { + if (typeof message === "string") return; + const bytes = new Uint8Array(message); + if (ws.data.role === "host") { + const envelope = unpackEnvelope(bytes); + if (!envelope) return; + if (envelope.peerId === 0) { + for (const guest of guests.values()) guest.send(bytes); + } else { + guests.get(envelope.peerId)?.send(bytes); + } + return; + } + rewriteEnvelopePeer(bytes, ws.data.peerId); + host?.send(bytes); + }, + close(ws: RelaySocket): void { + if (ws.data.role === "guest") { + guests.delete(ws.data.peerId); + host?.send(JSON.stringify({ t: "peer-left", peer: ws.data.peerId })); + } + }, + }, + }); + return { url: `ws://localhost:${server.port}`, stop: () => server.stop(true) }; +} + +interface CapturedPrompt { + details?: { from?: string; __pendingDisplayTag?: string }; +} + +interface StreamingHostHarness { + ctx: InteractiveModeContext; + prompts: CapturedPrompt[]; + enqueued: { text: string; mode: string; tag: string }[]; + nextPrompt(): Promise; +} + +/** Context double for a host whose agent is mid-turn (isStreaming === true). */ +function makeStreamingHostContext(): StreamingHostHarness { + const prompts: CapturedPrompt[] = []; + const enqueued: { text: string; mode: string; tag: string }[] = []; + const promptWaiters: ((prompt: CapturedPrompt) => void)[] = []; + const ctx = { + settings: { get: () => "" }, + sessionManager: { + getSessionId: () => "sess-1", + getCwd: () => "/tmp", + snapshotForReplication: () => ({ + header: { type: "session", id: "sess-1", timestamp: new Date().toISOString(), cwd: "/tmp" }, + entries: [], + }), + onEntryAppended: undefined, + }, + session: { + isStreaming: true, + get queuedMessageCount(): number { + return enqueued.length; + }, + sessionName: "test", + model: undefined, + thinkingLevel: undefined, + subscribe: () => () => {}, + emitNotice: () => {}, + enqueueCustomMessageDisplay: (text: string, mode: string) => { + const tag = `tag-${enqueued.length + 1}`; + enqueued.push({ text, mode, tag }); + return tag; + }, + promptCustomMessage: (message: CapturedPrompt) => { + const captured: CapturedPrompt = { details: message.details }; + prompts.push(captured); + for (const waiter of promptWaiters.splice(0)) waiter(captured); + return Promise.resolve(); + }, + }, + eventBus: undefined, + statusLine: { + setCollabStatus: () => {}, + invalidate: () => {}, + getCachedContextBreakdown: () => ({ usedTokens: 0, contextWindow: 0 }), + }, + ui: { requestRender: () => {} }, + updatePendingMessagesDisplay: () => {}, + showStatus: () => {}, + collabHost: undefined, + } as unknown as InteractiveModeContext; + const nextPrompt = (): Promise => { + const { promise, resolve } = Promise.withResolvers(); + promptWaiters.push(resolve); + return promise; + }; + return { ctx, prompts, enqueued, nextPrompt }; +} + +interface TestGuest { + socket: CollabSocket; + nextFrame(): Promise; +} + +async function joinAsGuest(link: string, name: string): Promise { + const parsed = parseCollabLink(link); + if ("error" in parsed) throw new Error(parsed.error); + const writeToken = parsed.writeToken ? Buffer.from(parsed.writeToken).toString("base64url") : undefined; + const key = await importRoomKey(parsed.key); + const socket = new CollabSocket({ wsUrl: parsed.wsUrl, role: "guest", key }); + const queue: CollabFrame[] = []; + const waiters: ((frame: CollabFrame) => void)[] = []; + socket.onFrame = frame => { + const waiter = waiters.shift(); + if (waiter) waiter(frame); + else queue.push(frame); + }; + socket.onOpen = () => socket.send({ t: "hello", proto: COLLAB_PROTO, name, writeToken }); + socket.connect(); + const nextFrame = (): Promise => { + const queued = queue.shift(); + if (queued) return Promise.resolve(queued); + const { promise, resolve } = Promise.withResolvers(); + waiters.push(resolve); + return promise; + }; + return { socket, nextFrame }; +} + +const cleanups: (() => void | Promise)[] = []; + +afterEach(async () => { + for (const cleanup of cleanups.splice(0).reverse()) await cleanup(); +}); + +describe("collab mid-turn guest prompts", () => { + it("registers the steer as a queued message and reports it to guests via state", async () => { + const relay = startTestRelay(); + cleanups.push(relay.stop); + const harness = makeStreamingHostContext(); + const host = new CollabHost(harness.ctx); + await host.start(relay.url); + cleanups.push(() => host.stop("test done")); + + const guest = await joinAsGuest(host.link, "writer"); + cleanups.push(() => guest.socket.close()); + const welcome = await guest.nextFrame(); + if (welcome.t !== "welcome") throw new Error(`expected welcome, got ${welcome.t}`); + + const prompted = harness.nextPrompt(); + guest.socket.send({ t: "prompt", text: "steer the host" }); + const prompt = await prompted; + + // Display twin registered before dispatch, tag forwarded in details so + // the session dequeues the entry when the agent consumes the message. + expect(harness.enqueued).toEqual([{ text: "steer the host", mode: "steer", tag: "tag-1" }]); + expect(prompt.details).toEqual({ from: "writer", __pendingDisplayTag: "tag-1" }); + + // The queued steer must reach guests through state.queuedMessageCount — + // that field drives the web composer's "queued ×N" badge. + let sawQueuedCount = false; + for (let i = 0; i < 10 && !sawQueuedCount; i++) { + const frame = await guest.nextFrame(); + if (frame.t === "state" && frame.state.queuedMessageCount === 1) sawQueuedCount = true; + } + expect(sawQueuedCount).toBe(true); + }); +});