diff --git a/docs/remote-access.md b/docs/remote-access.md index 55a7bc822..52768889c 100644 --- a/docs/remote-access.md +++ b/docs/remote-access.md @@ -37,6 +37,8 @@ The **Explorer** sidebar and Go to File use the normal file views for host folde The composer’s **+** menu supports file and image attachments, Plan mode, and saved drafts when the host advertises these capabilities. Attachments are copied to the host’s private data directory before the turn or draft is recorded; each file is limited to 20 MiB. Image previews are restored from the host when you reopen a conversation. A draft can be sent or removed from its transcript card. Plan mode uses the host provider and produces a reviewable plan card whose Build action continues on the host with the session’s current model. Update older hosts to enable these menu actions. +Messages sent while a turn is running wait in a queue on the host, which sends each one when the turn before it finishes, also while this computer is disconnected. Queued messages can be edited or removed until they are sent. Stopping a turn pauses the queue; **Resume** continues the interrupted work first, then the queue. Host providers cannot steer a running turn. Older hosts reject messages sent during a turn; update the host to queue them. + Features that read or run on this computer are not available in these projects: `@` file mentions, skills and slash commands other than `/plan` and `/compact`, operator mode, and terminals. Worktree deletion and the local worktree settings page are not available remotely yet. Plans from the transcript open normal read-only plan tabs. Source files stay on the host; this feature shares host-owned sessions, not working-directory synchronization. ## Manual connection (advanced / development) @@ -100,7 +102,7 @@ The release workflow publishes `monocode-host-{darwin,linux}-{arm64,x64}.tar.gz` Supported: persistent remote text conversations with all ten local provider adapters, follow-up turns, approvals, questions where the provider offers them, cancellation, per-device revocation, reconnect, remote file browsing and text editing, Git status, file diffs, and history, staging and commits, branch selection and creation, worktree selection and creation, and a tracked Git diff against HEAD. OpenCode's server and event stream stay on the host's loopback interface. Cursor's optional enrichment from its native session database is not available on the headless host; basic transcript and subagent events still work. The desktop polls the host and downloads only transcript blocks that changed since its last update (every 0.75 s while a session runs, 3 s otherwise). The desktop rejects any single host response over 16 MiB. A sync above 4 MiB, such as reopening a very long transcript or one very large tool output, is sent as a series of bounded pieces of one consistent revision. Transcript size is therefore not limited by the response cap. The host writes streamed output in 120 ms batches and keeps a bounded event journal. -Remote history appears in the Sessions sidebar of each project on a machine. Host snapshots also populate the app's normal session state while the tab is open; they are not written to the local session store. Remote `/compact` uses the provider's context compaction, and `/plan` selects the host provider's plan mode. Queued follow-ups and editing the last message still need host commands. Other local slash commands and skill expansion are not yet available remotely. Worktree deletion, terminals, generated image output, named provider accounts, `/operator`, automations, orchestration, LAN discovery, and account-based tunnels are not implemented yet. Other remote prompts are sent directly to the provider. +Remote history appears in the Sessions sidebar of each project on a machine. Host snapshots also populate the app's normal session state while the tab is open; they are not written to the local session store. Remote `/compact` uses the provider's context compaction, and `/plan` selects the host provider's plan mode. Editing the last message still needs a host command. Other local slash commands and skill expansion are not yet available remotely. Worktree deletion, terminals, generated image output, named provider accounts, `/operator`, automations, orchestration, LAN discovery, and account-based tunnels are not implemented yet. Other remote prompts are sent directly to the provider. The headless host runs the reused TypeScript adapters with a Node process backend. It proves the execution boundary without introducing the planned Rust daemon/worker IPC yet. Node is included in host release archives, separately from the desktop application. diff --git a/host/engine.test.ts b/host/engine.test.ts index 3f5b58687..a7296f850 100644 --- a/host/engine.test.ts +++ b/host/engine.test.ts @@ -698,6 +698,108 @@ 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("keeps a draft saved while the queue is paused through the queued turns", 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" }); + engine.command({ type: "cancel", commandId: "stop", sessionId: id, runId: store.session(id).runId }); + await vi.waitFor(() => expect(store.session(id).status).toBe("idle")); + engine.command({ type: "draft", commandId: "draft", sessionId: id, text: "Later" }); + + engine.command({ type: "resumeQueue", commandId: "resume", sessionId: id }); + await vi.waitFor(() => expect(turns).toHaveLength(2)); + turns[1].finish(); + await vi.waitFor(() => expect(turns).toHaveLength(3)); + expect(turns[2].input.text).toBe("Next"); + expect(store.session(id).session.blocks).toContainEqual( + expect.objectContaining({ id: "draft", text: "Later", draft: true }), + ); + 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 +814,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 +833,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 +999,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..80263c81f 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,87 @@ 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, + keepDrafts: true, + }, + )); } else { if (value.runId !== command.runId || value.status !== "running") throw new Error( @@ -711,6 +740,130 @@ 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; + /** A queued turn starts on its own, so it leaves the user's draft. */ + keepDrafts?: boolean; + }, + ): { value: HostSession; effect: (saved: HostSession) => void } { + const { blockId, prompt, intent, attachments = [], plan, keepDrafts } = + 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) => keepDrafts || !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, + keepDrafts: true, + }, + ); + 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 +992,8 @@ export class HostEngine { latest, this.closing || active.persistenceFailed ? "interrupted" : "idle", message, + undefined, + active.cancelled, ), { type: "settled", error, cancelled: active.cancelled }, ); @@ -849,6 +1004,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 +1039,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/connections.ts b/src/features/connections/model/connections.ts index b2305d7f9..f86e27431 100644 --- a/src/features/connections/model/connections.ts +++ b/src/features/connections/model/connections.ts @@ -94,6 +94,9 @@ export const pendingRemoteFollowup = (project: string, environment: string, id: return value ? readPendingEntry(value).followup : undefined; }; +export const hasPendingRemoteCommand = (project: string, environment: string, id: string) => + localStorage.getItem(`${pendingPrefix(project, environment)}${id}`) !== null; + export const pendingRemoteCommand = ( project: string, environment: string, 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..f99e93ab8 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,91 @@ 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("returns a queued message the host refused to the composer", async () => { + const original = vi.mocked(invoke).getMockImplementation()!; + vi.mocked(invoke).mockImplementation(async (command, input) => { + const request = input as { method?: string; params?: HostCommand } | undefined; + if (request?.method === "commands.dispatch" && request.params?.type === "enqueue") + throw new Error("Host rejected request: The queue is full"); + return original(command, input); + }); + 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"); + await vi.waitFor(() => + expect(container.querySelector("textarea")!.value).toBe("Second"), + ); + expect(container.textContent).toContain("The queue is full"); +}); + +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..ca85c25af 100644 --- a/src/features/connections/ui/RemoteSession.tsx +++ b/src/features/connections/ui/RemoteSession.tsx @@ -23,6 +23,7 @@ import { useProjectBranchesState } from "../../source-control/hooks/useProjectBr import { registerRemoteSessionActions } from "../model/remoteSessionActions"; import { clearPendingRemoteCommand, + hasPendingRemoteCommand, loadRemoteSession, OPEN_CONNECTIONS_EVENT, pendingRemoteCommand, @@ -323,6 +324,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 +817,64 @@ 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", + onRejected?: () => void, + ): boolean => { + const version = bindingVersion.current; + const commandId = crypto.randomUUID(); + preparingRef.current = true; + void uploadRemoteAttachments(machine.id, attachments) + .then((refs) => { + if (!alive.current || version !== bindingVersion.current) return; + return run({ + type: "enqueue", + commandId, + sessionId: id, + text, + attachments: refs, + intent, + }); + }) + .catch((reason) => { + if (alive.current && version === bindingVersion.current) + setError(String(reason)); + }) + .then((receipt) => { + // An unconfirmed send stays pending for Retry; anything else that + // did not reach the queue goes back to the composer, unless the tab + // now shows another conversation. + if ( + !receipt && + alive.current && + version === bindingVersion.current && + !hasPendingRemoteCommand(project.key, machine.environmentId, commandId) + ) + onRejected?.(); + }) + .finally(() => { + if (version === bindingVersion.current) 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 +912,6 @@ function ConnectedRemoteSession({ sending || preparingRef.current || pending || - busy || (!text.trim() && !attachments.length) ) return false; @@ -860,6 +919,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, options?.onRejected) + ); const turn = optimisticTurn( text, attachments, @@ -1245,11 +1314,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/model/session.ts b/src/features/sessions/model/session.ts index 585682dbc..849b9d428 100644 --- a/src/features/sessions/model/session.ts +++ b/src/features/sessions/model/session.ts @@ -82,6 +82,8 @@ export type ComposerTurnOptions = { resendEdited?: boolean; /** Restore an edited prompt when the resend rejects asynchronously. */ onResendRejected?: (recovery: EditedResendRejection) => void; + /** Restore the prompt when an accepted turn is dropped before it is recorded. */ + onRejected?: () => void; /** Promote an existing unsent transcript block instead of appending a turn. */ draftBlockId?: string; }; diff --git a/src/features/sessions/ui/Composer.tsx b/src/features/sessions/ui/Composer.tsx index aa56fae02..5a90fdedd 100644 --- a/src/features/sessions/ui/Composer.tsx +++ b/src/features/sessions/ui/Composer.tsx @@ -465,14 +465,16 @@ function MessageQueue({ {label} - + {onSteer ? ( + + ) : null}