From 6534dd16b54f5751824e4b96d90185d8d745ccc3 Mon Sep 17 00:00:00 2001 From: david Date: Mon, 5 Oct 2026 17:09:48 +0800 Subject: [PATCH 1/4] feat(remote): queue follow-ups sent during a running remote turn Remote (SSH) sessions rejected any message sent while a turn was running, for every provider. The host now keeps a per-session queue: follow-ups sent mid-turn are stored on the host and dispatched in order as each turn ends, so they still go out while the desktop is disconnected. - host: add enqueue, dequeue, editQueued and resumeQueue commands; share turn start between sends, queued dispatch and resume (startTurn) - host: pause the queue when a turn is stopped or interrupted; Resume continues the interrupted work first, as local sessions do - host: advertise the sessions.queue capability; older hosts keep the previous behavior - desktop: route busy submits to the host queue and wire queue edit, remove and resume; hide Steer where the provider cannot steer - desktop: do not copy the host's queue into local session state, which would have sent queued messages a second time Co-Authored-By: Claude Opus 5.5 --- host/engine.test.ts | 92 ++++++ host/engine.ts | 306 +++++++++++++----- host/server.ts | 1 + src/features/connections/model/protocol.ts | 25 ++ .../model/remoteSessionState.test.ts | 4 + .../connections/model/remoteSessionState.ts | 5 +- .../connections/ui/RemoteSession.test.ts | 86 ++++- src/features/connections/ui/RemoteSession.tsx | 68 +++- src/features/sessions/ui/Composer.tsx | 18 +- src/features/sessions/ui/SessionPane.tsx | 9 +- 10 files changed, 522 insertions(+), 92 deletions(-) diff --git a/host/engine.test.ts b/host/engine.test.ts index 3f5b58687..a07132dfc 100644 --- a/host/engine.test.ts +++ b/host/engine.test.ts @@ -698,6 +698,88 @@ describe("headless session ownership", () => { expect(provider.answer).toHaveBeenCalledTimes(1); }); + it("queues follow-ups during a turn and sends them in order as each turn ends", async () => { + const { engine, store, turns, id } = setup(); + engine.command({ type: "send", commandId: "send", sessionId: id, text: "Work" }); + await vi.waitFor(() => expect(turns).toHaveLength(1)); + engine.command({ type: "enqueue", commandId: "q1", sessionId: id, text: "Next", intent: "plan" }); + engine.command({ type: "enqueue", commandId: "q2", sessionId: id, text: "Then" }); + expect(store.session(id).session).toMatchObject({ + queuedMessages: [{ id: "q1", text: "Next" }, { id: "q2", text: "Then" }], + queueStatus: "active", + }); + expect(turns).toHaveLength(1); + + turns[0].finish(); + await vi.waitFor(() => expect(turns).toHaveLength(2)); + expect(turns[1].input).toMatchObject({ text: "Next", intent: "plan" }); + expect(store.session(id).session.blocks.at(-1)).toMatchObject({ id: "q1", role: "user" }); + expect(store.session(id).session.queuedMessages).toEqual([ + expect.objectContaining({ id: "q2" }), + ]); + + turns[1].finish(); + await vi.waitFor(() => expect(turns).toHaveLength(3)); + expect(turns[2].input.text).toBe("Then"); + expect(store.session(id).session.queuedMessages).toBeUndefined(); + expect(store.session(id).session.queueStatus).toBeUndefined(); + turns[2].finish(); + }); + + it("sends a follow-up at once when it reaches the host after the turn ended", async () => { + const { engine, turns, id } = setup(); + engine.command({ type: "enqueue", commandId: "late", sessionId: id, text: "Next" }); + await vi.waitFor(() => expect(turns).toHaveLength(1)); + expect(turns[0].input.text).toBe("Next"); + turns[0].finish(); + }); + + it("edits and removes queued follow-ups until they are sent", async () => { + const { engine, turns, id } = setup(); + engine.command({ type: "send", commandId: "send", sessionId: id, text: "Work" }); + await vi.waitFor(() => expect(turns).toHaveLength(1)); + engine.command({ type: "enqueue", commandId: "q1", sessionId: id, text: "Next" }); + engine.command({ type: "enqueue", commandId: "q2", sessionId: id, text: "Then" }); + engine.command({ type: "editQueued", commandId: "edit", sessionId: id, messageId: "q1", text: "Edited" }); + expect(() => + engine.command({ type: "editQueued", commandId: "blank", sessionId: id, messageId: "q1", text: " " }), + ).toThrow("Invalid prompt"); + engine.command({ type: "dequeue", commandId: "drop", sessionId: id, messageId: "q2" }); + + turns[0].finish(); + await vi.waitFor(() => expect(turns).toHaveLength(2)); + expect(turns[1].input.text).toBe("Edited"); + expect(() => + engine.command({ type: "dequeue", commandId: "late-drop", sessionId: id, messageId: "q1" }), + ).toThrow("no longer waiting"); + turns[1].finish(); + }); + + it("pauses the queue when the turn is stopped and resumes the work before it", async () => { + const { engine, store, turns, id } = setup(); + engine.command({ type: "send", commandId: "send", sessionId: id, text: "Work" }); + await vi.waitFor(() => expect(turns).toHaveLength(1)); + engine.command({ type: "enqueue", commandId: "q1", sessionId: id, text: "Next" }); + expect(() => + engine.command({ type: "resumeQueue", commandId: "early", sessionId: id }), + ).toThrow("not paused"); + engine.command({ type: "cancel", commandId: "stop", sessionId: id, runId: store.session(id).runId }); + await vi.waitFor(() => expect(store.session(id).status).toBe("idle")); + expect(store.session(id).session.queueStatus).toBe("paused"); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(turns).toHaveLength(1); + + engine.command({ type: "resumeQueue", commandId: "resume", sessionId: id }); + await vi.waitFor(() => expect(turns).toHaveLength(2)); + expect(turns[1].input.text).toBe("Continue from where you left off."); + expect(store.session(id).session.queueStatus).toBe("resuming"); + + turns[1].finish(); + await vi.waitFor(() => expect(turns).toHaveLength(3)); + expect(turns[2].input.text).toBe("Next"); + turns[2].finish(); + }); + it("recovers interrupted durable state without replaying an uncertain provider send", async () => { const { store, provider, id } = setup(); const value = store.session(id); @@ -712,6 +794,8 @@ describe("headless session ownership", () => { ...value.session, busy: true, providerSessionId: "retained", + queuedMessages: [{ id: "queued", text: "Next", attachments: [] }], + queueStatus: "active", blocks: [ { id: "interrupted-turn", @@ -729,6 +813,8 @@ describe("headless session ownership", () => { expect(store.session(id).status).toBe("interrupted"); expect(store.session(id).session.busy).toBe(false); expect(store.session(id).session.blocks[0].durationMs).toBe(2_000); + // Queued follow-ups wait for the user to inspect the interrupted turn. + expect(store.session(id).session.queueStatus).toBe("paused"); expect(provider.send).not.toHaveBeenCalled(); expect(provider.bind).toHaveBeenCalledWith( id, @@ -893,5 +979,11 @@ describe("headless session ownership", () => { reply: { kind: "answered", answers: { a: [42] } }, }), ).toThrow(); + expect(() => + parseCommand({ type: "enqueue", commandId: "x", sessionId: "y", text: " " }), + ).toThrow("Invalid prompt"); + expect(() => + parseCommand({ type: "dequeue", commandId: "x", sessionId: "y" }), + ).toThrow("queued message ID"); }); }); diff --git a/host/engine.ts b/host/engine.ts index 27d72f91c..c5784dd84 100644 --- a/host/engine.ts +++ b/host/engine.ts @@ -19,8 +19,11 @@ import { canReplaceSessionTitle, formatSessionTitle, titleFromPrompt, + type Attachment, + type Block, type Session, } from "../src/features/sessions/model/session"; +import { CONTINUE_PROMPT } from "../src/features/sessions/model/inFlight"; import { namedWorktreeBranch } from "../src/features/source-control/model/worktrees"; import { isRemoteProvider, @@ -55,6 +58,9 @@ const text = (value: unknown, label: string, max = 128): string => { return value; }; +const prompt = (value: unknown): value is string => + typeof value === "string" && value.length <= 256_000 && !value.includes("\0"); + function modelSettings(value: unknown): Record { if (value == null) return {}; if (!value || typeof value !== "object" || Array.isArray(value)) @@ -124,19 +130,17 @@ export function parseCommand(input: unknown): HostCommand { }; } if (v.type === "compact") return { type: "compact", commandId, sessionId }; - if (v.type === "send" || v.type === "draft") { + if (v.type === "send" || v.type === "draft" || v.type === "enqueue") { const attachments = parseRemoteAttachments(v.attachments); if ( - typeof v.text !== "string" || - v.text.length > 256_000 || - v.text.includes("\0") || + !prompt(v.text) || (!v.text.trim() && attachments.length === 0 && !(v.type === "send" && v.draftBlockId !== undefined)) ) throw new Error("Invalid prompt"); if ( - v.type === "send" && + v.type !== "draft" && v.intent !== undefined && !["default", "plan", "build"].includes(String(v.intent)) ) @@ -152,7 +156,7 @@ export function parseCommand(input: unknown): HostCommand { sessionId, text: v.text, ...(attachments.length ? { attachments } : {}), - ...(v.type === "send" && v.intent + ...(v.type !== "draft" && v.intent ? { intent: v.intent as "default" | "plan" | "build" } : {}), ...(v.type === "send" && v.draftBlockId !== undefined @@ -170,6 +174,25 @@ export function parseCommand(input: unknown): HostCommand { sessionId, draftBlockId: text(v.draftBlockId, "draft block ID"), }; + if (v.type === "dequeue") + return { + type: "dequeue", + commandId, + sessionId, + messageId: text(v.messageId, "queued message ID"), + }; + if (v.type === "editQueued") { + // Emptiness depends on the queued attachments, checked against the session. + if (!prompt(v.text)) throw new Error("Invalid prompt"); + return { + type: "editQueued", + commandId, + sessionId, + messageId: text(v.messageId, "queued message ID"), + text: v.text, + }; + } + if (v.type === "resumeQueue") return { type: "resumeQueue", commandId, sessionId }; const runId = text(v.runId, "run ID"); if (v.type === "cancel") return { type: "cancel", commandId, sessionId, runId }; @@ -316,6 +339,8 @@ export class HostEngine { return await action(); } finally { this.switchingProjects.delete(projectId); + for (const session of this.store.summaries(projectId)) + this.dispatchQueued(session.id); } } @@ -455,7 +480,9 @@ export class HostEngine { } else { value = this.store.session(command.sessionId); if ( - (command.type === "send" || command.type === "compact") && + (command.type === "send" || + command.type === "compact" || + command.type === "resumeQueue") && this.switchingProjects.has(value.projectId) ) throw new Error("Wait for the branch switch to finish"); @@ -558,85 +585,83 @@ export class HostEngine { ? (draft?.attachments ?? resolveAttachments(this.store, command.attachments ?? [])) : []; - const runId = randomUUID(); - const firstTurn = - command.type === "send" && - !value.session.blocks.some((block) => !block.draft); - const placeholderTitle = - value.session.title === "New remote session" || - canReplaceSessionTitle( - value.session.title, - value.session.harness, - HARNESS_LABEL[value.session.harness], - ); - const model = resolveModel( - value.session.harness, - value.session.model, + ({ value, effect } = this.startTurn(value, { + blockId: command.commandId, + prompt: command.type === "compact" ? null : command.text, + intent: command.type === "send" ? command.intent : undefined, + attachments, + plan, + })); + } else if (command.type === "enqueue") { + const attachments = resolveAttachments( + this.store, + command.attachments ?? [], ); value = { ...value, - status: "running", - runId, session: { ...value.session, - busy: true, - pendingQuestion: undefined, - title: - firstTurn && placeholderTitle - ? titleFromPrompt( - command.text, - value.session.harness, - attachments, - ) - : value.session.title, - blocks: [ - ...value.session.blocks - .filter((block) => !block.draft) - .map((block) => - block === plan - ? { - ...block, - plan: { - ...(block.plan ?? { status: "ready" as const }), - status: "building" as const, - approvedText: block.text, - }, - } - : block, - ), + queuedMessages: [ + ...(value.session.queuedMessages ?? []), { id: command.commandId, - role: "user", - text: command.type === "compact" ? "/compact" : command.text, - ...(attachments.length ? { attachments } : {}), - startedAt: Date.now(), - turnModel: { - harness: value.session.harness, - id: value.session.model, - name: - model.id === value.session.model - ? model.name - : value.session.model.replace(/^[^:]+:/, ""), - }, + text: command.text, + attachments, + ...(command.intent ? { intent: command.intent } : {}), }, ], + queueStatus: value.session.queueStatus ?? "active", }, }; - effect = (saved) => { - this.run( - saved, - command.type === "compact" ? null : command.text, - command.type === "send" ? command.intent : undefined, - attachments, - ); - if (firstTurn && command.type === "send") { - this.generateFirstTurnNames( - saved, - command.text, - placeholderTitle, - ); - } + // The turn it was meant to follow may have ended in the meantime. + effect = () => this.dispatchQueued(command.sessionId); + } else if ( + command.type === "dequeue" || + command.type === "editQueued" + ) { + const queued = value.session.queuedMessages ?? []; + const target = queued.find( + (message) => message.id === command.messageId, + ); + if (!target) + throw new Error("This queued message is no longer waiting"); + if ( + command.type === "editQueued" && + !command.text.trim() && + !target.attachments.length + ) + throw new Error("Invalid prompt"); + const next = + command.type === "dequeue" + ? queued.filter((message) => message !== target) + : queued.map((message) => + message === target + ? { ...message, text: command.text } + : message, + ); + value = { + ...value, + session: { + ...value.session, + queuedMessages: next.length ? next : undefined, + queueStatus: next.length ? value.session.queueStatus : undefined, + }, }; + } else if (command.type === "resumeQueue") { + if ( + value.status === "running" || + value.session.queueStatus !== "paused" || + !value.session.queuedMessages?.length + ) + throw new Error("The queue is not paused"); + // As locally: continue the interrupted work, then the queue. + ({ value, effect } = this.startTurn( + { + ...value, + session: { ...value.session, queueStatus: "resuming" }, + }, + { blockId: command.commandId, prompt: CONTINUE_PROMPT }, + )); } else { if (value.runId !== command.runId || value.status !== "running") throw new Error( @@ -711,6 +736,126 @@ export class HostEngine { return receipt; } + /** Marks the session running with the turn's message; the effect starts + * the provider once that state is saved. */ + private startTurn( + value: HostSession, + turn: { + blockId: string; + /** `null` compacts the context instead of sending a message. */ + prompt: string | null; + intent?: "default" | "plan" | "build"; + attachments?: Attachment[]; + plan?: Block; + }, + ): { value: HostSession; effect: (saved: HostSession) => void } { + const { blockId, prompt, intent, attachments = [], plan } = turn; + const runId = randomUUID(); + const firstTurn = + prompt !== null && !value.session.blocks.some((block) => !block.draft); + const placeholderTitle = + value.session.title === "New remote session" || + canReplaceSessionTitle( + value.session.title, + value.session.harness, + HARNESS_LABEL[value.session.harness], + ); + const model = resolveModel(value.session.harness, value.session.model); + return { + value: { + ...value, + status: "running", + runId, + session: { + ...value.session, + busy: true, + pendingQuestion: undefined, + title: + firstTurn && placeholderTitle + ? titleFromPrompt(prompt, value.session.harness, attachments) + : value.session.title, + blocks: [ + ...value.session.blocks + .filter((block) => !block.draft) + .map((block) => + block === plan + ? { + ...block, + plan: { + ...(block.plan ?? { status: "ready" as const }), + status: "building" as const, + approvedText: block.text, + }, + } + : block, + ), + { + id: blockId, + role: "user", + text: prompt ?? "/compact", + ...(attachments.length ? { attachments } : {}), + startedAt: Date.now(), + turnModel: { + harness: value.session.harness, + id: value.session.model, + name: + model.id === value.session.model + ? model.name + : value.session.model.replace(/^[^:]+:/, ""), + }, + }, + ], + }, + }, + effect: (saved) => { + this.run(saved, prompt, intent, attachments); + if (firstTurn) + this.generateFirstTurnNames(saved, prompt, placeholderTitle); + }, + }; + } + + /** Sends the next queued follow-up once its session is idle. */ + private dispatchQueued(id: string): void { + if (this.closing || this.running.has(id)) return; + try { + const value = this.store.session(id); + const { queuedMessages = [], queueStatus } = value.session; + const [head, ...rest] = queuedMessages; + if ( + !head || + value.status === "running" || + queueStatus === "paused" || + queueStatus === "resuming" || + this.switchingProjects.has(value.projectId) || + !this.providers[value.session.harness as RemoteProvider] + ) + return; + const { value: started, effect } = this.startTurn( + { + ...value, + session: { + ...value.session, + queuedMessages: rest.length ? rest : undefined, + queueStatus: rest.length ? queueStatus : undefined, + }, + }, + { + blockId: head.id, + prompt: head.text, + intent: head.intent === "orchestrate" ? undefined : head.intent, + attachments: head.attachments, + }, + ); + effect(this.save(started, { type: "queue.dispatched", messageId: head.id })); + } catch (error) { + console.error( + "Sending a queued message failed:", + error instanceof Error ? error.message : "unknown error", + ); + } + } + private generateFirstTurnNames( value: HostSession, message: string, @@ -839,6 +984,8 @@ export class HostEngine { latest, this.closing || active.persistenceFailed ? "interrupted" : "idle", message, + undefined, + active.cancelled, ), { type: "settled", error, cancelled: active.cancelled }, ); @@ -849,6 +996,7 @@ export class HostEngine { const persisted = this.store.session(session.id).session; if (persisted.providerSessionId) provider.bind(session.id, persisted.providerSessionId, persisted.cwd); + this.dispatchQueued(session.id); }) .catch((error) => { clearTimeout(this.live.get(session.id)?.timer); @@ -883,10 +1031,20 @@ export class HostEngine { status: "idle" | "interrupted", message?: string, endedAt = Date.now(), + cancelled = false, ): HostSession { const stopped = stopStreaming(value.session, endedAt); + // Follow-ups wait after an interruption, as locally, until resumed. + const queueStatus = !stopped.queuedMessages?.length + ? undefined + : cancelled || status === "interrupted" + ? ("paused" as const) + : stopped.queueStatus === "resuming" + ? ("active" as const) + : stopped.queueStatus; const session = { ...stopped, + queueStatus, blocks: stopped.blocks.map((block) => block.role === "plan" && block.plan?.status === "building" ? { diff --git a/host/server.ts b/host/server.ts index 62e6947cc..e87c2c8f8 100644 --- a/host/server.ts +++ b/host/server.ts @@ -288,6 +288,7 @@ export function createHostServer( "attachments.read", "sessions.draft", "sessions.plan", + "sessions.queue", ], }; break; diff --git a/src/features/connections/model/protocol.ts b/src/features/connections/model/protocol.ts index 21ce9f96f..d7acf1a89 100644 --- a/src/features/connections/model/protocol.ts +++ b/src/features/connections/model/protocol.ts @@ -184,6 +184,31 @@ export type HostCommand = sessionId: string; draftBlockId: string; } + /** Follow-ups wait on the host, which sends each when the turn before it + * finishes, so they still go out while this computer is disconnected. */ + | { + type: "enqueue"; + commandId: string; + sessionId: string; + text: string; + attachments?: RemoteAttachment[]; + intent?: "default" | "plan" | "build"; + } + | { + type: "dequeue"; + commandId: string; + sessionId: string; + messageId: string; + } + | { + type: "editQueued"; + commandId: string; + sessionId: string; + messageId: string; + text: string; + } + /** Continues the interrupted turn, then the paused queue. */ + | { type: "resumeQueue"; commandId: string; sessionId: string } | { type: "cancel"; commandId: string; sessionId: string; runId: string } | { type: "approve"; diff --git a/src/features/connections/model/remoteSessionState.test.ts b/src/features/connections/model/remoteSessionState.test.ts index a3d2181c7..c9188ebbe 100644 --- a/src/features/connections/model/remoteSessionState.test.ts +++ b/src/features/connections/model/remoteSessionState.test.ts @@ -20,6 +20,8 @@ it("keeps a host transcript in normal session state under its local tab ID", () { id: "turn", role: "user", text: "Fix it" }, { id: "plan", role: "plan", text: "Build the fix" }, ], + queuedMessages: [{ id: "next", text: "Then test", attachments: [] }], + queueStatus: "active", }, } as HostSession; const session = remoteSessionState(shell, snapshot, { @@ -35,5 +37,7 @@ it("keeps a host transcript in normal session state under its local tab ID", () title: "Fix the build", blocks: snapshot.session.blocks, }); + expect(session.queuedMessages).toBeUndefined(); + expect(session.queueStatus).toBeUndefined(); expect(shouldPersistSession(session)).toBe(false); }); diff --git a/src/features/connections/model/remoteSessionState.ts b/src/features/connections/model/remoteSessionState.ts index e5537211d..b11d1c97a 100644 --- a/src/features/connections/model/remoteSessionState.ts +++ b/src/features/connections/model/remoteSessionState.ts @@ -9,7 +9,10 @@ export function remoteSessionState( snapshot: HostSession, project: RemoteProject, ): Session { - const host = snapshot.session; + // The host sends its own queued follow-ups; copied here, this app would send + // them a second time. + const { queuedMessages: _queued, queueStatus: _status, ...host } = + snapshot.session; return { ...shell, ...host, diff --git a/src/features/connections/ui/RemoteSession.test.ts b/src/features/connections/ui/RemoteSession.test.ts index 9fdbc6204..66c8d78be 100644 --- a/src/features/connections/ui/RemoteSession.test.ts +++ b/src/features/connections/ui/RemoteSession.test.ts @@ -163,7 +163,7 @@ beforeEach(() => { environmentId: "env", name: "home", providers, - capabilities: ["attachments.upload", "sessions.plan", "sessions.draft"], + capabilities: ["attachments.upload", "sessions.plan", "sessions.draft", "sessions.queue"], }; if (method === "models.list") { if (catalog instanceof Error) throw catalog.message; @@ -328,6 +328,32 @@ function dispatch(command: HostCommand) { ], }, }; + } else if (host && command.type === "enqueue") { + host = { + ...host, + revision: host.revision + 1, + session: { + ...host.session, + queuedMessages: [ + ...(host.session.queuedMessages ?? []), + { id: command.commandId, text: command.text, attachments: [] }, + ], + queueStatus: host.session.queueStatus ?? "active", + }, + }; + } else if (host && command.type === "dequeue") { + const queuedMessages = host.session.queuedMessages?.filter( + (message) => message.id !== command.messageId, + ); + host = { + ...host, + revision: host.revision + 1, + session: { + ...host.session, + queuedMessages: queuedMessages?.length ? queuedMessages : undefined, + queueStatus: queuedMessages?.length ? host.session.queueStatus : undefined, + }, + }; } else if (host && command.type === "removeDraft") { host = { ...host, @@ -1029,6 +1055,64 @@ it("holds a settings change during a running turn and applies it afterwards", as ); }); +it("queues a message sent during a running turn on the host", async () => { + await render(); + await send("First"); + host = { + ...host!, + revision: host!.revision + 1, + status: "running", + runId: "run", + session: { ...host!.session, busy: true }, + }; + await vi.waitFor(() => expect(byLabel("Stop")).not.toBeNull(), { + timeout: 4_000, + }); + await send("Second"); + expect(commands.at(-1)).toMatchObject({ + type: "enqueue", + sessionId: "host-session", + text: "Second", + }); + await vi.waitFor( + () => expect(byLabel("Remove queued message")).not.toBeNull(), + { timeout: 4_000 }, + ); + expect(container.querySelector("[data-message-queue]")?.textContent).toContain("Second"); + // Host providers cannot steer a running turn. + expect(container.querySelector("[data-message-queue]")?.textContent).not.toContain("Steer"); + + const queued = commands.at(-1)!.commandId; + await act(async () => byLabel("Remove queued message")!.click()); + await settle(); + expect(commands.at(-1)).toMatchObject({ type: "dequeue", messageId: queued }); +}); + +it("resumes a queue the host paused after a stop", async () => { + await render(); + await send("First"); + host = { + ...host!, + revision: host!.revision + 1, + session: { + ...host!.session, + queuedMessages: [{ id: "next", text: "Next", attachments: [] }], + queueStatus: "paused", + }, + }; + const resume = () => + [...container.querySelectorAll("button")].find( + (button) => button.textContent === "Resume", + ); + await vi.waitFor(() => expect(resume()).toBeTruthy(), { timeout: 4_000 }); + await act(async () => resume()!.click()); + await settle(); + expect(commands.at(-1)).toMatchObject({ + type: "resumeQueue", + sessionId: "host-session", + }); +}); + it.each([false, true])("retries a lost create response without duplicating the first turn (remount: %s)", async (remount) => { const original = vi.mocked(invoke).getMockImplementation()!; let accepted: ReturnType | undefined; diff --git a/src/features/connections/ui/RemoteSession.tsx b/src/features/connections/ui/RemoteSession.tsx index 781e4abfa..2ceb2e11f 100644 --- a/src/features/connections/ui/RemoteSession.tsx +++ b/src/features/connections/ui/RemoteSession.tsx @@ -323,6 +323,7 @@ function ConnectedRemoteSession({ !hasHostBlock(pending.commandId); const busy = !!hostSession?.busy || unseenActive || (startingActive && !starting?.draft) || pendingSendActive; + const queueSupported = !!descriptor?.capabilities.includes("sessions.queue"); // An accepted turn stays on screen until a sync shows the host's copy, so // the transcript never drops it for a moment in between. useEffect(() => { @@ -815,6 +816,50 @@ function ConnectedRemoteSession({ } }; + // The host holds follow-ups sent during a turn and sends each in order. + const enqueue = ( + id: string, + text: string, + attachments: Attachment[], + intent: "default" | "plan" | "build", + ): boolean => { + const version = bindingVersion.current; + preparingRef.current = true; + void uploadRemoteAttachments(machine.id, attachments) + .then((refs) => { + if (!alive.current || version !== bindingVersion.current) return; + return run({ + type: "enqueue", + commandId: crypto.randomUUID(), + sessionId: id, + text, + attachments: refs, + intent, + }); + }) + .catch((reason) => { + if (alive.current && version === bindingVersion.current) + setError(String(reason)); + }) + .finally(() => { + preparingRef.current = false; + }); + return true; + }; + const queueCommand = ( + command: + | { type: "dequeue"; messageId: string } + | { type: "editQueued"; messageId: string; text: string } + | { type: "resumeQueue" }, + ) => { + if (!hostSession || !online || pending) return; + void run({ + ...command, + commandId: crypto.randomUUID(), + sessionId: hostSession.id, + }); + }; + const message = ( id: string, text: string, @@ -852,7 +897,6 @@ function ConnectedRemoteSession({ sending || preparingRef.current || pending || - busy || (!text.trim() && !attachments.length) ) return false; @@ -860,6 +904,16 @@ function ConnectedRemoteSession({ options?.intent === "plan" || options?.intent === "build" ? options.intent : "default"; + if (busy) + return ( + !!hostSession && + queueSupported && + !asDraft && + !planBlockId && + !options?.draftBlockId && + message(hostSession.id, text, undefined, [], intent).type === "send" && + enqueue(hostSession.id, text, attachments, intent) + ); const turn = optimisticTurn( text, attachments, @@ -1245,11 +1299,15 @@ function ConnectedRemoteSession({ return true; }, onPlaceSessionInFolder: noop, - onDeleteQueuedMessage: noop, - onEditQueuedMessage: noop, + onDeleteQueuedMessage: (_, messageId) => + queueCommand({ type: "dequeue", messageId }), + onEditQueuedMessage: (_, messageId, text) => + queueCommand({ type: "editQueued", messageId, text }), + // The host may send the message being edited; saving then reports that. onQueuedMessageEditingChange: noop, - onSteerQueuedMessage: noop, - onResumeQueue: noop, + // Host providers cannot add to a running turn. + onSteerQueuedMessage: undefined, + onResumeQueue: () => queueCommand({ type: "resumeQueue" }), onUsageLimitResume: noop, onUsageLimitResumeAtReset: noop, onUsageLimitDismiss: noop, diff --git a/src/features/sessions/ui/Composer.tsx b/src/features/sessions/ui/Composer.tsx index b93bce467..c9ed41471 100644 --- a/src/features/sessions/ui/Composer.tsx +++ b/src/features/sessions/ui/Composer.tsx @@ -464,14 +464,16 @@ function MessageQueue({ {label} - + {onSteer ? ( + + ) : null}