diff --git a/packages/cli/src/run/noninteractive.ts b/packages/cli/src/run/noninteractive.ts index 8bada9c8d78..102294f64b0 100644 --- a/packages/cli/src/run/noninteractive.ts +++ b/packages/cli/src/run/noninteractive.ts @@ -5,6 +5,7 @@ import type { LocationRef, OpenCodeClient, SessionMessageAssistantTool, + SessionMessageInfo, } from "@opencode-ai/client/promise" import { SessionMessage } from "@opencode-ai/schema/session-message" import { EOL } from "node:os" @@ -75,6 +76,9 @@ export async function runNonInteractivePrompt(input: Input) { const messageID = SessionMessage.ID.create() const starts = new Map() const tools = new Map() + const renderedText = new Map() + const renderedReasoning = new Map() + const renderedTools = new Set() let submitted = false let promoted = false let emittedError = false @@ -82,6 +86,8 @@ export async function runNonInteractivePrompt(input: Input) { let formCancelled = false let interrupted = false let v1InvalidOutput = false + let prePromotionError: { message: string; [key: string]: unknown } | undefined + let finalizing = false let admission: AbortController | undefined let pendingStep: { timestamp: number; part: Record; label: string } | undefined @@ -104,6 +110,17 @@ export async function runNonInteractivePrompt(input: Input) { UI.empty() } + const writeReasoning = (part: { text: string; [key: string]: unknown }, timestamp: number) => { + if (emit("reasoning", timestamp, { part })) return + const text = part.text.trim() + if (!text) return + const line = `Thinking: ${text}` + if (!process.stdout.isTTY) return void process.stdout.write(line + EOL) + UI.empty() + UI.println(`${UI.Style.TEXT_DIM}\u001b[3m${line}\u001b[0m${UI.Style.TEXT_NORMAL}`) + UI.empty() + } + const flushStep = () => { if (!pendingStep) return const value = pendingStep @@ -181,6 +198,7 @@ export async function runNonInteractivePrompt(input: Input) { if (event.type === "session.input.promoted") { if (event.data.inputID === messageID) { promoted = true + prePromotionError = undefined continue } } @@ -191,7 +209,12 @@ export async function runNonInteractivePrompt(input: Input) { ) { return } + if (!promoted && event.type === "session.execution.failed") { + prePromotionError = event.data.error + continue + } if (!promoted) continue + if (finalizing) continue if (event.type === "session.step.started") { const part = { @@ -219,12 +242,16 @@ export async function runNonInteractivePrompt(input: Input) { if (event.type === "session.text.started") { flushStep() - starts.set("text", { id: partID(event.id), timestamp: time }) + starts.set(`text\u0000${contentKey(event.data.assistantMessageID, event.data.ordinal)}`, { + id: partID(event.id), + timestamp: time, + }) continue } if (event.type === "session.text.ended") { - const started = starts.get("text") - starts.delete("text") + const key = contentKey(event.data.assistantMessageID, event.data.ordinal) + const started = starts.get(`text\u0000${key}`) + starts.delete(`text\u0000${key}`) const part = { id: started?.id ?? partID(event.id), sessionID: input.sessionID, @@ -233,18 +260,23 @@ export async function runNonInteractivePrompt(input: Input) { text: event.data.text, time: { start: started?.timestamp ?? time, end: time }, } + renderedText.set(key, event.data.text) writeText(part, time) continue } if (event.type === "session.reasoning.started") { flushStep() - starts.set("reasoning", { id: partID(event.id), timestamp: time }) + starts.set(`reasoning\u0000${contentKey(event.data.assistantMessageID, event.data.ordinal)}`, { + id: partID(event.id), + timestamp: time, + }) continue } if (event.type === "session.reasoning.ended" && input.thinking) { - const started = starts.get("reasoning") - starts.delete("reasoning") + const key = contentKey(event.data.assistantMessageID, event.data.ordinal) + const started = starts.get(`reasoning\u0000${key}`) + starts.delete(`reasoning\u0000${key}`) const part = { id: started?.id ?? partID(event.id), sessionID: input.sessionID, @@ -254,17 +286,8 @@ export async function runNonInteractivePrompt(input: Input) { metadata: event.data.state, time: { start: started?.timestamp ?? time, end: time }, } - if (emit("reasoning", time, { part })) continue - const text = part.text.trim() - if (!text) continue - const line = `Thinking: ${text}` - if (!process.stdout.isTTY) { - process.stdout.write(line + EOL) - continue - } - UI.empty() - UI.println(`${UI.Style.TEXT_DIM}\u001b[3m${line}\u001b[0m${UI.Style.TEXT_NORMAL}`) - UI.empty() + renderedReasoning.set(key, event.data.text) + writeReasoning(part, time) continue } @@ -360,6 +383,7 @@ export async function runNonInteractivePrompt(input: Input) { }, } tools.delete(key) + renderedTools.add(key) if (!emit("tool_use", time, { part })) await input.renderTool(tool) continue } @@ -405,6 +429,7 @@ export async function runNonInteractivePrompt(input: Input) { }, } tools.delete(key) + renderedTools.add(key) if (input.compatibility === "v1" && (permissionRejected || formCancelled)) continue if (!emit("tool_use", time, { part })) { if (toolOutputText(current.tool, current.content).trim()) @@ -480,6 +505,122 @@ export async function runNonInteractivePrompt(input: Input) { } } + const projectedMessages = async () => { + const messages: SessionMessageInfo[] = [] + let cursor: string | undefined + while (true) { + const page = await input.client.message.list( + cursor + ? { sessionID: input.sessionID, limit: 200, cursor } + : { sessionID: input.sessionID, limit: 200, order: "desc" }, + ) + for (const message of page.data) { + if (message.id === messageID) return { found: true, messages: messages.toReversed() } + messages.push(message) + } + cursor = page.cursor.next ?? undefined + if (!cursor) return { found: false, messages: [] } + } + } + + const reconcile = async () => { + const projected = await projectedMessages() + for (const message of projected.messages) { + if (message.type !== "assistant") continue + const timestamp = message.time.completed ?? message.time.created + let textOrdinal = 0 + let reasoningOrdinal = 0 + for (const item of message.content) { + if (item.type === "text") { + const ordinal = textOrdinal++ + const key = contentKey(message.id, ordinal) + const rendered = renderedText.get(key) ?? "" + if (rendered === item.text || !item.text.startsWith(rendered)) continue + const text = item.text.slice(rendered.length) + writeText( + { + id: projectedPartID(message.id, `text-${ordinal}`), + sessionID: input.sessionID, + messageID: message.id, + type: "text", + text, + time: { start: message.time.created, end: timestamp }, + }, + timestamp, + ) + renderedText.set(key, item.text) + continue + } + if (item.type === "reasoning") { + const ordinal = reasoningOrdinal++ + if (!input.thinking) continue + const key = contentKey(message.id, ordinal) + const rendered = renderedReasoning.get(key) ?? "" + if (rendered === item.text || !item.text.startsWith(rendered)) continue + const text = item.text.slice(rendered.length) + const part = { + id: projectedPartID(message.id, `reasoning-${ordinal}`), + sessionID: input.sessionID, + messageID: message.id, + type: "reasoning", + text, + metadata: item.state, + time: { start: message.time.created, end: timestamp }, + } + renderedReasoning.set(key, item.text) + writeReasoning(part, timestamp) + continue + } + + const key = toolKey(message.id, item.id) + if (renderedTools.has(key) || item.state.status === "streaming" || item.state.status === "running") continue + const part: MiniToolPart = { + id: projectedPartID(message.id, `tool-${item.id}`), + sessionID: input.sessionID, + messageID: message.id, + type: "tool", + callID: item.id, + tool: item.name, + state: + item.state.status === "completed" + ? { + status: "completed", + input: item.state.input, + output: toolOutputText(item.name, item.state.content), + title: item.name, + metadata: { structured: item.state.structured, content: item.state.content, result: item.state.result }, + time: { start: item.time.ran ?? item.time.created, end: item.time.completed ?? timestamp }, + } + : { + status: "error", + input: item.state.input, + error: item.state.error.message, + metadata: { structured: item.state.structured, content: item.state.content, result: item.state.result }, + time: { start: item.time.ran ?? item.time.created, end: item.time.completed ?? timestamp }, + }, + } + renderedTools.add(key) + if (emit("tool_use", timestamp, { part })) continue + if (item.state.status === "completed") { + await input.renderTool(item) + continue + } + if (toolOutputText(item.name, item.state.content).trim()) { + await input.renderTool({ ...item, state: { ...item.state, status: "completed" } }) + } + await input.renderToolError(item) + UI.error(item.state.error.message) + } + + if (message.error && !emittedError) { + emittedError = true + process.exitCode = 1 + if (!emit("error", timestamp, { error: message.error })) UI.error(message.error.message) + } + } + return projected.found + } + const interrupt = () => { if (interrupted) process.exit(130) interrupted = true @@ -559,11 +700,27 @@ export async function runNonInteractivePrompt(input: Input) { ? globals.data.filter((form) => form.sessionID === GLOBAL_FORM_SESSION_ID).map(cancelForm) : []), ]) - await completed + if (input.compatibility === "v1") { + await completed + return + } + + const waiting = input.client.session.wait({ sessionID: input.sessionID }) + await Promise.race([waiting, completed.then(() => waiting)]) + finalizing = true + controller.abort() + const found = await reconcile() + if (!found && !interrupted && !permissionRejected && !formCancelled && !emittedError) { + const error = prePromotionError ?? { type: "unknown", message: "Prompt was not promoted" } + emittedError = true + process.exitCode = 1 + if (!emit("error", Date.now(), { error })) UI.error(error.message) + } } finally { process.off("SIGINT", interrupt) controller.abort() - await stream.return?.(undefined).catch(() => {}) + if (input.compatibility === "v1") await stream.return?.(undefined).catch(() => {}) + else void stream.return?.(undefined).catch(() => {}) } } @@ -595,6 +752,14 @@ function toolKey(messageID: string, callID: string) { return `${messageID}\u0000${callID}` } +function contentKey(messageID: string, ordinal: number) { + return `${messageID}\u0000${ordinal}` +} + +function projectedPartID(messageID: string, part: string) { + return `prt_${messageID.replace(/^msg_/, "")}_${part}` +} + function fallbackTool(event: { id: string created: number diff --git a/packages/cli/test/run/noninteractive.test.ts b/packages/cli/test/run/noninteractive.test.ts index f77aa7a6791..f27869a991d 100644 --- a/packages/cli/test/run/noninteractive.test.ts +++ b/packages/cli/test/run/noninteractive.test.ts @@ -1,5 +1,10 @@ import { afterEach, describe, expect, mock, spyOn, test } from "bun:test" -import { OpenCode, type EventSubscribeOutput, type SessionMessageAssistantTool } from "@opencode-ai/client/promise" +import { + OpenCode, + type EventSubscribeOutput, + type SessionMessageAssistantTool, + type SessionMessageInfo, +} from "@opencode-ai/client/promise" import { runNonInteractivePrompt } from "../../src/run/noninteractive" type V2Event = EventSubscribeOutput @@ -162,10 +167,13 @@ async function run(input: { cancel?: (input: { sessionID: string; formID: string }) => Promise renderTool?: (part: SessionMessageAssistantTool) => Promise renderToolError?: (part: SessionMessageAssistantTool) => Promise + messages?: (inputID: string) => SessionMessageInfo[] + wait?: () => Promise }) { const sdk = OpenCode.make({ baseUrl: "https://opencode.test" }) const values: V2Event[] = [{ id: "evt_connected", type: "server.connected", data: {} }] let wake: (() => void) | undefined + const wait = Promise.withResolvers() const stream = (async function* (): AsyncGenerator { while (true) { const value = values.shift() @@ -175,6 +183,7 @@ async function run(input: { }) continue } + if (value.type.startsWith("session.execution.")) setTimeout(wait.resolve, 0) yield value } })() @@ -193,8 +202,19 @@ async function run(input: { }) as never, ) spyOn(sdk.form, "cancel").mockImplementation((request) => (input.cancel?.(request) ?? ok(undefined)) as never) + let promptID = "msg_prompt" + spyOn(sdk.session, "wait").mockImplementation(() => input.wait?.() ?? wait.promise) + spyOn(sdk.message, "list").mockImplementation(() => + ok({ + data: input.messages?.(promptID) ?? [ + { id: promptID, type: "user", text: "hello", time: { created: 1 } }, + ], + cursor: {}, + }), + ) spyOn(sdk.session, "prompt").mockImplementation((request) => { const messageID = request.id ?? "msg_prompt" + promptID = messageID values.push(...input.turn(messageID)) wake?.() wake = undefined @@ -244,6 +264,63 @@ afterEach(() => { }) describe("runNonInteractivePrompt", () => { + test("uses session.wait then reconciles projected output without a terminal event", async () => { + const idle = Promise.withResolvers() + let done = false + const task = capture({ + format: "json", + turn: (messageID) => [prompted(messageID)], + wait: () => idle.promise, + messages: (messageID) => [ + { + id: "msg_assistant", + type: "assistant", + agent: "build", + model: { providerID: "test", id: "test-model" }, + content: [{ type: "text", text: "projected answer" }], + finish: "stop", + time: { created: 2, completed: 3 }, + }, + { id: messageID, type: "user", text: "hello", time: { created: 1 } }, + ], + }).then((output) => { + done = true + return output + }) + + await Bun.sleep(0) + await Bun.sleep(0) + expect(done).toBe(false) + idle.resolve() + const output = await task + expect( + output.stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)), + ).toEqual([expect.objectContaining({ type: "text", part: expect.objectContaining({ text: "projected answer" }) })]) + }) + + test("reports an observed execution failure before prompt promotion", async () => { + const output = await capture({ + format: "json", + turn: () => [executionFailed("instructions unavailable")], + messages: () => [], + }) + + expect( + output.stdout + .split("\n") + .filter(Boolean) + .map((line) => JSON.parse(line)), + ).toEqual([ + expect.objectContaining({ + type: "error", + error: { type: "provider.transport", message: "instructions unavailable" }, + }), + ]) + }) + test("cancels session and global form blockers and exits on pre-promotion interrupt", async () => { const sdk = await run({ pendingForms: [form("frm_pending", "ses_1"), form("frm_pending_global", "global")], @@ -307,6 +384,9 @@ describe("runNonInteractivePrompt", () => { }), ]) expect(output.stderr).toBe("") + const sdk = await run({ compatibility: "v1", turn: (messageID) => [prompted(messageID), settled()] }) + expect(sdk.session.wait).not.toHaveBeenCalled() + expect(sdk.message.list).not.toHaveBeenCalled() }) test("V1 default output flushes step_start before an unrelated execution failure", async () => { diff --git a/packages/tui/src/config/v1/keybind.ts b/packages/tui/src/config/v1/keybind.ts index d549fe07a52..8ad014488e0 100644 --- a/packages/tui/src/config/v1/keybind.ts +++ b/packages/tui/src/config/v1/keybind.ts @@ -96,7 +96,7 @@ export const Definitions = { session_interrupt: keybind("escape", "Interrupt current session"), session_background: keybind("ctrl+b", "Background blocking session tools"), session_compact: keybind("c", "Compact the session"), - session_queued_prompts: keybind("q", "Manage queued prompts"), + session_queued_prompts: keybind("q", "View pending work"), session_child_first: keybind("down,down", "Toggle subagent picker"), session_child_cycle: keybind("right", "Go to next child session"), session_child_cycle_reverse: keybind("left", "Go to previous child session"), diff --git a/packages/tui/src/mini/footer.command.tsx b/packages/tui/src/mini/footer.command.tsx index 9906286ad48..ff75321bf64 100644 --- a/packages/tui/src/mini/footer.command.tsx +++ b/packages/tui/src/mini/footer.command.tsx @@ -465,8 +465,8 @@ export function RunCommandMenuBody(props: { { action: "queued" as const, category: "Agent", - display: "Manage queued prompts", - footer: `${props.queued().length} queued`, + display: "View pending work", + footer: `${props.queued().length} pending`, keywords: props .queued() .map((item) => item.prompt.text) @@ -673,15 +673,13 @@ export function RunQueuedPromptSelectBody(props: { theme: Accessor prompts: Accessor onClose: () => void - onEdit: (prompt: FooterQueuedPrompt) => void | Promise - onDelete: (prompt: FooterQueuedPrompt) => void | Promise onRows?: (rows: number) => void }) { const entries = createMemo(() => props.prompts().map((prompt) => ({ category: "", display: prompt.prompt.text.replaceAll("\n", " "), - footer: "queued · ctrl+e edit · ctrl+d remove", + footer: prompt.delivery, keywords: prompt.prompt.text, prompt, })), @@ -690,29 +688,13 @@ export function RunQueuedPromptSelectBody(props: { entries, limit: SUBAGENT_LIST_ROWS, onClose: props.onClose, - onSelect: (item) => props.onEdit(item.prompt), + onSelect: props.onClose, onRows: props.onRows, - onKey: (event, item) => { - const ctrl = event.ctrl && !event.meta && !event.shift && !event.super - if (item && (event.name === "delete" || (ctrl && event.name === "d"))) { - event.preventDefault() - props.onDelete(item.prompt) - return true - } - - if (item && ctrl && event.name === "e") { - event.preventDefault() - props.onEdit(item.prompt) - return true - } - - return false - }, }) return ( void>() - private queuedRemoves = new Set<(messageID: string) => boolean | Promise>() private closes = new Set<() => void>() // Microtask-coalesced commit queue. Flushed on next microtask or on close/destroy. private queue: StreamCommit[] = [] @@ -226,7 +219,6 @@ export class RunFooter implements FooterApi { const [state, setState] = createSignal({ phase: "idle", status: "", - queue: 0, model: options.modelLabel, usage: "", first: options.first, @@ -328,7 +320,6 @@ export class RunFooter implements FooterApi { onStatus: footer.setStatus, onSubagentSelect: options.onSubagentSelect, onSubagentInterrupt: options.onSubagentInterrupt, - onQueuedRemove: footer.handleQueuedRemove, }) }, }), @@ -355,13 +346,6 @@ export class RunFooter implements FooterApi { } } - public onQueuedRemove(fn: (messageID: string) => boolean | Promise): () => void { - this.queuedRemoves.add(fn) - return () => { - this.queuedRemoves.delete(fn) - } - } - public onClose(fn: () => void): () => void { if (this.isClosed) { fn() @@ -487,7 +471,6 @@ export class RunFooter implements FooterApi { const state = { phase: next.phase ?? prev.phase, status: typeof next.status === "string" ? next.status : prev.status, - queue: typeof next.queue === "number" ? Math.max(0, next.queue) : prev.queue, model: typeof next.model === "string" ? next.model : prev.model, usage: typeof next.usage === "string" ? next.usage : prev.usage, first: typeof next.first === "boolean" ? next.first : prev.first, @@ -665,11 +648,6 @@ export class RunFooter implements FooterApi { this.requestExitHandler = fn } - private handleQueuedRemove = async (messageID: string): Promise => { - const fn = [...this.queuedRemoves][0] - return fn ? await fn(messageID) : false - } - private handleInputClear = (): void => { this.clearInterruptTimer() this.clearExitTimer() @@ -1080,7 +1058,6 @@ export class RunFooter implements FooterApi { for (const timeout of this.themeRefreshTimeouts) clearTimeout(timeout) this.themeRefreshTimeouts.length = 0 this.prompts.clear() - this.queuedRemoves.clear() this.closes.clear() this.scrollback.destroy() for (const theme of [...this.themes]) this.destroyTheme(theme) diff --git a/packages/tui/src/mini/footer.view.tsx b/packages/tui/src/mini/footer.view.tsx index 3e25a07d406..811ddf838fe 100644 --- a/packages/tui/src/mini/footer.view.tsx +++ b/packages/tui/src/mini/footer.view.tsx @@ -102,7 +102,6 @@ type RunFooterViewProps = { onStatus: (text: string) => void onSubagentSelect?: (sessionID: string | undefined) => void onSubagentInterrupt?: (sessionID: string) => void - onQueuedRemove: (messageID: string) => Promise } export function RunFooterView(props: RunFooterViewProps) { @@ -179,7 +178,6 @@ export function RunFooterView(props: RunFooterViewProps) { const busy = createMemo(() => props.state().phase === "running") const armed = createMemo(() => props.state().interrupt > 0) const exiting = createMemo(() => props.state().exit > 0) - const queue = createMemo(() => props.state().queue) const usage = createMemo(() => props.state().usage) const interruptLabel = createMemo(() => { if (!interrupt()) { @@ -414,7 +412,7 @@ export function RunFooterView(props: RunFooterViewProps) { items.push({ kind: "background", key: backgroundShortcut(), label: "background" }) } if (queuedPrompts().length > 0 && queuedShortcut()) { - items.push({ kind: "queued", key: queuedShortcut(), label: `${queue()} queued` }) + items.push({ kind: "queued", key: queuedShortcut(), label: `${queuedPrompts().length} pending` }) } if (activeTabs().length > 0 && subagentShortcut()) { items.push({ kind: "subagents", key: subagentShortcut(), label: "subagents" }) @@ -495,7 +493,7 @@ export function RunFooterView(props: RunFooterViewProps) { commands: [ { id: "session.queued_prompts", - title: "Manage queued prompts", + title: "View pending work", group: "Session", run: openQueuedMenu, }, @@ -656,12 +654,6 @@ export function RunFooterView(props: RunFooterViewProps) { theme={theme} prompts={queuedPrompts} onClose={closePanel} - onDelete={(item) => void props.onQueuedRemove(item.messageID)} - onEdit={async (item) => { - if (!(await props.onQueuedRemove(item.messageID))) return - closePanel() - queueMicrotask(() => composer.replacePrompt(item.prompt)) - }} onRows={setSubagentMenuRows} /> diff --git a/packages/tui/src/mini/runtime.queue.ts b/packages/tui/src/mini/runtime.queue.ts index cc58f1ebf96..8141f135ea3 100644 --- a/packages/tui/src/mini/runtime.queue.ts +++ b/packages/tui/src/mini/runtime.queue.ts @@ -1,8 +1,8 @@ // Serial prompt queue for direct interactive mode. // -// Prompts arrive from the footer (user types and hits enter) and queue up -// here. The queue drains one turn at a time; ordinary prompts waiting behind -// an active ordinary turn are exposed for edit/removal until they begin. +// Prompts arrive from the footer (user types and hits enter) and local +// operations drain one at a time. Ordinary prompts submitted during an active +// ordinary turn are admitted immediately to the server's durable queue. // // The queue also handles /exit, /quit, and /new commands, empty-prompt rejection, // and tracks per-turn wall-clock duration for the footer status line. @@ -11,84 +11,54 @@ import { SessionMessage } from "@opencode-ai/schema/session-message" import { Locale } from "../util/locale" import { isExitCommand, isNewCommand } from "./prompt.shared" -import type { FooterApi, FooterEvent, FooterQueuedPrompt, RunPrompt } from "./types" +import type { FooterApi, FooterEvent, RunPrompt } from "./types" type Trace = { write(type: string, data?: unknown): void } -type Deferred = { - promise: Promise - resolve: (value: T | PromiseLike) => void - reject: (error?: unknown) => void -} - export type QueueInput = { footer: FooterApi initialInput?: string trace?: Trace - onSend?: (prompt: RunPrompt) => void + onSend?: (prompt: RunPrompt, delivery: "steer" | "queue") => void + onAdmissionError?: (prompt: RunPrompt, error: unknown) => void | Promise onNewSession?: () => void | Promise - run: (prompt: RunPrompt, signal: AbortSignal) => Promise + admit: (prompt: RunPrompt, signal: AbortSignal) => Promise + settle: () => Promise + run: (prompt: RunPrompt, signal: AbortSignal, admitted: () => void) => Promise } type State = { queue: RunPrompt[] - queued: FooterQueuedPrompt[] active?: RunPrompt + admission?: Promise ctrl?: AbortController closed: boolean } -function defer(): Deferred { - let resolve!: (value: T | PromiseLike) => void - let reject!: (error?: unknown) => void - const promise = new Promise((next, fail) => { - resolve = next - reject = fail - }) - - return { promise, resolve, reject } -} - // Runs the prompt queue until the footer closes. // // Subscribes to footer prompt events and drains operations through input.run(). -// Ordinary prompts submitted during an ordinary active turn remain local and -// are exposed by the footer for edit/removal until their turn begins. +// Ordinary prompts submitted during an ordinary active turn are admitted as +// durable queued work instead of remaining editable process-local state. export async function runPromptQueue(input: QueueInput): Promise { - const stop = defer<{ type: "closed" }>() - const done = defer() + const stop = Promise.withResolvers<{ type: "closed" }>() + const done = Promise.withResolvers() const state: State = { queue: [], - queued: [], closed: input.footer.isClosed, } let draining: Promise | undefined + let admissions = Promise.resolve() + let admissionVersion = 0 + const admissionController = new AbortController() const emit = (next: FooterEvent, row: Record) => { input.trace?.write("ui.patch", row) input.footer.event(next) } - const syncQueue = () => { - const queue = state.queue.length - emit({ type: "queue", queue }, { queue }) - emit( - { - type: "queued.prompts", - prompts: [...state.queued], - }, - { queued: state.queued.length }, - ) - } - - const removeLocalQueued = (queued: FooterQueuedPrompt) => { - if (!state.queued.includes(queued)) return - state.queued = state.queued.filter((item) => item !== queued) - syncQueue() - } - const finish = () => { if (!state.closed || draining) { return @@ -104,8 +74,8 @@ export async function runPromptQueue(input: QueueInput): Promise { state.closed = true state.queue.length = 0 - state.queued.length = 0 state.ctrl?.abort() + admissionController.abort() stop.resolve({ type: "closed" }) finish() } @@ -123,11 +93,7 @@ export async function runPromptQueue(input: QueueInput): Promise { continue } - const queued = state.queued.find((item) => item.prompt === prompt) - if (queued) removeLocalQueued(queued) - if (prompt.mode !== "shell" && isNewCommand(prompt.text)) { - syncQueue() if (!input.onNewSession) { emit( { @@ -149,13 +115,11 @@ export async function runPromptQueue(input: QueueInput): Promise { patch: { phase: "running", status: "starting new session", - queue: state.queue.length, }, }, { phase: "running", status: "starting new session", - queue: state.queue.length, }, ) await input.onNewSession() @@ -167,24 +131,23 @@ export async function runPromptQueue(input: QueueInput): Promise { ? prompt : { ...prompt, - messageID: prompt.messageID ?? queued?.messageID ?? SessionMessage.ID.create(), + messageID: prompt.messageID ?? SessionMessage.ID.create(), } state.active = sent emit( - { - type: "turn.send", - queue: state.queue.length, - }, + { type: "turn.send" }, { phase: "running", status: "sending prompt", - queue: state.queue.length, }, ) const start = Date.now() const ctrl = new AbortController() + const admission = Promise.withResolvers() + const version = admissionVersion state.ctrl = ctrl + state.admission = admission.promise try { await input.footer.idle() @@ -203,13 +166,13 @@ export async function runPromptQueue(input: QueueInput): Promise { input.trace?.write("ui.commit", commit) input.footer.append(commit) } - input.onSend?.(sent) + input.onSend?.(sent, "steer") if (state.closed) { break } - const task = input.run(sent, ctrl.signal).then( + const task = input.run(sent, ctrl.signal, admission.resolve).then( () => ({ type: "done" as const }), (error) => ({ type: "error" as const, error }), ) @@ -223,10 +186,21 @@ export async function runPromptQueue(input: QueueInput): Promise { if (next.type === "error") { throw next.error } + if (sent.mode !== "shell" && admissionVersion !== version) { + do { + const current = admissionVersion + await admissions + if (state.closed) break + await input.settle() + if (current === admissionVersion) break + } while (!state.closed) + } } finally { + admission.resolve() if (state.ctrl === ctrl) { state.ctrl = undefined } + if (state.admission === admission.promise) state.admission = undefined if (sent.mode !== "shell") { const duration = Locale.duration(Math.max(0, Date.now() - start)) @@ -249,14 +223,10 @@ export async function runPromptQueue(input: QueueInput): Promise { } finally { draining = undefined emit( - { - type: "turn.idle", - queue: state.queue.length, - }, + { type: "turn.idle" }, { phase: "idle", status: "", - queue: state.queue.length, }, ) } @@ -279,23 +249,22 @@ export async function runPromptQueue(input: QueueInput): Promise { if ( active && active.mode !== "shell" && - !active.command && prompt.mode !== "shell" && - !prompt.command && + prompt.command?.source !== "skill" && !isNewCommand(prompt.text) ) { - const queued: FooterQueuedPrompt = { - messageID: SessionMessage.ID.create(), - prompt, - } - state.queued = [...state.queued, queued] - state.queue.push(prompt) - syncQueue() + const sent = { ...prompt, messageID: SessionMessage.ID.create() } + const admission = state.admission + admissionVersion += 1 + input.onSend?.(sent, "queue") + admissions = admissions + .then(() => admission) + .then(() => input.admit(sent, admissionController.signal)) + .catch((error) => (state.closed ? undefined : input.onAdmissionError?.(sent, error))) return } state.queue.push(prompt) - syncQueue() if (prompt.mode !== "shell" && isNewCommand(prompt.text)) { drain() return @@ -319,14 +288,6 @@ export async function runPromptQueue(input: QueueInput): Promise { const offClose = input.footer.onClose(() => { close() }) - const offRemoveQueued = input.footer.onQueuedRemove((messageID) => { - const queued = state.queued.find((item) => item.messageID === messageID) - if (!queued) return false - state.queue = state.queue.filter((prompt) => prompt !== queued.prompt) - removeLocalQueued(queued) - return true - }) - try { if (state.closed) { return @@ -341,8 +302,8 @@ export async function runPromptQueue(input: QueueInput): Promise { } finally { offPrompt() offClose() - offRemoveQueued() close() await draining?.catch(() => {}) + await admissions } } diff --git a/packages/tui/src/mini/runtime.ts b/packages/tui/src/mini/runtime.ts index 482f9da66fc..444c5b24520 100644 --- a/packages/tui/src/mini/runtime.ts +++ b/packages/tui/src/mini/runtime.ts @@ -780,6 +780,22 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep }, RESIZE_DELAY) }) + const renderPromptError = async (prompt: RunPrompt, error: unknown, signal?: AbortSignal) => { + if (signal?.aborted || footer.isClosed) return + const text = + (await state.stream?.then((item) => item.mod).catch(() => undefined))?.formatUnknownError(error) ?? + (error instanceof Error ? error.message : String(error)) + const commit = { + kind: "error", + text, + phase: "start", + source: "system", + messageID: prompt.messageID, + } as const + rememberLocal(commit) + footer.append(commit) + } + const runQueue = async () => { await firstPaint if (footer.isClosed) return @@ -798,10 +814,10 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep footer, initialInput: input.initialInput, trace: log, - onSend: (prompt) => { + onSend: (prompt, delivery) => { state.shown = true state.history.push(prompt) - if (prompt.mode !== "shell") { + if (prompt.mode !== "shell" && delivery === "steer") { rememberLocal({ kind: "user", text: prompt.text, @@ -811,6 +827,24 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep }) } }, + admit: async (prompt, signal) => { + await state.switching?.catch(() => {}) + const next = await ensureStream() + await next.handle.queuePromptTurn({ + agent: state.agent, + model: state.model, + variant: state.activeVariant, + prompt, + files: input.files, + includeFiles: false, + signal, + }) + }, + onAdmissionError: renderPromptError, + settle: async () => { + const next = await ensureStream() + await next.handle.waitForIdle() + }, onNewSession: createSession ? async () => { try { @@ -856,6 +890,7 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep }, }) footer.event({ type: "stream.view", view: { type: "prompt" } }) + footer.event({ type: "queued.prompts", prompts: [] }) footer.event({ type: "stream.patch", patch: { @@ -891,7 +926,7 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep } } : undefined, - run: async (prompt, signal) => { + run: async (prompt, signal, admitted) => { if (state.demo && (await state.demo.prompt(prompt, signal))) { return } @@ -900,15 +935,18 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep try { const next = await ensureStream() - await next.handle.runPromptTurn({ - agent: state.agent, - model: state.model, - variant: state.activeVariant, - prompt, - files: input.files, - includeFiles, - signal, - }) + await next.handle.runPromptTurn( + { + agent: state.agent, + model: state.model, + variant: state.activeVariant, + prompt, + files: input.files, + includeFiles, + signal, + }, + admitted, + ) if (prompt.messageID) { state.localRows = state.localRows.filter( (row) => row.commit.kind !== "user" || row.commit.messageID !== prompt.messageID, @@ -918,22 +956,7 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep // pending for the next prompt-shaped turn. if (prompt.mode !== "shell" && prompt.command?.source !== "skill") includeFiles = false } catch (error) { - if (signal.aborted || footer.isClosed) { - return - } - - const text = - (await state.stream?.then((item) => item.mod).catch(() => undefined))?.formatUnknownError(error) ?? - (error instanceof Error ? error.message : String(error)) - const commit = { - kind: "error", - text, - phase: "start", - source: "system", - messageID: prompt.messageID, - } as const - rememberLocal(commit) - footer.append(commit) + await renderPromptError(prompt, error, signal) } }, }) diff --git a/packages/tui/src/mini/stream-v2.transport.ts b/packages/tui/src/mini/stream-v2.transport.ts index ae07370de9a..886c983cab7 100644 --- a/packages/tui/src/mini/stream-v2.transport.ts +++ b/packages/tui/src/mini/stream-v2.transport.ts @@ -6,6 +6,7 @@ import type { PermissionV2Request, SessionMessageAssistantTool, SessionMessageInfo, + SessionPendingInfo, } from "@opencode-ai/client/promise" import { Event } from "@opencode-ai/schema/event" import { SessionMessage } from "@opencode-ai/schema/session-message" @@ -20,6 +21,7 @@ import type { LocalReplayRow, MiniPermissionRequest, MiniFormRequest, + FooterQueuedPrompt, RunFilePart, RunInput, RunPrompt, @@ -64,7 +66,9 @@ export type SessionResizeReplayInput = { } export type SessionTransport = { - runPromptTurn(input: SessionTurnInput): Promise + runPromptTurn(input: SessionTurnInput, admitted?: () => void): Promise + queuePromptTurn(input: SessionTurnInput): Promise + waitForIdle(): Promise interruptActiveTurn(): Promise selectSubagent(sessionID: string | undefined): void replayOnResize(input: SessionResizeReplayInput): Promise @@ -74,11 +78,12 @@ export type SessionTransport = { type Wait = { messageID: string + failureMessageID: string promoted: boolean + promotionObserved: boolean interrupted: boolean failureRendered: boolean - resolve: () => void - reject: (error: unknown) => void + terminalError?: Error } // One active session.shell call. The HTTP response is the completion signal; @@ -134,6 +139,8 @@ type State = { rootActive: boolean buffered?: ReplayBuffer errors: Set + pending: Map + admitted: Set } const money = new Intl.NumberFormat("en-US", { style: "currency", currency: "USD" }) @@ -167,6 +174,16 @@ function errorMessage(error: { message?: string; _tag?: string }) { return error.message || error._tag || "Session execution failed" } +function pendingPrompt(item: SessionPendingInfo): FooterQueuedPrompt | undefined { + if (item.type !== "user") return undefined + return { + messageID: item.id, + prompt: { messageID: item.id, text: item.data.text, parts: [] }, + delivery: item.delivery, + admittedSeq: item.admittedSeq, + } +} + function wait(delay: number, signal: AbortSignal) { return new Promise((resolve) => { const timer = setTimeout(done, delay) @@ -369,6 +386,7 @@ export async function createSessionTransport(input: StreamInput): Promise controller.abort(), { once: true }) const state: State = { permissions: [], @@ -389,6 +407,8 @@ export async function createSessionTransport(input: StreamInput): Promise void let readyReject!: (error: unknown) => void @@ -443,6 +463,30 @@ export async function createSessionTransport(input: StreamInput): Promise { + const prompts = [...state.pending.values()].toSorted((left, right) => left.admittedSeq - right.admittedSeq) + input.trace?.write("ui.patch", { pending: prompts.length }) + input.footer.event({ type: "queued.prompts", prompts }) + } + + const mergePending = (item: SessionPendingInfo) => { + const prompt = pendingPrompt(item) + if (!prompt || state.messageIDs.has(prompt.messageID)) return + state.admitted.add(prompt.messageID) + state.pending.set(prompt.messageID, prompt) + syncPending() + } + + const promoteWait = (wait: Wait, observed: boolean, messageID = wait.messageID) => { + const transition = messageID !== wait.failureMessageID || (observed ? !wait.promotionObserved : !wait.promoted) + wait.promoted = true + if (observed) wait.promotionObserved = true + if (!transition) return + wait.failureMessageID = messageID + wait.failureRendered = false + wait.terminalError = undefined + } + const syncBlockers = () => { if (state.closed || controller.signal.aborted || input.footer.isClosed) return const descendant = subagents.snapshot() @@ -528,15 +572,19 @@ export async function createSessionTransport(input: StreamInput): Promise { if (message.type === "user") { const waiting = state.wait?.messageID === message.id - if (waiting && state.wait) state.wait.promoted = true - if (!render || state.messageIDs.has(message.id)) return + const admitted = state.admitted.delete(message.id) + if (state.wait && (admitted || (waiting && state.wait.failureMessageID === message.id))) + promoteWait(state.wait, false, message.id) + if (state.pending.delete(message.id)) syncPending() + if (state.messageIDs.has(message.id)) return state.messageIDs.add(message.id) + if (!render) return if (reuseVisibleWait && waiting) return write([{ kind: "user", source: "system", text: message.text, phase: "start", messageID: message.id }]) return } if (message.type === "skill") { - if (state.wait?.messageID === message.id) state.wait.promoted = true + if (state.wait?.messageID === message.id) promoteWait(state.wait, false) if (!render || state.skillMessages.has(message.id)) { state.skillMessages.add(message.id) return @@ -616,8 +664,10 @@ export async function createSessionTransport(input: StreamInput): Promise + ( + await client.message.list( + { sessionID: input.sessionID, limit: input.replayLimit ?? 200, order: "desc" }, + { signal }, + ) + ).data.toReversed() + + const settleSession = async (client: OpenCodeClient) => { + await client.session.wait({ sessionID: input.sessionID }, { signal: controller.signal }) + for (const message of await projectedMessages(client, controller.signal)) renderMessage(message, true, true) + state.rootActive = false + write([], { phase: "idle", status: blockerStatus(state.view) }) + await input.footer.idle() + } + const resolvePermissionSources = async ( client: OpenCodeClient, permissions: PermissionV2Request[], @@ -672,8 +738,9 @@ export async function createSessionTransport(input: StreamInput): Promise { const client = attempt.client const options = { signal: attempt.signal } - const [messages, permissions, forms, globals, active] = await Promise.all([ - client.message.list({ sessionID: input.sessionID, limit: input.replayLimit ?? 200, order: "desc" }, options), + const [projected, pending, permissions, forms, globals, active] = await Promise.all([ + projectedMessages(client, attempt.signal), + client.session.pending.list({ sessionID: input.sessionID }, options), client.permission.list({ sessionID: input.sessionID }, options), client.form.list({ sessionID: input.sessionID }, options), input.location @@ -687,7 +754,11 @@ export async function createSessionTransport(input: StreamInput): Promise { + const prompt = pendingPrompt(item) + return prompt ? [[prompt.messageID, prompt] as const] : [] + })) + syncPending() state.permissions = permissions pruneToolSources() for (const message of projected) renderMessage(message, next.render, next.reuseVisibleWait) @@ -714,11 +785,6 @@ export async function createSessionTransport(input: StreamInput): Promise { @@ -756,9 +822,37 @@ export async function createSessionTransport(input: StreamInput): Promise Promise }, + client: OpenCodeClient, + send: () => Promise, + onAdmitted?: () => void, ) => { - let resolve!: () => void - let reject!: (error: unknown) => void - const done = new Promise((ok, fail) => { - resolve = ok - reject = fail - }) const active: Wait = { messageID, - promoted: turn.promoted === true, + failureMessageID: messageID, + promoted: false, + promotionObserved: false, interrupted: false, failureRendered: false, - resolve, - reject, } state.wait = active const interrupt = () => { @@ -1277,14 +1358,26 @@ export async function createSessionTransport(input: StreamInput): Promise | undefined + const admitPrompt = async (next: SessionTurnInput, client: OpenCodeClient, delivery: "steer" | "queue") => { + const messageID = next.prompt.messageID + if (!messageID) throw new Error("Prompt message ID is required") + const command = next.prompt.command + const attachments = await prepareAttachments(next, command ? "command" : "prompt", input.readTextFile) + const agents = promptAgents(next) + if (!command) { + input.trace?.write("send.prompt", { sessionID: input.sessionID, messageID, delivery }) + return client.session.prompt( + { + sessionID: input.sessionID, + id: messageID, + text: [next.prompt.text, ...attachments.text].join("\n\n"), + files: attachments.files.length ? attachments.files : undefined, + agents: agents.length ? agents : undefined, + delivery, + }, + { signal: next.signal }, + ) + } + + const selected = await resolveSelectedModel(input, client, next) + if (next.variant && !selected) throw new Error("Cannot select a variant before selecting a model") + input.trace?.write("send.command", { sessionID: input.sessionID, messageID, command: command.name, delivery }) + return client.session.command( + { + sessionID: input.sessionID, + id: messageID, + command: command.name, + arguments: command.arguments, + agent: next.agent, + model: selected, + files: attachments.files.length ? attachments.files : undefined, + agents: agents.length ? agents : undefined, + delivery, + }, + { signal: next.signal }, + ) + } + const replayOnResize = (next: SessionResizeReplayInput) => { queuedResizeReplay = next if (resizeReplay) return resizeReplay @@ -1388,7 +1521,20 @@ export async function createSessionTransport(input: StreamInput): Promise + await runTurnWait( + next, + messageID, + client, + () => client.session.skill( { sessionID: input.sessionID, id: messageID, skill: command.name }, { signal: next.signal }, ), - }) + admitted, + ) return } if (command) { - const selected = await resolveSelectedModel(input, client, next) - if (next.variant && !selected) throw new Error("Cannot select a variant before selecting a model") - // Agent and model ride the command payload; the server switches only - // when the command itself does not pin them. - const attachments = await prepareAttachments(next, "command") - const agents = promptAgents(next) - input.trace?.write("send.command", { sessionID: input.sessionID, messageID, command: command.name }) - await runTurnWait(next, messageID, { - send: () => - client.session.command( - { - sessionID: input.sessionID, - id: messageID, - command: command.name, - arguments: command.arguments, - agent: next.agent, - model: selected, - files: attachments.files.length ? attachments.files : undefined, - agents: agents.length ? agents : undefined, - delivery: "steer", - }, - { signal: next.signal }, - ), - }) + await runTurnWait(next, messageID, client, () => admitPrompt(next, client, "steer"), admitted) return } @@ -1447,23 +1574,7 @@ export async function createSessionTransport(input: StreamInput): Promise - client.session.prompt( - { - sessionID: input.sessionID, - id: messageID, - text: [next.prompt.text, ...attachments.text].join("\n\n"), - files: attachments.files.length ? attachments.files : undefined, - agents: agents.length ? agents : undefined, - delivery: "steer", - }, - { signal: next.signal }, - ), - }) + await runTurnWait(next, messageID, client, () => admitPrompt(next, client, "steer"), admitted) }, async interruptActiveTurn() { // A running shell holds no drain, so session.interrupt cannot reach it; diff --git a/packages/tui/src/mini/types.ts b/packages/tui/src/mini/types.ts index 296fc225a9d..35831d8f999 100644 --- a/packages/tui/src/mini/types.ts +++ b/packages/tui/src/mini/types.ts @@ -84,6 +84,8 @@ export type RunPrompt = { export type FooterQueuedPrompt = { messageID: string prompt: RunPrompt + delivery: "steer" | "queue" + admittedSeq: number } export type RunAgent = { @@ -161,7 +163,6 @@ export type FooterPhase = "idle" | "running" export type FooterState = { phase: FooterPhase status: string - queue: number model: string usage: string first: boolean @@ -338,10 +339,6 @@ export type FooterEvent = variants: string[] current: string | undefined } - | { - type: "queue" - queue: number - } | { type: "queued.prompts" prompts: FooterQueuedPrompt[] @@ -355,14 +352,8 @@ export type FooterEvent = model: string selection: NonNullable } - | { - type: "turn.send" - queue: number - } - | { - type: "turn.idle" - queue: number - } + | { type: "turn.send" } + | { type: "turn.idle" } | { type: "turn.duration" duration: string @@ -438,7 +429,6 @@ export type LocalReplayRow = { export type FooterApi = { readonly isClosed: boolean onPrompt(fn: (input: RunPrompt) => void): () => void - onQueuedRemove(fn: (messageID: string) => boolean | Promise): () => void onClose(fn: () => void): () => void event(next: FooterEvent): void append(commit: StreamCommit): void diff --git a/packages/tui/test/mini/fixture/footer-api.ts b/packages/tui/test/mini/fixture/footer-api.ts index 223a5145b18..bb1a05fb46d 100644 --- a/packages/tui/test/mini/fixture/footer-api.ts +++ b/packages/tui/test/mini/fixture/footer-api.ts @@ -2,7 +2,6 @@ import type { FooterApi, FooterEvent, RunPrompt, StreamCommit } from "../../../s export function createFooterApiFixture(input: { events?: FooterEvent[]; commits?: StreamCommit[] } = {}) { const prompts = new Set<(input: RunPrompt) => void>() - const queuedRemoves = new Set<(messageID: string) => boolean | Promise>() const closes = new Set<() => void>() const events = input.events ?? [] const commits = input.commits ?? [] @@ -17,10 +16,6 @@ export function createFooterApiFixture(input: { events?: FooterEvent[]; commits? prompts.add(fn) return () => prompts.delete(fn) }, - onQueuedRemove(fn) { - queuedRemoves.add(fn) - return () => queuedRemoves.delete(fn) - }, onClose(fn) { if (closed) { fn() @@ -46,7 +41,6 @@ export function createFooterApiFixture(input: { events?: FooterEvent[]; commits? destroy() { api.close() prompts.clear() - queuedRemoves.clear() closes.clear() }, } @@ -60,8 +54,5 @@ export function createFooterApiFixture(input: { events?: FooterEvent[]; commits? const prompt: RunPrompt = mode ? { text, parts: [], mode } : { text, parts: [] } for (const fn of [...prompts]) fn(prompt) }, - removeQueued(messageID: string) { - for (const fn of [...queuedRemoves]) void fn(messageID) - }, } } diff --git a/packages/tui/test/mini/footer-keymap.test.tsx b/packages/tui/test/mini/footer-keymap.test.tsx index 12e4b7590ce..b7904da120b 100644 --- a/packages/tui/test/mini/footer-keymap.test.tsx +++ b/packages/tui/test/mini/footer-keymap.test.tsx @@ -12,7 +12,6 @@ test("down opens subagents from an empty prompt", async () => { const [state] = createSignal({ phase: "idle", status: "", - queue: 0, model: "gpt-5", usage: "", first: false, @@ -70,7 +69,6 @@ test("down opens subagents from an empty prompt", async () => { onRows={() => {}} onLayout={() => {}} onStatus={() => {}} - onQueuedRemove={async () => true} /> ) diff --git a/packages/tui/test/mini/footer.view.test.tsx b/packages/tui/test/mini/footer.view.test.tsx index 7fce03d8797..842653adf98 100644 --- a/packages/tui/test/mini/footer.view.test.tsx +++ b/packages/tui/test/mini/footer.view.test.tsx @@ -92,7 +92,6 @@ function footerState(input: Partial = {}) { return createSignal({ phase: "idle", status: "", - queue: 0, model: "gpt-5", usage: "", first: false, @@ -158,7 +157,6 @@ async function renderFooter( onRows={() => {}} onLayout={() => {}} onStatus={() => {}} - onQueuedRemove={async () => true} /> ) @@ -680,10 +678,15 @@ test("direct subagent panel closes when moving up from the first item", async () } }) -test("direct queued prompt panel renders pending prompt actions", async () => { - const [prompts] = createSignal([{ messageID: "m-1", prompt: { text: "fix the auth test", parts: [] } }]) - const edited: string[] = [] - const deleted: string[] = [] +test("direct pending panel shows durable delivery without edit actions", async () => { + const [prompts] = createSignal([ + { + messageID: "m-1", + prompt: { text: "fix the auth test", parts: [] }, + delivery: "queue" as const, + admittedSeq: 1, + }, + ]) const app = await testRender( () => ( @@ -692,12 +695,6 @@ test("direct queued prompt panel renders pending prompt actions", async () => { theme={() => RUN_THEME_FALLBACK.footer} prompts={prompts} onClose={() => {}} - onEdit={(prompt) => { - edited.push(prompt.messageID) - }} - onDelete={(prompt) => { - deleted.push(prompt.messageID) - }} /> ), @@ -709,16 +706,14 @@ test("direct queued prompt panel renders pending prompt actions", async () => { const frame = app.captureCharFrame() const list = panelMenu(app.renderer.root) - expect(frame).toContain("Queued prompts") + expect(frame).toContain("Pending work") expect(frame).toContain("fix the auth test") - expect(frame).toContain("queued") + expect(frame).toContain("queue") expect(frame).not.toContain("┌") expect(frame).not.toContain("┃") expectPaletteList(list, 0) - app.mockInput.pressKey("e", { ctrl: true }) - app.mockInput.pressKey("DELETE") - expect(edited).toEqual(["m-1"]) - expect(deleted).toEqual(["m-1"]) + expect(frame).not.toContain("edit") + expect(frame).not.toContain("remove") } finally { app.renderer.destroy() } @@ -999,11 +994,10 @@ test.skip("direct footer clears the synthetic skills draft when the panel closes } }) -test("direct footer shows editable prompts and additional queued work while running", async () => { +test("direct footer shows authoritative pending work while running", async () => { const [state] = createSignal({ phase: "running", status: "", - queue: 3, model: "gpt-5", usage: "", first: false, @@ -1036,7 +1030,14 @@ test("direct footer shows editable prompts and additional queued work while runn state={state} view={view} subagent={subagents} - queuedPrompts={() => [{ messageID: "m-queued", prompt: { text: "follow up", parts: [] } }]} + queuedPrompts={() => [ + { + messageID: "m-queued", + prompt: { text: "follow up", parts: [] }, + delivery: "queue", + admittedSeq: 1, + }, + ]} theme={() => RUN_THEME_FALLBACK} tuiConfig={tuiConfig} onSubmit={() => true} @@ -1053,7 +1054,6 @@ test("direct footer shows editable prompts and additional queued work while runn onRows={() => {}} onLayout={() => {}} onStatus={() => {}} - onQueuedRemove={async () => true} /> ) @@ -1088,9 +1088,9 @@ test("direct footer shows editable prompts and additional queued work while runn expect(spinner).toBeDefined() expect(frame).toContain("a-model-name-long-enough-to-force-responsive-truncation") - expect(frame).toContain("3 queued") + expect(frame).toContain("1 pending") expect(frame).toContain("ctrl+b background") - expect(frame).toContain("ctrl+x q 3 queued") + expect(frame).toContain("ctrl+x q 1 pending") expect(frame).toContain("↓ subagents") expect(frame).toContain("ctrl+p cmd") expect(frame).toContain("a-model-name-long-enough-to-force-responsive-truncation") diff --git a/packages/tui/test/mini/runtime.queue.test.ts b/packages/tui/test/mini/runtime.queue.test.ts index 165fd1e8338..78ce2913c5a 100644 --- a/packages/tui/test/mini/runtime.queue.test.ts +++ b/packages/tui/test/mini/runtime.queue.test.ts @@ -1,8 +1,16 @@ import { describe, expect, test } from "bun:test" -import { runPromptQueue } from "../../src/mini/runtime.queue" +import { runPromptQueue as runPromptQueueBase, type QueueInput } from "../../src/mini/runtime.queue" import type { RunPrompt } from "../../src/mini/types" import { createFooterApiFixture } from "./fixture/footer-api" +function runPromptQueue(input: Omit & Partial>) { + return runPromptQueueBase({ + admit: async () => {}, + settle: async () => {}, + ...input, + }) +} + describe("run runtime queue", () => { test("ignores empty prompts", async () => { const ui = createFooterApiFixture() @@ -171,206 +179,90 @@ describe("run runtime queue", () => { ]) }) - test("passes prompts to onSend", async () => { + test("durably admits in-flight follow-ups in submission order", async () => { const ui = createFooterApiFixture() - const seen: string[] = [] - - await runPromptQueue({ - footer: ui.api, - initialInput: " hello ", - onSend: (input) => { - seen.push(input.text) - }, - run: async () => { - ui.api.close() - }, - }) - - expect(seen).toEqual([" hello "]) - }) - - test("appends the user row before the turn starts", async () => { - const ui = createFooterApiFixture() - - await runPromptQueue({ - footer: ui.api, - initialInput: "/fmt bash", - run: async () => { - expect(ui.commits).toEqual([ - { - kind: "user", - text: "/fmt bash", - phase: "start", - source: "system", - messageID: expect.any(String), - }, - ]) - ui.api.close() - }, - }) - }) - - test("runs queued prompts in order", async () => { - const ui = createFooterApiFixture() - const seen: string[] = [] - let wake: (() => void) | undefined - const gate = new Promise((resolve) => { - wake = resolve - }) + const admitted: string[] = [] + const gate = Promise.withResolvers() const task = runPromptQueue({ footer: ui.api, - run: async (input) => { - seen.push(input.text) - if (seen.length === 1) { - await gate - return - } - - ui.api.close() + run: async (input, _signal, onAdmitted) => { + admitted.push(`${input.text}:steer`) + onAdmitted() + await gate.promise }, + admit: async (input) => { + admitted.push(`${input.text}:queue`) + }, + settle: async () => ui.api.close(), }) ui.submit("one") ui.submit("two") - await Promise.resolve() - expect(seen).toEqual(["one"]) - - wake?.() - await task - - expect(seen).toEqual(["one", "two"]) - }) - - test("exposes ordinary in-flight prompts for removal before sending", async () => { - const ui = createFooterApiFixture() - const turns: RunPrompt[] = [] - let wake: (() => void) | undefined - const gate = new Promise((resolve) => { - wake = resolve - }) - - const task = runPromptQueue({ - footer: ui.api, - run: async (input) => { - turns.push(input) - await gate - }, - }) - - ui.submit("one") - ui.submit("two") - await Promise.resolve() - await Promise.resolve() - - expect(turns.map((item) => item.text)).toEqual(["one"]) - expect(turns[0]?.messageID).toEqual(expect.any(String)) + ui.submit("three") + while (admitted.length < 3) await Bun.sleep(0) + expect(admitted).toEqual(["one:steer", "two:queue", "three:queue"]) expect(ui.commits.map((item) => item.text)).toEqual(["one"]) - const first = ui.events.find((item) => item.type === "queued.prompts") - const event = ui.events.findLast((item) => item.type === "queued.prompts") - expect(first?.type === "queued.prompts" ? first.prompts : []).toEqual([]) - expect( - first?.type === "queued.prompts" && event?.type === "queued.prompts" ? first.prompts === event.prompts : true, - ).toBe(false) - expect(ui.events.findLast((item) => item.type === "queue")).toEqual({ type: "queue", queue: 1 }) - expect(event?.type === "queued.prompts" ? event.prompts.map((item) => item.prompt.text) : []).toEqual(["two"]) - if (event?.type === "queued.prompts") ui.removeQueued(event.prompts[0]!.messageID) - await Promise.resolve() - wake?.() - ui.api.close() + gate.resolve() await task - expect(turns.map((item) => item.text)).toEqual(["one"]) }) - test("removing one managed queued prompt preserves the others", async () => { + test("continues durable admission after one fails", async () => { const ui = createFooterApiFixture() - const turns: string[] = [] - let wake: (() => void) | undefined - const gate = new Promise((resolve) => { - wake = resolve - }) - + const admitted: string[] = [] + const errors: string[] = [] + const gate = Promise.withResolvers() const task = runPromptQueue({ footer: ui.api, - run: async (input) => { - turns.push(input.text) - if (input.text === "active") await gate - if (input.text === "queued three") ui.api.close() + run: async (_input, _signal, admitted) => { + admitted() + await gate.promise }, - }) - - ui.submit("active") - ui.submit("queued one") - ui.submit("queued two") - ui.submit("queued three") - await Promise.resolve() - await Promise.resolve() - - const event = ui.events.findLast((item) => item.type === "queued.prompts") - if (event?.type === "queued.prompts") { - const second = event.prompts.find((item) => item.prompt.text === "queued two") - if (second) ui.removeQueued(second.messageID) - } - - wake?.() - await task - expect(turns).toEqual(["active", "queued one", "queued three"]) - }) - - test("drains a prompt queued during an in-flight turn", async () => { - const ui = createFooterApiFixture() - const seen: string[] = [] - let wake: (() => void) | undefined - const gate = new Promise((resolve) => { - wake = resolve - }) - - const task = runPromptQueue({ - footer: ui.api, - run: async (input) => { - seen.push(input.text) - if (seen.length === 1) { - await gate - return - } - - ui.api.close() + admit: async (input) => { + if (input.text === "two") throw new Error("admission failed") + admitted.push(input.text) }, + onAdmissionError: (_prompt, error) => { + errors.push(error instanceof Error ? error.message : String(error)) + }, + settle: async () => ui.api.close(), }) ui.submit("one") - await Promise.resolve() - expect(seen).toEqual(["one"]) - - wake?.() - await Promise.resolve() ui.submit("two") + ui.submit("three") + while (admitted.length === 0) await Bun.sleep(0) + gate.resolve() await task - expect(seen).toEqual(["one", "two"]) + expect(errors).toEqual(["admission failed"]) + expect(admitted).toEqual(["three"]) }) - test("close aborts the active run and drops pending queued work", async () => { + test("close aborts an in-flight durable admission", async () => { const ui = createFooterApiFixture() - const seen: string[] = [] - let hit = false + let admissionHit = false + const admissionStarted = Promise.withResolvers() const task = runPromptQueue({ footer: ui.api, - run: async (input, signal) => { - seen.push(input.text) + run: async (_input, signal, admitted) => { + admitted() + await new Promise((resolve) => signal.addEventListener("abort", () => resolve(), { once: true })) + }, + admit: async (_prompt, signal) => { + admissionStarted.resolve() await new Promise((resolve) => { if (signal.aborted) { - hit = true + admissionHit = true resolve() return } - signal.addEventListener( "abort", () => { - hit = true + admissionHit = true resolve() }, { once: true }, @@ -382,11 +274,11 @@ describe("run runtime queue", () => { ui.submit("one") await Promise.resolve() ui.submit("two") + await admissionStarted.promise ui.api.close() await task - expect(hit).toBe(true) - expect(seen).toEqual(["one"]) + expect(admissionHit).toBe(true) }) test("propagates run errors", async () => { diff --git a/packages/tui/test/mini/runtime.test.ts b/packages/tui/test/mini/runtime.test.ts index 3c6e629ea88..2185aec127a 100644 --- a/packages/tui/test/mini/runtime.test.ts +++ b/packages/tui/test/mini/runtime.test.ts @@ -93,6 +93,8 @@ describe("run interactive runtime", () => { streamStarted.resolve() return { runPromptTurn: async () => {}, + queuePromptTurn: async () => {}, + waitForIdle: async () => {}, interruptActiveTurn: async () => {}, selectSubagent: () => {}, settleForm: (sessionID: string, formID: string) => settled.push({ sessionID, formID }), @@ -432,6 +434,8 @@ describe("run interactive runtime", () => { setTimeout(() => input.footer.close(), 0) return { runPromptTurn: async () => {}, + queuePromptTurn: async () => {}, + waitForIdle: async () => {}, interruptActiveTurn: async () => {}, selectSubagent: () => {}, replayOnResize: async () => false, diff --git a/packages/tui/test/mini/stream-v2.transport.test.ts b/packages/tui/test/mini/stream-v2.transport.test.ts index c7c951a3e11..72dedd94d7d 100644 --- a/packages/tui/test/mini/stream-v2.transport.test.ts +++ b/packages/tui/test/mini/stream-v2.transport.test.ts @@ -127,6 +127,8 @@ function sdk(input: { globals?: FormInfo[] globalLocation?: { directory: string; workspaceID?: string } permissions?: Record + pending?: Record>> + wait?: () => Promise }) { const client = OpenCode.make({ baseUrl: "https://opencode.test" }) let subscription = 0 @@ -159,6 +161,8 @@ function sdk(input: { }), ) spyOn(client.session, "active").mockImplementation(() => ok(input.active?.() ?? {})) + spyOn(client.session.pending, "list").mockImplementation((request) => ok(input.pending?.[request.sessionID] ?? [])) + spyOn(client.session, "wait").mockImplementation(() => input.wait?.() ?? ok(undefined)) spyOn(client.session, "message").mockImplementation((request) => { const message = input.messages?.[request.sessionID]?.find((item) => item.id === request.messageID) return message ? (ok(message) as never) : Promise.reject(new Error(`message not found: ${request.messageID}`)) @@ -417,35 +421,23 @@ describe("V2 mini transport", () => { await transport.close() }) - test("finalizes an idle projection before reducing live output", async () => { + test("waits authoritatively and reconciles the projected terminal suffix", async () => { const events = feed() events.push(connected()) + const settled = defer() + const messages: SessionMessages = [] const client = sdk({ streams: [events], - messages: { - ses_1: [ - { - id: "msg_old", - type: "assistant", - agent: "build", - model: { providerID: "test", id: "model" }, - content: [{ type: "text", text: "[Link](https://example.com)" }], - time: { created: 1 }, - }, - ], - }, + messages: { ses_1: messages }, + wait: () => settled.promise, }) const ui = footer() - const idle = spyOn(ui.api, "idle") const transport = await createSessionTransport({ sdk: client, sessionID: "ses_1", - thinking: true, - replay: true, + thinking: false, footer: ui.api, }) - expect(ui.commits.map((item) => item.text)).toEqual(["[Link](https://example.com)"]) - expect(idle).toHaveBeenCalledTimes(1) let admitted = false spyOn(client.session, "prompt").mockImplementation((request) => { @@ -480,7 +472,7 @@ describe("V2 mini transport", () => { sessionID: "ses_1", assistantMessageID: "msg_assistant", ordinal: 0, - delta: "answer", + delta: "ans", }, }) events.push({ @@ -490,10 +482,222 @@ describe("V2 mini transport", () => { durable: durable("ses_1"), data: { sessionID: "ses_1" }, }) + let done = false + void turn.then(() => { + done = true + }) + await Bun.sleep(0) + expect(done).toBe(false) + messages.push( + { id: "msg_prompt", type: "user", text: "hello", time: { created: 2 } }, + { + id: "msg_assistant", + type: "assistant", + agent: "build", + model: { providerID: "test", id: "model" }, + content: [{ type: "text", text: "answer" }], + time: { created: 3, completed: 4 }, + }, + ) + settled.resolve() await turn - expect(ui.commits.map((item) => item.text)).toEqual(["[Link](https://example.com)", "answer"]) - expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "idle", status: "" } }) + expect(ui.commits.map((item) => item.text)).toEqual(["ans", "wer"]) + await transport.close() + }) + + test("shows durable pending delivery and appends queued input on promotion", async () => { + const events = feed() + events.push(connected()) + const client = sdk({ + streams: [events], + pending: { + ses_1: [ + { + admittedSeq: 1, + id: "msg_queued", + sessionID: "ses_1", + timeCreated: 1, + type: "user", + data: { text: "follow up" }, + delivery: "queue", + }, + ], + }, + }) + const ui = footer() + const transport = await createSessionTransport({ + sdk: client, + sessionID: "ses_1", + thinking: false, + footer: ui.api, + }) + const pending = () => + ui.events + .findLast((item) => item.type === "queued.prompts") + ?.prompts.map((item) => [item.messageID, item.delivery, item.admittedSeq]) + + expect(pending()).toEqual([["msg_queued", "queue", 1]]) + events.push({ + id: "evt_promoted", + created: 2, + type: "session.input.promoted", + durable: durable("ses_1", 2), + data: { sessionID: "ses_1", inputID: "msg_queued" }, + }) + while (!ui.commits.some((item) => item.messageID === "msg_queued")) await Bun.sleep(0) + + expect(ui.commits).toContainEqual( + expect.objectContaining({ kind: "user", messageID: "msg_queued", text: "follow up" }), + ) + expect(pending()).toEqual([]) + const prompt = spyOn(client.session, "prompt").mockImplementation((request) => + ok({ ...promptAdmission(request), admittedSeq: 2 }) as never, + ) + await transport.queuePromptTurn({ + agent: undefined, + model: undefined, + variant: undefined, + prompt: { messageID: "msg_next", text: "another", parts: [] }, + files: [], + includeFiles: false, + }) + expect(prompt).toHaveBeenCalledWith(expect.objectContaining({ delivery: "queue" }), expect.anything()) + events.push({ + id: "evt_earlier_admission", + created: 3, + type: "session.input.admitted", + durable: durable("ses_1", 1), + data: { + sessionID: "ses_1", + inputID: "msg_earlier", + input: { type: "user", data: { text: "earlier" }, delivery: "steer" }, + }, + }) + while (true) { + const pending = ui.events.findLast((item) => item.type === "queued.prompts") + if (pending?.type === "queued.prompts" && pending.prompts.length >= 2) break + await Bun.sleep(0) + } + expect(pending()).toEqual([ + ["msg_earlier", "steer", 1], + ["msg_next", "queue", 2], + ]) + await transport.close() + }) + + test("reports an observed execution failure before prompt promotion", async () => { + const events = feed() + events.push(connected()) + const idle = defer() + const client = sdk({ streams: [events], messages: { ses_1: [] }, wait: () => idle.promise }) + const ui = footer() + const transport = await createSessionTransport({ + sdk: client, + sessionID: "ses_1", + thinking: false, + footer: ui.api, + }) + let admitted = false + spyOn(client.session, "prompt").mockImplementation((request) => { + admitted = true + return ok(promptAdmission(request)) as never + }) + + const turn = transport.runPromptTurn({ + agent: undefined, + model: undefined, + variant: undefined, + prompt: { messageID: "msg_prompt", text: "hello", parts: [] }, + files: [], + includeFiles: false, + }) + while (!admitted) await Bun.sleep(0) + events.push({ + id: "evt_failed", + created: 2, + type: "session.execution.failed", + durable: durable("ses_1", 2), + data: { sessionID: "ses_1", error: { type: "unknown", message: "instructions unavailable" } }, + }) + await Bun.sleep(0) + idle.resolve() + + await turn + expect(ui.commits).toContainEqual( + expect.objectContaining({ kind: "error", messageID: "msg_prompt", text: "instructions unavailable" }), + ) + await transport.close() + }) + + test("attributes an execution-only failure to the latest promoted prompt", async () => { + const events = feed() + events.push(connected()) + const idle = defer() + const messages: SessionMessages = [] + const client = sdk({ streams: [events], messages: { ses_1: messages }, wait: () => idle.promise }) + const ui = footer() + const transport = await createSessionTransport({ + sdk: client, + sessionID: "ses_1", + thinking: false, + footer: ui.api, + }) + let admitted = false + spyOn(client.session, "prompt").mockImplementation((request) => { + admitted = true + return ok(promptAdmission(request)) as never + }) + + const turn = transport.runPromptTurn({ + agent: undefined, + model: undefined, + variant: undefined, + prompt: { messageID: "msg_prompt", text: "hello", parts: [] }, + files: [], + includeFiles: false, + }) + while (!admitted) await Bun.sleep(0) + events.push({ + id: "evt_prompt_promoted", + created: 2, + type: "session.input.promoted", + durable: durable("ses_1", 2), + data: { sessionID: "ses_1", inputID: "msg_prompt" }, + }) + await transport.queuePromptTurn({ + agent: undefined, + model: undefined, + variant: undefined, + prompt: { messageID: "msg_queued", text: "follow up", parts: [] }, + files: [], + includeFiles: false, + }) + events.push({ + id: "evt_queued_promoted", + created: 3, + type: "session.input.promoted", + durable: durable("ses_1", 3), + data: { sessionID: "ses_1", inputID: "msg_queued" }, + }) + events.push({ + id: "evt_failed", + created: 4, + type: "session.execution.failed", + durable: durable("ses_1", 4), + data: { sessionID: "ses_1", error: { type: "unknown", message: "model unavailable" } }, + }) + await Bun.sleep(0) + messages.push( + { id: "msg_prompt", type: "user", text: "hello", time: { created: 2 } }, + { id: "msg_queued", type: "user", text: "follow up", time: { created: 3 } }, + ) + idle.resolve() + + await turn + expect(ui.commits).toContainEqual( + expect.objectContaining({ kind: "error", messageID: "msg_queued", text: "model unavailable" }), + ) await transport.close() }) @@ -786,11 +990,12 @@ describe("V2 mini transport", () => { await transport.close() }) - test("rebootstraps after disconnect and completes a promoted turn from idle active state", async () => { + test("reconnects and hydrates without completing before session.wait", async () => { const first = feed() const second = feed() first.push(connected("evt_connected_1")) second.push(connected("evt_connected_2")) + const idle = defer() let running = true const client = sdk({ streams: [first, second], @@ -799,6 +1004,7 @@ describe("V2 mini transport", () => { if (running) active.ses_1 = { type: "running" } return active }, + wait: () => idle.promise, }) let projected = false spyOn(client.message, "list").mockImplementation(() => @@ -844,11 +1050,25 @@ describe("V2 mini transport", () => { while (!admitted) await Bun.sleep(0) projected = true running = false + second.push({ + id: "evt_prior_failed", + created: 1, + type: "session.execution.failed", + durable: durable("ses_1", 1), + data: { sessionID: "ses_1", error: { type: "unknown", message: "prior execution failed" } }, + }) + second.push({ + id: "evt_prompted", + created: 2, + type: "session.input.promoted", + durable: durable("ses_1", 2), + data: { sessionID: "ses_1", inputID: "msg_prompt" }, + }) first.close() + while (!ui.events.some((event) => event.type === "stream.patch" && event.patch.status === "reconnecting")) + await Bun.sleep(0) + idle.resolve() await turn - - expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "running", status: "reconnecting" } }) - expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "idle", status: "" } }) await transport.close() }) @@ -1789,52 +2009,6 @@ describe("V2 mini transport", () => { await transport.close() }) - test("resolves an interrupted turn even when promotion never arrived", async () => { - const events = feed() - events.push(connected()) - const client = sdk({ - streams: [events], - active: () => ({ ses_1: { type: "running" } }), - }) - const ui = footer() - const transport = await createSessionTransport({ - sdk: client, - sessionID: "ses_1", - thinking: false, - footer: ui.api, - }) - let admitted = false - // The generated method has conditional return types for throwOnError; this mock represents the successful branch. - // @ts-expect-error successful SDK response is valid for both modes at runtime - spyOn(client.session, "prompt").mockImplementation((request) => { - admitted = true - return ok({ data: promptAdmission(request) }) - }) - const interrupted = spyOn(client.session, "interrupt").mockImplementation(() => ok(undefined)) - - const turn = transport.runPromptTurn({ - agent: undefined, - model: undefined, - variant: undefined, - prompt: { messageID: "msg_prompt", text: "hello", parts: [] }, - files: [], - includeFiles: true, - }) - while (!admitted) await Bun.sleep(0) - await transport.interruptActiveTurn() - events.push({ - id: "evt_settled", - created: 0, - type: "session.execution.interrupted", - durable: durable("ses_1"), - data: { sessionID: "ses_1", reason: "user" }, - }) - await turn - - expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1" }) - await transport.close() - }) - test("falls back to the default model when selecting a variant on a fresh session", async () => { const events = feed() events.push(connected()) @@ -1901,7 +2075,8 @@ describe("V2 mini transport", () => { test("interrupts the current Session when an active turn is aborted", async () => { const events = feed() events.push(connected()) - const client = sdk({ streams: [events] }) + const idle = defer() + const client = sdk({ streams: [events], wait: () => idle.promise }) const ui = footer() const transport = await createSessionTransport({ sdk: client, @@ -1940,13 +2115,7 @@ describe("V2 mini transport", () => { }) await Bun.sleep(0) controller.abort() - events.push({ - id: "evt_settled", - created: 0, - type: "session.execution.interrupted", - durable: durable("ses_1"), - data: { sessionID: "ses_1", reason: "user" }, - }) + idle.resolve() await turn expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1" }) @@ -2493,90 +2662,6 @@ describe("V2 mini transport", () => { await transport.close() }) - test("does not resolve a skill turn before the matching activation is observed", async () => { - const events = feed() - events.push(connected()) - const client = sdk({ streams: [events] }) - const ui = footer() - const transport = await createSessionTransport({ - sdk: client, - sessionID: "ses_1", - thinking: false, - footer: ui.api, - }) - let sent = false - spyOn(client.session, "skill").mockImplementation(() => { - sent = true - return ok(undefined) as never - }) - - let done = false - const turn = transport - .runPromptTurn({ - agent: undefined, - model: undefined, - variant: undefined, - prompt: { - messageID: "msg_skill", - text: "/tigerstyle", - parts: [], - command: { name: "tigerstyle", arguments: "", source: "skill" }, - }, - files: [], - includeFiles: true, - }) - .then(() => { - done = true - }) - while (!sent) await Bun.sleep(0) - events.push({ - id: "evt_other", - created: 0, - type: "session.skill.activated", - durable: durable("ses_1"), - data: { - sessionID: "ses_1", - id: "other", - name: "other", - text: "other instructions", - }, - }) - events.push({ - id: "evt_unrelated_settled", - created: 0, - type: "session.execution.succeeded", - durable: durable("ses_1"), - data: { sessionID: "ses_1" }, - }) - await Bun.sleep(0) - await Bun.sleep(0) - expect(done).toBe(false) - - events.push({ - id: "evt_skill", - created: 0, - type: "session.skill.activated", - durable: durable("ses_1"), - data: { - sessionID: "ses_1", - id: "tigerstyle", - name: "tigerstyle", - text: "skill instructions", - }, - }) - events.push({ - id: "evt_skill_settled", - created: 0, - type: "session.execution.succeeded", - durable: durable("ses_1"), - data: { sessionID: "ses_1" }, - }) - await turn - - expect(done).toBe(true) - await transport.close() - }) - test("refreshes catalogs on connection and location-scoped invalidations", async () => { const events = feed() events.push(connected())