From d81ae0f0d41a14646a16d43948c8d3532cd1fe61 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Thu, 6 Aug 2026 21:28:32 -0400 Subject: [PATCH] feat(tui): manage queued prompts --- .../context/server-session-v2-reducer.test.ts | 76 +++++++++++++ .../src/context/server-session-v2-reducer.ts | 5 + packages/client/src/effect/api/api.ts | 86 ++++++++++----- .../client/src/effect/generated/client.ts | 60 +++++++---- .../client/src/promise/generated/client.ts | 26 +++++ .../client/src/promise/generated/types.ts | 38 +++++++ packages/client/test/promise.test.ts | 21 ++++ packages/core/src/session.ts | 44 ++++++++ packages/core/src/session/message-updater.ts | 2 + packages/core/src/session/pending.ts | 69 ++++++++++++ packages/core/src/session/projector.ts | 12 +++ packages/core/test/session-prompt.test.ts | 68 ++++++++++++ packages/protocol/src/groups/session.ts | 30 ++++++ packages/schema/src/session-event.ts | 27 ++++- packages/schema/test/event-manifest.test.ts | 2 + packages/server/src/handlers/session.ts | 52 +++++++++ packages/tui/src/component/prompt/index.tsx | 8 +- packages/tui/src/config/v1/keybind.ts | 2 + packages/tui/src/context/data.tsx | 62 +++++++++-- packages/tui/src/mini/footer.command.tsx | 37 ++++++- packages/tui/src/mini/footer.prompt.tsx | 43 ++++++-- packages/tui/src/mini/footer.ts | 4 + packages/tui/src/mini/footer.view.tsx | 41 +++++-- packages/tui/src/mini/runtime.lifecycle.ts | 4 + packages/tui/src/mini/runtime.ts | 10 ++ packages/tui/src/mini/stream-v2.subagent.ts | 4 + packages/tui/src/mini/stream-v2.transport.ts | 23 ++++ packages/tui/src/routes/session/index.tsx | 42 +++++++- packages/tui/test/cli/tui/data.test.tsx | 86 +++++++++++++++ packages/tui/test/mini/footer.view.test.tsx | 101 +++++++++++++++--- .../tui/test/mini/stream-v2.transport.test.ts | 42 +++++++- 31 files changed, 1024 insertions(+), 103 deletions(-) diff --git a/packages/app/src/context/server-session-v2-reducer.test.ts b/packages/app/src/context/server-session-v2-reducer.test.ts index 96100d6367..1b68e6d272 100644 --- a/packages/app/src/context/server-session-v2-reducer.test.ts +++ b/packages/app/src/context/server-session-v2-reducer.test.ts @@ -153,4 +153,80 @@ describe("v2 session reducer", () => { expect(result).toMatchObject({ sessionID: "ses_1", missing: "msg_user", touched: [] }) }) + + test("removes cancelled input from the pending promotion fold", () => { + const reducer = createV2SessionReducer() + reducer.reduce( + [], + event({ + ...base, + id: "evt_admitted", + type: "session.input.admitted", + data: { + sessionID: "ses_1", + inputID: "msg_user", + input: { type: "user", delivery: "queue", data: { text: "cancel me" } }, + }, + }), + ) + reducer.reduce( + [], + event({ + ...base, + id: "evt_cancelled", + type: "session.input.cancelled", + data: { sessionID: "ses_1", inputID: "msg_user" }, + }), + ) + + const result = reducer.reduce( + [], + event({ + ...base, + id: "evt_promoted", + type: "session.input.promoted", + data: { sessionID: "ses_1", inputID: "msg_user" }, + }), + ) + + expect(result).toMatchObject({ missing: "msg_user" }) + }) + + test("keeps steered input available to the promotion fold", () => { + const reducer = createV2SessionReducer() + reducer.reduce( + [], + event({ + ...base, + id: "evt_admitted", + type: "session.input.admitted", + data: { + sessionID: "ses_1", + inputID: "msg_user", + input: { type: "user", delivery: "queue", data: { text: "steer me" } }, + }, + }), + ) + reducer.reduce( + [], + event({ + ...base, + id: "evt_steered", + type: "session.input.steered", + data: { sessionID: "ses_1", inputID: "msg_user" }, + }), + ) + + const result = reducer.reduce( + [], + event({ + ...base, + id: "evt_promoted", + type: "session.input.promoted", + data: { sessionID: "ses_1", inputID: "msg_user" }, + }), + ) + + expect(result?.messages).toMatchObject([{ id: "msg_user", type: "user", text: "steer me" }]) + }) }) diff --git a/packages/app/src/context/server-session-v2-reducer.ts b/packages/app/src/context/server-session-v2-reducer.ts index 12f792d212..7c47f5e6c5 100644 --- a/packages/app/src/context/server-session-v2-reducer.ts +++ b/packages/app/src/context/server-session-v2-reducer.ts @@ -29,6 +29,11 @@ export function createV2SessionReducer() { case "session.input.admitted": pending.set(key(sessionID, event.data.inputID), event.data.input) return result([...source]) + case "session.input.cancelled": + pending.delete(key(sessionID, event.data.inputID)) + return + case "session.input.steered": + return case "session.input.promoted": { const input = pending.get(key(sessionID, event.data.inputID)) pending.delete(key(sessionID, event.data.inputID)) diff --git a/packages/client/src/effect/api/api.ts b/packages/client/src/effect/api/api.ts index 9af36aaff6..53be1a88f2 100644 --- a/packages/client/src/effect/api/api.ts +++ b/packages/client/src/effect/api/api.ts @@ -251,38 +251,48 @@ export type Endpoint5_21Input = { readonly sessionID: Session.ID } export type Endpoint5_21Output = ReadonlyArray export type SessionPendingListOperation = (input: Endpoint5_21Input) => Effect.Effect -export type Endpoint5_22Input = { readonly sessionID: Session.ID } -export type Endpoint5_22Output = ReadonlyArray -export type SessionInstructionsEntryListOperation = ( +export type Endpoint5_22Input = { readonly sessionID: Session.ID; readonly inputID: SessionMessage.ID } +export type Endpoint5_22Output = void +export type SessionPendingCancelOperation = ( input: Endpoint5_22Input, ) => Effect.Effect -export type Endpoint5_23Input = { +export type Endpoint5_23Input = { readonly sessionID: Session.ID; readonly inputID: SessionMessage.ID } +export type Endpoint5_23Output = void +export type SessionPendingSteerOperation = (input: Endpoint5_23Input) => Effect.Effect + +export type Endpoint5_24Input = { readonly sessionID: Session.ID } +export type Endpoint5_24Output = ReadonlyArray +export type SessionInstructionsEntryListOperation = ( + input: Endpoint5_24Input, +) => Effect.Effect + +export type Endpoint5_25Input = { readonly sessionID: Session.ID readonly key: InstructionEntry.Key readonly value: Schema.Json } -export type Endpoint5_23Output = void +export type Endpoint5_25Output = void export type SessionInstructionsEntryPutOperation = ( - input: Endpoint5_23Input, -) => Effect.Effect + input: Endpoint5_25Input, +) => Effect.Effect -export type Endpoint5_24Input = { readonly sessionID: Session.ID; readonly key: InstructionEntry.Key } -export type Endpoint5_24Output = void +export type Endpoint5_26Input = { readonly sessionID: Session.ID; readonly key: InstructionEntry.Key } +export type Endpoint5_26Output = void export type SessionInstructionsEntryRemoveOperation = ( - input: Endpoint5_24Input, -) => Effect.Effect + input: Endpoint5_26Input, +) => Effect.Effect -export type Endpoint5_25Input = { readonly sessionID: Session.ID; readonly prompt: string } -export type Endpoint5_25Output = { readonly text: string } -export type SessionGenerateOperation = (input: Endpoint5_25Input) => Effect.Effect +export type Endpoint5_27Input = { readonly sessionID: Session.ID; readonly prompt: string } +export type Endpoint5_27Output = { readonly text: string } +export type SessionGenerateOperation = (input: Endpoint5_27Input) => Effect.Effect -export type Endpoint5_26Input = { +export type Endpoint5_28Input = { readonly sessionID: Session.ID readonly after?: Event.Seq | undefined readonly follow?: boolean | undefined } -export type Endpoint5_26Output = +export type Endpoint5_28Output = | ( | { readonly id: Event.ID @@ -392,6 +402,24 @@ export type Endpoint5_26Output = readonly input: SessionPending.Message } } + | { + readonly id: Event.ID + readonly created: DateTime.Utc + readonly metadata?: { readonly [x: string]: unknown } | undefined + readonly type: "session.input.cancelled" + readonly durable: { readonly aggregateID: string; readonly seq: Event.Seq; readonly version: Event.Version } + readonly location?: Location.Ref | undefined + readonly data: { readonly sessionID: Session.ID; readonly inputID: SessionMessage.ID } + } + | { + readonly id: Event.ID + readonly created: DateTime.Utc + readonly metadata?: { readonly [x: string]: unknown } | undefined + readonly type: "session.input.steered" + readonly durable: { readonly aggregateID: string; readonly seq: Event.Seq; readonly version: Event.Version } + readonly location?: Location.Ref | undefined + readonly data: { readonly sessionID: Session.ID; readonly inputID: SessionMessage.ID } + } | { readonly id: Event.ID readonly created: DateTime.Utc @@ -850,19 +878,19 @@ export type Endpoint5_26Output = } ) | EventLog.Synced -export type SessionLogOperation = (input: Endpoint5_26Input) => Stream.Stream +export type SessionLogOperation = (input: Endpoint5_28Input) => Stream.Stream -export type Endpoint5_27Input = { readonly sessionID: Session.ID } -export type Endpoint5_27Output = void -export type SessionInterruptOperation = (input: Endpoint5_27Input) => Effect.Effect +export type Endpoint5_29Input = { readonly sessionID: Session.ID } +export type Endpoint5_29Output = void +export type SessionInterruptOperation = (input: Endpoint5_29Input) => Effect.Effect -export type Endpoint5_28Input = { readonly sessionID: Session.ID } -export type Endpoint5_28Output = void -export type SessionBackgroundOperation = (input: Endpoint5_28Input) => Effect.Effect +export type Endpoint5_30Input = { readonly sessionID: Session.ID } +export type Endpoint5_30Output = void +export type SessionBackgroundOperation = (input: Endpoint5_30Input) => Effect.Effect -export type Endpoint5_29Input = { readonly sessionID: Session.ID; readonly messageID: SessionMessage.ID } -export type Endpoint5_29Output = SessionMessage.Info -export type SessionMessageOperation = (input: Endpoint5_29Input) => Effect.Effect +export type Endpoint5_31Input = { readonly sessionID: Session.ID; readonly messageID: SessionMessage.ID } +export type Endpoint5_31Output = SessionMessage.Info +export type SessionMessageOperation = (input: Endpoint5_31Input) => Effect.Effect export interface SessionApi { readonly list: SessionListOperation @@ -888,7 +916,11 @@ export interface SessionApi { readonly commit: SessionRevertCommitOperation } readonly context: SessionContextOperation - readonly pending: { readonly list: SessionPendingListOperation } + readonly pending: { + readonly list: SessionPendingListOperation + readonly cancel: SessionPendingCancelOperation + readonly steer: SessionPendingSteerOperation + } readonly instructions: { readonly entry: { readonly list: SessionInstructionsEntryListOperation diff --git a/packages/client/src/effect/generated/client.ts b/packages/client/src/effect/generated/client.ts index 55a50933cd..54ea227cb2 100644 --- a/packages/client/src/effect/generated/client.ts +++ b/packages/client/src/effect/generated/client.ts @@ -76,6 +76,10 @@ import type { Endpoint5_28Output, Endpoint5_29Input, Endpoint5_29Output, + Endpoint5_30Input, + Endpoint5_30Output, + Endpoint5_31Input, + Endpoint5_31Output, Endpoint6_0Input, Endpoint6_0Output, Endpoint7_0Input, @@ -501,37 +505,51 @@ const Endpoint5_21 = (raw: RawClient["server.session"]) => (input: Endpoint5_21I const Endpoint5_22 = (raw: RawClient["server.session"]) => (input: Endpoint5_22Input) => preserveEffect()( + raw["session.pending.cancel"]({ params: { sessionID: input["sessionID"], inputID: input["inputID"] } }).pipe( + Effect.mapError(mapClientError), + ), + ) + +const Endpoint5_23 = (raw: RawClient["server.session"]) => (input: Endpoint5_23Input) => + preserveEffect()( + raw["session.pending.steer"]({ params: { sessionID: input["sessionID"], inputID: input["inputID"] } }).pipe( + Effect.mapError(mapClientError), + ), + ) + +const Endpoint5_24 = (raw: RawClient["server.session"]) => (input: Endpoint5_24Input) => + preserveEffect()( raw["session.instructions.entry.list"]({ params: { sessionID: input["sessionID"] } }).pipe( Effect.mapError(mapClientError), Effect.map((value) => value.data), ), ) -const Endpoint5_23 = (raw: RawClient["server.session"]) => (input: Endpoint5_23Input) => - preserveEffect()( +const Endpoint5_25 = (raw: RawClient["server.session"]) => (input: Endpoint5_25Input) => + preserveEffect()( raw["session.instructions.entry.put"]({ params: { sessionID: input["sessionID"], key: input["key"] }, payload: { value: input["value"] }, }).pipe(Effect.mapError(mapClientError)), ) -const Endpoint5_24 = (raw: RawClient["server.session"]) => (input: Endpoint5_24Input) => - preserveEffect()( +const Endpoint5_26 = (raw: RawClient["server.session"]) => (input: Endpoint5_26Input) => + preserveEffect()( raw["session.instructions.entry.remove"]({ params: { sessionID: input["sessionID"], key: input["key"] } }).pipe( Effect.mapError(mapClientError), ), ) -const Endpoint5_25 = (raw: RawClient["server.session"]) => (input: Endpoint5_25Input) => - preserveEffect()( +const Endpoint5_27 = (raw: RawClient["server.session"]) => (input: Endpoint5_27Input) => + preserveEffect()( raw["session.generate"]({ params: { sessionID: input["sessionID"] }, payload: { prompt: input["prompt"] } }).pipe( Effect.mapError(mapClientError), Effect.map((value) => value.data), ), ) -const Endpoint5_26 = (raw: RawClient["server.session"]) => (input: Endpoint5_26Input) => - preserveStream()( +const Endpoint5_28 = (raw: RawClient["server.session"]) => (input: Endpoint5_28Input) => + preserveStream()( Stream.unwrap( raw["session.log"]({ params: { sessionID: input["sessionID"] }, @@ -543,18 +561,18 @@ const Endpoint5_26 = (raw: RawClient["server.session"]) => (input: Endpoint5_26I ), ) -const Endpoint5_27 = (raw: RawClient["server.session"]) => (input: Endpoint5_27Input) => - preserveEffect()( +const Endpoint5_29 = (raw: RawClient["server.session"]) => (input: Endpoint5_29Input) => + preserveEffect()( raw["session.interrupt"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError)), ) -const Endpoint5_28 = (raw: RawClient["server.session"]) => (input: Endpoint5_28Input) => - preserveEffect()( +const Endpoint5_30 = (raw: RawClient["server.session"]) => (input: Endpoint5_30Input) => + preserveEffect()( raw["session.background"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError)), ) -const Endpoint5_29 = (raw: RawClient["server.session"]) => (input: Endpoint5_29Input) => - preserveEffect()( +const Endpoint5_31 = (raw: RawClient["server.session"]) => (input: Endpoint5_31Input) => + preserveEffect()( raw["session.message"]({ params: { sessionID: input["sessionID"], messageID: input["messageID"] } }).pipe( Effect.mapError(mapClientError), Effect.map((value) => value.data), @@ -581,13 +599,13 @@ const adaptGroup5 = (raw: RawClient["server.session"]) => ({ wait: Endpoint5_16(raw), revert: { stage: Endpoint5_17(raw), clear: Endpoint5_18(raw), commit: Endpoint5_19(raw) }, context: Endpoint5_20(raw), - pending: { list: Endpoint5_21(raw) }, - instructions: { entry: { list: Endpoint5_22(raw), put: Endpoint5_23(raw), remove: Endpoint5_24(raw) } }, - generate: Endpoint5_25(raw), - log: Endpoint5_26(raw), - interrupt: Endpoint5_27(raw), - background: Endpoint5_28(raw), - message: Endpoint5_29(raw), + pending: { list: Endpoint5_21(raw), cancel: Endpoint5_22(raw), steer: Endpoint5_23(raw) }, + instructions: { entry: { list: Endpoint5_24(raw), put: Endpoint5_25(raw), remove: Endpoint5_26(raw) } }, + generate: Endpoint5_27(raw), + log: Endpoint5_28(raw), + interrupt: Endpoint5_29(raw), + background: Endpoint5_30(raw), + message: Endpoint5_31(raw), }) const Endpoint6_0 = (raw: RawClient["server.message"]) => (input: Endpoint6_0Input) => diff --git a/packages/client/src/promise/generated/client.ts b/packages/client/src/promise/generated/client.ts index 16147d51e7..a9fa625327 100644 --- a/packages/client/src/promise/generated/client.ts +++ b/packages/client/src/promise/generated/client.ts @@ -54,6 +54,10 @@ import type { SessionContextOutput, SessionPendingListInput, SessionPendingListOutput, + SessionPendingCancelInput, + SessionPendingCancelOutput, + SessionPendingSteerInput, + SessionPendingSteerOutput, SessionInstructionsEntryListInput, SessionInstructionsEntryListOutput, SessionInstructionsEntryPutInput, @@ -738,6 +742,28 @@ export function make(options: ClientOptions) { }, requestOptions, ).then((value) => value.data), + cancel: (input: SessionPendingCancelInput, requestOptions?: RequestOptions) => + request( + { + method: "DELETE", + path: `/api/session/${encodeURIComponent(input.sessionID)}/pending/${encodeURIComponent(input.inputID)}`, + successStatus: 204, + declaredStatuses: [409, 404, 400, 401], + empty: true, + }, + requestOptions, + ), + steer: (input: SessionPendingSteerInput, requestOptions?: RequestOptions) => + request( + { + method: "POST", + path: `/api/session/${encodeURIComponent(input.sessionID)}/pending/${encodeURIComponent(input.inputID)}/steer`, + successStatus: 204, + declaredStatuses: [409, 404, 400, 401], + empty: true, + }, + requestOptions, + ), }, instructions: { entry: { diff --git a/packages/client/src/promise/generated/types.ts b/packages/client/src/promise/generated/types.ts index 6d4046c557..951232f69c 100644 --- a/packages/client/src/promise/generated/types.ts +++ b/packages/client/src/promise/generated/types.ts @@ -500,6 +500,26 @@ export type SessionInputPromoted = { data: { sessionID: string; inputID: string } } +export type SessionInputCancelled = { + id: string + created: number + metadata?: { [x: string]: any } + type: "session.input.cancelled" + durable: { aggregateID: string; seq: number; version: 1 } + location?: LocationRef + data: { sessionID: string; inputID: string } +} + +export type SessionInputSteered = { + id: string + created: number + metadata?: { [x: string]: any } + type: "session.input.steered" + durable: { aggregateID: string; seq: number; version: 1 } + location?: LocationRef + data: { sessionID: string; inputID: string } +} + export type SessionExecutionStarted = { id: string created: number @@ -1964,6 +1984,8 @@ export type SessionEventDurable = | SessionForked | SessionInputPromoted | SessionInputAdmitted + | SessionInputCancelled + | SessionInputSteered | SessionExecutionStarted | SessionExecutionSucceeded | SessionExecutionFailed @@ -2016,6 +2038,8 @@ export type V2Event = | SessionForked | SessionInputPromoted | SessionInputAdmitted + | SessionInputCancelled + | SessionInputSteered | SessionExecutionStarted | SessionExecutionSucceeded | SessionExecutionFailed @@ -2934,6 +2958,20 @@ export type SessionPendingListInput = { readonly sessionID: { readonly sessionID export type SessionPendingListOutput = { data: Array }["data"] +export type SessionPendingCancelInput = { + readonly sessionID: { readonly sessionID: string; readonly inputID: string }["sessionID"] + readonly inputID: { readonly sessionID: string; readonly inputID: string }["inputID"] +} + +export type SessionPendingCancelOutput = void + +export type SessionPendingSteerInput = { + readonly sessionID: { readonly sessionID: string; readonly inputID: string }["sessionID"] + readonly inputID: { readonly sessionID: string; readonly inputID: string }["inputID"] +} + +export type SessionPendingSteerOutput = void + export type SessionInstructionsEntryListInput = { readonly sessionID: { readonly sessionID: string }["sessionID"] } export type SessionInstructionsEntryListOutput = { data: Array }["data"] diff --git a/packages/client/test/promise.test.ts b/packages/client/test/promise.test.ts index 14a495cb54..41ac107c5c 100644 --- a/packages/client/test/promise.test.ts +++ b/packages/client/test/promise.test.ts @@ -32,6 +32,7 @@ test("exposes every standard HTTP API group", () => { "projectCopy", "vcs", "debug", + "migration", "websearch", "config", ]) @@ -356,6 +357,26 @@ test("session.pending.list uses the public HTTP contract", async () => { expect(requests).toEqual([{ method: "GET", url: "http://localhost:3000/api/session/ses_test/pending" }]) }) +test("session.pending mutations use the public HTTP contract", async () => { + const requests: Array<{ method: string; url: string }> = [] + const client = OpenCode.make({ + baseUrl: "http://localhost:3000", + fetch: async (input, init) => { + const request = input instanceof Request ? input : new Request(input, init) + requests.push({ method: request.method, url: request.url }) + return new Response(null, { status: 204 }) + }, + }) + + await client.session.pending.cancel({ sessionID: "ses_test", inputID: "msg_cancel" }) + await client.session.pending.steer({ sessionID: "ses_test", inputID: "msg_steer" }) + + expect(requests).toEqual([ + { method: "DELETE", url: "http://localhost:3000/api/session/ses_test/pending/msg_cancel" }, + { method: "POST", url: "http://localhost:3000/api/session/ses_test/pending/msg_steer/steer" }, + ]) +}) + test("event.subscribe exposes the Promise event stream wire projection", async () => { const client = OpenCode.make({ baseUrl: "http://localhost:3000", diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index fc6fc41e0e..61b75dc645 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -133,6 +133,13 @@ export class CompactionConflictError extends Schema.TaggedErrorClass()("Session.BusyError", { sessionID: SessionSchema.ID, }) {} +export class PendingInputConflictError extends Schema.TaggedErrorClass()( + "Session.PendingInputConflictError", + { + sessionID: SessionSchema.ID, + inputID: SessionMessage.ID, + }, +) {} export class SkillNotFoundError extends Schema.TaggedErrorClass()("Session.SkillNotFoundError", { skill: Skill.ID, }) {} @@ -181,6 +188,14 @@ export interface Interface { * unhandled compaction barriers. */ readonly pending: (sessionID: SessionSchema.ID) => Effect.Effect + readonly cancelPending: (input: { + sessionID: SessionSchema.ID + inputID: SessionMessage.ID + }) => Effect.Effect + readonly steerPending: (input: { + sessionID: SessionSchema.ID + inputID: SessionMessage.ID + }) => Effect.Effect /** * Durable, ordered session log read. Replays durable session bus after * the exclusive `after` cursor, emits a `Synced` marker at the captured @@ -507,6 +522,35 @@ const layer = Layer.effect( yield* result.get(sessionID) return yield* SessionPending.list(db, sessionID) }), + cancelPending: Effect.fn("Session.cancelPending")((input) => + Effect.uninterruptible( + Effect.gen(function* () { + yield* result.get(input.sessionID) + yield* SessionPending.cancel(bus, { sessionID: input.sessionID, id: input.inputID }).pipe( + Effect.catchDefect((defect) => + defect instanceof SessionPending.LifecycleConflict + ? new PendingInputConflictError(input) + : Effect.die(defect), + ), + ) + }), + ), + ), + steerPending: Effect.fn("Session.steerPending")((input) => + Effect.uninterruptible( + Effect.gen(function* () { + yield* result.get(input.sessionID) + yield* SessionPending.steer(bus, { sessionID: input.sessionID, id: input.inputID }).pipe( + Effect.catchDefect((defect) => + defect instanceof SessionPending.LifecycleConflict + ? new PendingInputConflictError(input) + : Effect.die(defect), + ), + ) + yield* execution.wake(input.sessionID) + }), + ), + ), log: (input) => Stream.unwrap( result diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index 51ff426bc4..5cafa6f895 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -90,6 +90,8 @@ export function update(adapter: Adapter, event: SessionEvent.DurableEvent) { "session.forked": () => Effect.void, "session.input.promoted": () => Effect.void, "session.input.admitted": () => Effect.void, + "session.input.cancelled": () => Effect.void, + "session.input.steered": () => Effect.void, "session.execution.started": () => Effect.void, "session.execution.succeeded": () => clearCurrentRetry, "session.execution.failed": () => clearCurrentRetry, diff --git a/packages/core/src/session/pending.ts b/packages/core/src/session/pending.ts index 2b9bdac8be..6310cbae20 100644 --- a/packages/core/src/session/pending.ts +++ b/packages/core/src/session/pending.ts @@ -312,6 +312,51 @@ export const projectPromoted = Effect.fn("SessionPending.projectPromoted")(funct return stored }) +export const projectCancelled = Effect.fn("SessionPending.projectCancelled")(function* ( + db: DatabaseService, + input: { + readonly id: SessionMessage.ID + readonly sessionID: SessionSchema.ID + }, +) { + const deleted = yield* db + .delete(SessionPendingTable) + .where( + and( + eq(SessionPendingTable.id, input.id), + eq(SessionPendingTable.session_id, input.sessionID), + eq(SessionPendingTable.delivery, "queue"), + ), + ) + .returning({ id: SessionPendingTable.id }) + .get() + .pipe(Effect.orDie) + if (!deleted) return yield* Effect.die(new LifecycleConflict({ id: input.id })) +}) + +export const projectSteered = Effect.fn("SessionPending.projectSteered")(function* ( + db: DatabaseService, + input: { + readonly id: SessionMessage.ID + readonly sessionID: SessionSchema.ID + }, +) { + const updated = yield* db + .update(SessionPendingTable) + .set({ delivery: "steer" }) + .where( + and( + eq(SessionPendingTable.id, input.id), + eq(SessionPendingTable.session_id, input.sessionID), + eq(SessionPendingTable.delivery, "queue"), + ), + ) + .returning({ id: SessionPendingTable.id }) + .get() + .pipe(Effect.orDie) + if (!updated) return yield* Effect.die(new LifecycleConflict({ id: input.id })) +}) + export const settleCompaction = Effect.fn("SessionPending.settleCompaction")(function* ( db: DatabaseService, input: { readonly sessionID: SessionSchema.ID }, @@ -389,6 +434,30 @@ export const equivalent = ( return false } +export const cancel = Effect.fn("SessionPending.cancel")(function* ( + bus: Bus.Interface, + input: { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID }, +) { + yield* inboxLocks.withLock(input.sessionID)( + bus.publish(SessionEvent.InputCancelled, { + sessionID: input.sessionID, + inputID: input.id, + }), + ) +}) + +export const steer = Effect.fn("SessionPending.steer")(function* ( + bus: Bus.Interface, + input: { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID }, +) { + yield* inboxLocks.withLock(input.sessionID)( + bus.publish(SessionEvent.InputSteered, { + sessionID: input.sessionID, + inputID: input.id, + }), + ) +}) + const publish = Effect.fn("SessionPending.publish")(function* ( db: DatabaseService, bus: Bus.Interface, diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index 5e12004d3e..d9c51baa2d 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -485,6 +485,18 @@ const layer = Layer.effectDiscard( .pipe(Effect.orDie) }), ) + yield* bus.project(SessionEvent.InputCancelled, (event) => + SessionPending.projectCancelled(db, { + id: event.data.inputID, + sessionID: event.data.sessionID, + }), + ) + yield* bus.project(SessionEvent.InputSteered, (event) => + SessionPending.projectSteered(db, { + id: event.data.inputID, + sessionID: event.data.sessionID, + }), + ) yield* bus.project(SessionEvent.Compaction.Admitted, (event) => Effect.gen(function* () { if (event.durable === undefined) diff --git a/packages/core/test/session-prompt.test.ts b/packages/core/test/session-prompt.test.ts index 457175c35c..462e00922c 100644 --- a/packages/core/test/session-prompt.test.ts +++ b/packages/core/test/session-prompt.test.ts @@ -1086,4 +1086,72 @@ describe("Session.pending", () => { expect(yield* session.pending(sessionID)).toEqual([]) }), ) + + it.effect("cancels only queued input and allows its ID to be admitted again", () => + Effect.gen(function* () { + yield* setup + const session = yield* Session.Service + const inputID = SessionMessage.ID.make("msg_cancelled_queue") + yield* session.prompt({ + id: inputID, + sessionID, + text: "Queue this", + delivery: "queue", + resume: false, + }) + + yield* session.cancelPending({ sessionID, inputID }) + + expect(yield* session.pending(sessionID)).toEqual([]) + expect(yield* eventCount(Bus.versionedType(SessionEvent.InputCancelled.type, 1))).toBe(1) + expect( + yield* session.cancelPending({ sessionID, inputID }).pipe(Effect.flip), + ).toMatchObject({ _tag: "Session.PendingInputConflictError", sessionID, inputID }) + expect(yield* eventCount(Bus.versionedType(SessionEvent.InputCancelled.type, 1))).toBe(1) + + const retried = yield* session.prompt({ + id: inputID, + sessionID, + text: "Queue this", + delivery: "queue", + resume: false, + }) + expect(retried).toMatchObject({ id: inputID, delivery: "queue" }) + }), + ) + + it.effect("changes only queued input to steer and wakes after the durable mutation", () => + Effect.gen(function* () { + yield* setup + const session = yield* Session.Service + const queued = yield* session.synthetic({ + sessionID, + text: "Steer this", + delivery: "queue", + resume: false, + }) + const alreadySteered = yield* session.prompt({ sessionID, text: "Already steer", resume: false }) + wakeCalls.length = 0 + + yield* session.steerPending({ sessionID, inputID: queued.id }) + + expect(yield* session.pending(sessionID)).toMatchObject([ + { id: queued.id, delivery: "steer" }, + { id: alreadySteered.id, delivery: "steer" }, + ]) + expect(wakeCalls).toEqual([sessionID]) + expect(yield* eventCount(Bus.versionedType(SessionEvent.InputSteered.type, 1))).toBe(1) + + wakeCalls.length = 0 + expect( + yield* session.steerPending({ sessionID, inputID: alreadySteered.id }).pipe(Effect.flip), + ).toMatchObject({ _tag: "Session.PendingInputConflictError", sessionID, inputID: alreadySteered.id }) + expect( + yield* session.cancelPending({ sessionID, inputID: alreadySteered.id }).pipe(Effect.flip), + ).toMatchObject({ _tag: "Session.PendingInputConflictError", sessionID, inputID: alreadySteered.id }) + expect(wakeCalls).toEqual([]) + expect(yield* eventCount(Bus.versionedType(SessionEvent.InputSteered.type, 1))).toBe(1) + expect(yield* eventCount(Bus.versionedType(SessionEvent.InputCancelled.type, 1))).toBe(0) + }), + ) }) diff --git a/packages/protocol/src/groups/session.ts b/packages/protocol/src/groups/session.ts index 3a4b4e99d8..5bf543073e 100644 --- a/packages/protocol/src/groups/session.ts +++ b/packages/protocol/src/groups/session.ts @@ -491,6 +491,36 @@ export const makeSessionGroup = (sessionLo }), ), ) + .add( + HttpApiEndpoint.delete("session.pending.cancel", "/api/session/:sessionID/pending/:inputID", { + params: { sessionID: Session.ID, inputID: SessionMessage.ID }, + success: HttpApiSchema.NoContent, + error: [ConflictError, SessionNotFoundError], + }) + .middleware(sessionLocationMiddleware) + .annotateMerge( + OpenApi.annotations({ + identifier: "v2.session.pending.cancel", + summary: "Cancel queued input", + description: "Cancel an input that is still queued for delivery.", + }), + ), + ) + .add( + HttpApiEndpoint.post("session.pending.steer", "/api/session/:sessionID/pending/:inputID/steer", { + params: { sessionID: Session.ID, inputID: SessionMessage.ID }, + success: HttpApiSchema.NoContent, + error: [ConflictError, SessionNotFoundError], + }) + .middleware(sessionLocationMiddleware) + .annotateMerge( + OpenApi.annotations({ + identifier: "v2.session.pending.steer", + summary: "Steer queued input", + description: "Change a queued input to steer delivery and wake session execution.", + }), + ), + ) .add( HttpApiEndpoint.get("session.instructions.entry.list", "/api/session/:sessionID/instructions/entries", { params: { sessionID: Session.ID }, diff --git a/packages/schema/src/session-event.ts b/packages/schema/src/session-event.ts index a47be31941..8cf8c55697 100644 --- a/packages/schema/src/session-event.ts +++ b/packages/schema/src/session-event.ts @@ -173,6 +173,26 @@ export const InputAdmitted = Event.durable({ }) export type InputAdmitted = typeof InputAdmitted.Type +export const InputCancelled = Event.durable({ + type: "session.input.cancelled", + ...options, + schema: { + sessionID: SessionID, + inputID: SessionMessage.ID, + }, +}) +export type InputCancelled = typeof InputCancelled.Type + +export const InputSteered = Event.durable({ + type: "session.input.steered", + ...options, + schema: { + sessionID: SessionID, + inputID: SessionMessage.ID, + }, +}) +export type InputSteered = typeof InputSteered.Type + export namespace Execution { export const Started = Event.durable({ type: "session.execution.started", ...options, schema: Base }) export type Started = typeof Started.Type @@ -580,6 +600,8 @@ export const Definitions = Event.inventory( Forked, InputPromoted, InputAdmitted, + InputCancelled, + InputSteered, Execution.Started, Execution.Succeeded, Execution.Failed, @@ -621,13 +643,16 @@ export const DurableDefinitions = Event.inventory( ...Definitions.filter((definition) => definition.durability === "durable"), UsageRecorded, ) +export const EphemeralDefinitions = Event.inventory( + ...Definitions.filter((definition) => definition.durability === "ephemeral"), +) export const Durable = Schema.Union(DurableDefinitions, { mode: "oneOf" }) .pipe(Schema.toTaggedUnion("type")) .annotate({ identifier: "Session.Event.Durable" }) export type DurableEvent = typeof Durable.Type -export const All = Schema.Union(Event.inventory(...Definitions, UsageRecorded), { mode: "oneOf" }).pipe( +export const All = Schema.Union([Durable, ...EphemeralDefinitions], { mode: "oneOf" }).pipe( Schema.toTaggedUnion("type"), ) export type Event = typeof All.Type diff --git a/packages/schema/test/event-manifest.test.ts b/packages/schema/test/event-manifest.test.ts index 5f2872f381..4d9c4a69cb 100644 --- a/packages/schema/test/event-manifest.test.ts +++ b/packages/schema/test/event-manifest.test.ts @@ -84,6 +84,8 @@ describe("public event manifest", () => { "session.forked.2", "session.input.promoted.1", "session.input.admitted.1", + "session.input.cancelled.1", + "session.input.steered.1", "session.execution.started.1", "session.execution.succeeded.1", "session.execution.failed.1", diff --git a/packages/server/src/handlers/session.ts b/packages/server/src/handlers/session.ts index c3c1559a70..99b435060e 100644 --- a/packages/server/src/handlers/session.ts +++ b/packages/server/src/handlers/session.ts @@ -609,6 +609,58 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl } }), ) + .handle( + "session.pending.cancel", + Effect.fn(function* (ctx) { + yield* session + .cancelPending({ sessionID: ctx.params.sessionID, inputID: ctx.params.inputID }) + .pipe( + Effect.catchTag( + "Session.NotFoundError", + (error) => + new SessionNotFoundError({ + sessionID: error.sessionID, + message: `Session not found: ${error.sessionID}`, + }), + ), + Effect.catchTag( + "Session.PendingInputConflictError", + (error) => + new ConflictError({ + resource: error.inputID, + message: `Pending input is no longer queued: ${error.inputID}`, + }), + ), + ) + return HttpApiSchema.NoContent.make() + }), + ) + .handle( + "session.pending.steer", + Effect.fn(function* (ctx) { + yield* session + .steerPending({ sessionID: ctx.params.sessionID, inputID: ctx.params.inputID }) + .pipe( + Effect.catchTag( + "Session.NotFoundError", + (error) => + new SessionNotFoundError({ + sessionID: error.sessionID, + message: `Session not found: ${error.sessionID}`, + }), + ), + Effect.catchTag( + "Session.PendingInputConflictError", + (error) => + new ConflictError({ + resource: error.inputID, + message: `Pending input is no longer queued: ${error.inputID}`, + }), + ), + ) + return HttpApiSchema.NoContent.make() + }), + ) .handle( "session.instructions.entry.list", Effect.fn(function* (ctx) { diff --git a/packages/tui/src/component/prompt/index.tsx b/packages/tui/src/component/prompt/index.tsx index a12fd84f0a..30293b7f75 100644 --- a/packages/tui/src/component/prompt/index.tsx +++ b/packages/tui/src/component/prompt/index.tsx @@ -59,6 +59,7 @@ export type PromptProps = { visible?: boolean disabled?: boolean onSubmit?: () => void + onEmptySubmit?: () => boolean | Promise ref?: (ref: PromptRef | undefined) => void hint?: JSX.Element right?: JSX.Element @@ -950,9 +951,12 @@ export function Prompt(props: PromptProps) { if (props.disabled) return false if (move.creating()) return false if (auto()?.visible) return false - if (!store.prompt.text) return false const trimmed = store.prompt.text.trim() - if (delivery === "queue" && (store.mode === "shell" || trimmed === "exit" || trimmed === "quit" || trimmed === ":q")) { + if (!trimmed) return delivery === "steer" ? (await props.onEmptySubmit?.()) === true : false + if ( + delivery === "queue" && + (store.mode === "shell" || trimmed === "exit" || trimmed === "quit" || trimmed === ":q") + ) { toast.show({ message: "This prompt cannot be queued", variant: "warning" }) return false } diff --git a/packages/tui/src/config/v1/keybind.ts b/packages/tui/src/config/v1/keybind.ts index ac048ff666..70655f5d40 100644 --- a/packages/tui/src/config/v1/keybind.ts +++ b/packages/tui/src/config/v1/keybind.ts @@ -104,6 +104,7 @@ export const Definitions = { session_background: keybind("ctrl+b", "Background blocking session tools"), session_compact: keybind("c", "Compact the session"), session_queued_prompts: keybind("q", "View pending work"), + queued_prompt_delete: keybind("ctrl+d", "Delete queued prompt"), session_child_first: keybind("down", "Toggle subagent picker"), session_parent: keybind("up", "Go to parent session"), session_pin_toggle: keybind("ctrl+f", "Pin or unpin session in the session list"), @@ -306,6 +307,7 @@ export const CommandMap = { session_background: "session.background", session_compact: "session.compact", session_queued_prompts: "session.queued_prompts", + queued_prompt_delete: "queued_prompt.delete", session_child_first: "session.child.first", session_parent: "session.parent", session_pin_toggle: "session.pin.toggle", diff --git a/packages/tui/src/context/data.tsx b/packages/tui/src/context/data.tsx index 92423a9bcb..1d1f995efb 100644 --- a/packages/tui/src/context/data.tsx +++ b/packages/tui/src/context/data.tsx @@ -168,6 +168,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({ function removePending(sessionID: string, inputID?: string) { if (!inputID) return + if (!store.session.pending[sessionID]?.some((item) => item.id === inputID)) return setStore( "session", "pending", @@ -176,6 +177,23 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({ ) } + function updatePending(sessionID: string, inputID: string, delivery: "steer" | "queue") { + if ( + !store.session.pending[sessionID]?.some( + (item) => item.id === inputID && item.type !== "compaction" && item.delivery !== delivery, + ) + ) + return + setStore( + "session", + "pending", + sessionID, + (store.session.pending[sessionID] ?? []).map((item) => + item.id === inputID && item.type !== "compaction" ? { ...item, delivery } : item, + ), + ) + } + const message = { update(sessionID: string, fn: (messages: SessionMessageInfo[], index: Map) => void) { setStore( @@ -222,6 +240,12 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({ (item): item is SessionMessageAssistantReasoning => item.type === "reasoning" && !item.time?.completed, ) }, + reindex(messages: SessionMessageInfo[], index: Map, start: number) { + for (let position = start; position < messages.length; position++) { + const item = messages[position] + if (item) index.set(item.id, position) + } + }, } function index(sessionID: string) { @@ -412,15 +436,37 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({ existing.time.created = event.created draft.splice(position, 1) draft.push(existing) - index.clear() - draft.forEach((message, indexValue) => index.set(message.id, indexValue)) + message.reindex(draft, index, position) }) - setStore( - "session", - "input", - event.data.sessionID, - (store.session.input[event.data.sessionID] ?? []).filter((id) => id !== event.data.inputID), - ) + if (store.session.input[event.data.sessionID]?.includes(event.data.inputID)) + setStore( + "session", + "input", + event.data.sessionID, + (store.session.input[event.data.sessionID] ?? []).filter((id) => id !== event.data.inputID), + ) + break + } + case "session.input.steered": + updatePending(event.data.sessionID, event.data.inputID, "steer") + break + case "session.input.cancelled": { + removePending(event.data.sessionID, event.data.inputID) + if (messageIndex.get(event.data.sessionID)?.has(event.data.inputID)) + message.update(event.data.sessionID, (draft, index) => { + const position = index.get(event.data.inputID) + if (position === undefined) return + draft.splice(position, 1) + index.delete(event.data.inputID) + message.reindex(draft, index, position) + }) + if (store.session.input[event.data.sessionID]?.includes(event.data.inputID)) + setStore( + "session", + "input", + event.data.sessionID, + (store.session.input[event.data.sessionID] ?? []).filter((id) => id !== event.data.inputID), + ) break } case "session.input.admitted": diff --git a/packages/tui/src/mini/footer.command.tsx b/packages/tui/src/mini/footer.command.tsx index 8564fbc082..b5a9825226 100644 --- a/packages/tui/src/mini/footer.command.tsx +++ b/packages/tui/src/mini/footer.command.tsx @@ -3,7 +3,9 @@ import { TextAttributes, type InputRenderable, type KeyEvent } from "@opentui/co import { useKeyboard, type JSX } from "@opentui/solid" import fuzzysort from "fuzzysort" import { createEffect, createMemo, createSignal, type Accessor } from "solid-js" +import { Keymap } from "../context/keymap" import { RunFooterMenu, createFooterMenuState, type RunFooterMenuItem } from "./footer.menu" +import { monoShortcut } from "./mono" import type { RunFooterTheme } from "./theme" import type { FooterQueuedPrompt, @@ -56,6 +58,10 @@ type SkillEntry = PanelEntry & { name: string } +type QueuedPromptEntry = PanelEntry & { + prompt: FooterQueuedPrompt +} + type SubagentEntry = PanelEntry & { sessionID: string current: boolean @@ -837,28 +843,48 @@ export function RunQueuedPromptSelectBody(props: { theme: Accessor prompts: Accessor onClose: () => void + onSteer: (prompt: FooterQueuedPrompt) => void + onDelete: (prompt: FooterQueuedPrompt) => void onRows?: (rows: number) => void mono?: boolean }) { - const entries = createMemo(() => + const entries = createMemo(() => props.prompts().map((prompt) => ({ category: "", display: prompt.prompt.text.replaceAll("\n", " "), - footer: prompt.delivery, + footer: "queued", keywords: prompt.prompt.text, + prompt, })), ) const controller = createSearchablePanelController({ entries, limit: SUBAGENT_LIST_ROWS, onClose: props.onClose, - onSelect: props.onClose, + onSelect: (item) => props.onSteer(item.prompt), onRows: props.onRows, }) + const shortcuts = Keymap.useShortcuts() + const deleteShortcut = () => monoShortcut(shortcuts.get("queued_prompt.delete") ?? "", props.mono ?? false) + Keymap.createLayer(() => ({ + priority: 1, + commands: [ + { + id: "queued_prompt.delete", + title: "Delete queued prompt", + group: "Prompt", + run() { + const item = controller.items()[controller.menu.selected()] + if (!item) return false + props.onDelete(item.prompt) + }, + }, + ], + })) return ( mono: Accessor history?: Accessor + queuedPrompts: Accessor + onQueuedPromptSteer: (inputID: string) => Promise onSubmit: (input: RunPrompt) => boolean | Promise onCycle: () => void onInterrupt: () => boolean @@ -1127,6 +1137,7 @@ export function createPromptState(input: PromptInput): PromptState { } } + let submitting = false const submitPrompt = (next: RunPrompt, delivery: "steer" | "queue" = "steer") => { if (!area || area.isDestroyed) { draft = promptCopy(next) @@ -1141,7 +1152,17 @@ export function createPromptState(input: PromptInput): PromptState { hide() } + if (submitting) return + if (!next.text.trim()) { + const queued = delivery === "steer" ? input.queuedPrompts()[0] : undefined + if (queued) { + submitting = true + void input.onQueuedPromptSteer(queued.messageID).finally(() => { + submitting = false + }) + return + } input.onStatus(input.state().phase === "running" ? "waiting for current response" : "empty prompt ignored") return } @@ -1181,18 +1202,22 @@ export function createPromptState(input: PromptInput): PromptState { : { ...next, delivery } const shellMode = next.mode === "shell" + submitting = true resetDraft() queueMicrotask(async () => { - if (await input.onSubmit(submit)) { - push(next) - if (shellMode) { - setShellMode(false) - draft = emptyPrompt(false) + try { + if (await input.onSubmit(submit)) { + push(next) + if (shellMode) { + setShellMode(false) + draft = emptyPrompt(false) + } + return } - return + restore(next) + } finally { + submitting = false } - - restore(next) }) } diff --git a/packages/tui/src/mini/footer.ts b/packages/tui/src/mini/footer.ts index 99c86dd769..50e028a71c 100644 --- a/packages/tui/src/mini/footer.ts +++ b/packages/tui/src/mini/footer.ts @@ -96,6 +96,8 @@ type RunFooterOptions = { onVariantSelect?: (variant: string | undefined) => CycleResult | void | Promise onInterrupt?: () => void onBackground?: () => void + onQueuedPromptSteer?: (inputID: string) => Promise + onQueuedPromptCancel?: (inputID: string) => Promise onEditorOpen: (input: { value: string }) => Promise onSubagentSelect?: (sessionID: string | undefined) => void onSubagentInterrupt?: (sessionID: string) => void @@ -343,6 +345,8 @@ export class RunFooter implements FooterApi { onCycle: footer.handleCycle, onInterrupt: footer.handleInterrupt, onBackground: options.onBackground, + onQueuedPromptSteer: options.onQueuedPromptSteer, + onQueuedPromptCancel: options.onQueuedPromptCancel, onEditorOpen: options.onEditorOpen, onInputClear: footer.handleInputClear, onExitRequest: footer.handleExit, diff --git a/packages/tui/src/mini/footer.view.tsx b/packages/tui/src/mini/footer.view.tsx index 7afa67a1e8..844177bad7 100644 --- a/packages/tui/src/mini/footer.view.tsx +++ b/packages/tui/src/mini/footer.view.tsx @@ -92,13 +92,15 @@ type RunFooterViewProps = { mono: boolean miniSettings: () => MiniSettings history?: () => RunPrompt[] - onSubmit: (input: RunPrompt) => boolean + onSubmit: (input: RunPrompt) => boolean | Promise onPermissionReply: (input: PermissionReply) => void | Promise onFormReply: (input: FormReply) => void | Promise onFormCancel: (input: FormCancel) => void | Promise onCycle: () => void onInterrupt: () => boolean onBackground?: () => void + onQueuedPromptSteer?: (inputID: string) => Promise + onQueuedPromptCancel?: (inputID: string) => Promise onEditorOpen: (input: { value: string }) => Promise onInputClear: () => void onExitRequest?: () => boolean @@ -132,6 +134,7 @@ export function RunFooterView(props: RunFooterViewProps) { const [route, setRoute] = createSignal({ type: "composer" }) const [subagentMenuRows, setSubagentMenuRows] = createSignal(RUN_SUBAGENT_PANEL_ROWS) const queuedPrompts = createMemo(() => props.queuedPrompts?.() ?? []) + const queue = createMemo(() => queuedPrompts().filter((item) => item.delivery === "queue")) const skills = createMemo(() => (props.commands() ?? []).filter((item) => item.source === "skill")) const prompt = createMemo(() => active().type === "prompt" && route().type === "composer") const selectingSubagent = createMemo(() => active().type === "prompt" && route().type === "subagent-menu") @@ -229,7 +232,7 @@ export function RunFooterView(props: RunFooterViewProps) { const details = [busy() ? "running" : "idle", `agent ${props.currentAgent()}`] if (current) details.push(variant ? `${current} ${variant}` : current) if (usage()) details.push(props.mono ? usage().replaceAll(" · ", " - ") : usage()) - if (queuedPrompts().length > 0) details.push(`${queuedPrompts().length} pending`) + if (queue().length > 0) details.push(`${queue().length} queued`) if (activeTabs().length > 0) details.push(`${activeTabs().length} subagent${activeTabs().length === 1 ? "" : "s"}`) return details.join(props.mono ? " - " : " · ") }) @@ -309,7 +312,7 @@ export function RunFooterView(props: RunFooterViewProps) { } const openQueuedMenu = () => { - if (queuedPrompts().length === 0) return + if (queue().length === 0) return setRoute({ type: "queued-menu" }) props.onSubagentSelect?.(undefined) } @@ -318,6 +321,18 @@ export function RunFooterView(props: RunFooterViewProps) { setRoute({ type: "composer" }) } + const queuedPromptAction = async (action: "steer" | "delete", inputID: string) => { + const run = action === "steer" ? props.onQueuedPromptSteer : props.onQueuedPromptCancel + if (!run) return false + const error = await run(inputID).then( + () => undefined, + (error) => error, + ) + if (!error) return true + props.onStatus(`failed to ${action} queued prompt: ${error instanceof Error ? error.message : String(error)}`) + return false + } + const openTab = (sessionID: string) => { setRoute({ type: "subagent", sessionID }) props.onSubagentSelect?.(sessionID) @@ -357,6 +372,8 @@ export function RunFooterView(props: RunFooterViewProps) { theme, mono: () => props.mono, history: props.history, + queuedPrompts: queue, + onQueuedPromptSteer: (inputID) => queuedPromptAction("steer", inputID), onSubmit: props.onSubmit, onCycle: props.onCycle, onInterrupt: props.onInterrupt, @@ -451,8 +468,8 @@ export function RunFooterView(props: RunFooterViewProps) { if (foregroundSubagents() && backgroundShortcut()) { items.push({ key: backgroundShortcut(), label: "background" }) } - if (queuedPrompts().length > 0 && queuedShortcut()) { - items.push({ key: queuedShortcut(), label: `${queuedPrompts().length} pending` }) + if (queue().length > 0 && queuedShortcut()) { + items.push({ key: queuedShortcut(), label: `${queue().length} queued` }) } if (activeTabs().length > 0 && subagentShortcut()) { items.push({ key: subagentShortcut(), label: "subagents" }) @@ -567,7 +584,7 @@ export function RunFooterView(props: RunFooterViewProps) { })) Keymap.createLayer(() => ({ - enabled: active().type === "prompt" && route().type === "composer" && queuedPrompts().length > 0, + enabled: active().type === "prompt" && route().type === "composer" && queue().length > 0, commands: [ { id: "session.queued_prompts", @@ -629,7 +646,7 @@ export function RunFooterView(props: RunFooterViewProps) { }) createEffect(() => { - if (route().type !== "queued-menu" || queuedPrompts().length > 0) return + if (route().type !== "queued-menu" || queue().length > 0) return closePanel() }) @@ -733,8 +750,16 @@ export function RunFooterView(props: RunFooterViewProps) { { + void queuedPromptAction("steer", item.messageID).then((steered) => { + if (steered) closePanel() + }) + }} + onDelete={(item) => { + void queuedPromptAction("delete", item.messageID) + }} onRows={setSubagentMenuRows} mono={props.mono} /> diff --git a/packages/tui/src/mini/runtime.lifecycle.ts b/packages/tui/src/mini/runtime.lifecycle.ts index d8ef63eb33..b04bb8e15b 100644 --- a/packages/tui/src/mini/runtime.lifecycle.ts +++ b/packages/tui/src/mini/runtime.lifecycle.ts @@ -70,6 +70,8 @@ export type LifecycleInput = { onVariantSelect?: (variant: string | undefined) => CycleResult | void | Promise onInterrupt?: () => void onBackground?: () => void + onQueuedPromptSteer?: (inputID: string) => Promise + onQueuedPromptCancel?: (inputID: string) => Promise onSubagentSelect?: (sessionID: string | undefined) => void onSubagentInterrupt?: (sessionID: string) => void } @@ -243,6 +245,8 @@ export async function createRuntimeLifecycle(input: LifecycleInput): Promise { if (closed || renderer.isDestroyed) { return diff --git a/packages/tui/src/mini/runtime.ts b/packages/tui/src/mini/runtime.ts index dc5f425947..33a1bef639 100644 --- a/packages/tui/src/mini/runtime.ts +++ b/packages/tui/src/mini/runtime.ts @@ -390,6 +390,16 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep log?.write("send.background", { sessionID: state.sessionID }) void state.sdk.session.background({ sessionID: state.sessionID }).catch(() => {}) }, + onQueuedPromptSteer: async (inputID) => { + if (!state.sessionID) return + log?.write("send.pending.steer", { sessionID: state.sessionID, inputID }) + await state.sdk.session.pending.steer({ sessionID: state.sessionID, inputID }) + }, + onQueuedPromptCancel: async (inputID) => { + if (!state.sessionID) return + log?.write("send.pending.cancel", { sessionID: state.sessionID, inputID }) + await state.sdk.session.pending.cancel({ sessionID: state.sessionID, inputID }) + }, onSubagentInterrupt: (sessionID) => { log?.write("send.subagent.interrupt", { sessionID }) void state.sdk.session.interrupt({ sessionID }).catch(() => {}) diff --git a/packages/tui/src/mini/stream-v2.subagent.ts b/packages/tui/src/mini/stream-v2.subagent.ts index d1ff8556d0..2dc2c76d8b 100644 --- a/packages/tui/src/mini/stream-v2.subagent.ts +++ b/packages/tui/src/mini/stream-v2.subagent.ts @@ -653,6 +653,10 @@ export function createSubagentTracker(input: SubagentTrackerInput): SubagentTrac } return } + if (event.type === "session.input.cancelled") { + child.prompts.delete(event.data.inputID) + return + } if (event.type === "session.step.started") { touch(child, event.created) if (child.label === FALLBACK_LABEL && event.data.agent) child.label = Locale.titlecase(event.data.agent) diff --git a/packages/tui/src/mini/stream-v2.transport.ts b/packages/tui/src/mini/stream-v2.transport.ts index a77ac64b1c..906f9766c6 100644 --- a/packages/tui/src/mini/stream-v2.transport.ts +++ b/packages/tui/src/mini/stream-v2.transport.ts @@ -934,6 +934,29 @@ export async function createSessionTransport(input: StreamInput): Promise { + const error = await client.api.session.pending.steer({ sessionID: route.sessionID, inputID }).then( + () => undefined, + (error) => error, + ) + if (!error) return true + toast.show({ title: "Failed to steer queued prompt", message: errorMessage(error), variant: "error" }) + return false + } + const cancelQueuedPrompt = async (inputID: string) => { + const error = await client.api.session.pending.cancel({ sessionID: route.sessionID, inputID }).then( + () => undefined, + (error) => error, + ) + if (!error) return true + toast.show({ title: "Failed to delete queued prompt", message: errorMessage(error), variant: "error" }) + return false + } const openQueuedPrompts = () => dialog.replace(() => ( { + void steerQueuedPrompt(option.value).then((steered) => { + if (steered) dialog.clear() + }) + }} + actions={[ + { + command: "queued_prompt.delete", + title: "delete", + onTrigger: (option) => { + void cancelQueuedPrompt(option.value).then((cancelled) => { + if (cancelled && queuedPrompts().length <= 1) dialog.clear() + }) + }, + }, + ]} + footerHints={[{ title: "steer", label: "enter" }]} /> )) const unavailable = (feature: string) => { @@ -899,7 +934,7 @@ export function Session() { }, { title: "View queued prompts", - id: "session.queue.list", + id: "session.queued_prompts", group: "Session", enabled: queuedPrompts().length > 0, run: openQueuedPrompts, @@ -1068,6 +1103,11 @@ export function Session() { onSubmit={() => { toBottom() }} + onEmptySubmit={async () => { + const next = queuedPrompts()[0] + if (!next) return false + return steerQueuedPrompt(next.id) + }} sessionID={route.sessionID} /> diff --git a/packages/tui/test/cli/tui/data.test.tsx b/packages/tui/test/cli/tui/data.test.tsx index 430775bd18..f5f3b337d4 100644 --- a/packages/tui/test/cli/tui/data.test.tsx +++ b/packages/tui/test/cli/tui/data.test.tsx @@ -914,6 +914,92 @@ test("completes exploration when a queued prompt is promoted", async () => { } }) +test("updates and removes queued inputs from durable lifecycle events", async () => { + const events = createEventStream() + const sessionID = "session-queue-management" + const calls = createFetch((url) => { + if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} }) + }, events) + let data!: ReturnType + let rows!: ReturnType + let client!: ReturnType + + function Probe() { + client = useClient() + data = useData() + rows = createSessionRows(() => sessionID) + return + } + + const app = await testRender(() => ( + + + + + + + + + + )) + + try { + await wait(() => client.connection.status() === "connected") + emitEvent(events, { + id: "evt_queue_admitted", + created: 1, + type: "session.input.admitted", + durable: durable(sessionID), + data: { + sessionID, + inputID: "message-queued", + input: { type: "user", data: { text: "Steer me" }, delivery: "queue" }, + }, + }) + await wait(() => data.session.pending.list(sessionID).length === 1) + expect(rows).not.toContainEqual({ type: "message", messageID: "message-queued" }) + + emitEvent(events, { + id: "evt_queue_steered", + created: 2, + type: "session.input.steered", + durable: durable(sessionID, 1), + data: { sessionID, inputID: "message-queued" }, + }) + await wait(() => + data.session.pending + .list(sessionID) + .some((item) => item.id === "message-queued" && item.type !== "compaction" && item.delivery === "steer"), + ) + expect(rows).toContainEqual({ type: "message", messageID: "message-queued" }) + + emitEvent(events, { + id: "evt_cancel_admitted", + created: 3, + type: "session.input.admitted", + durable: durable(sessionID, 2), + data: { + sessionID, + inputID: "message-cancelled", + input: { type: "user", data: { text: "Delete me" }, delivery: "queue" }, + }, + }) + await wait(() => data.session.pending.list(sessionID).length === 2) + emitEvent(events, { + id: "evt_queue_cancelled", + created: 4, + type: "session.input.cancelled", + durable: durable(sessionID, 3), + data: { sessionID, inputID: "message-cancelled" }, + }) + await wait(() => !data.session.input.has(sessionID, "message-cancelled")) + expect(data.session.pending.list(sessionID).map((item) => item.id)).toEqual(["message-queued"]) + expect(data.session.message.get(sessionID, "message-cancelled")).toBeUndefined() + } finally { + app.renderer.destroy() + } +}) + test("classifies live tool rows independently of their call ID", async () => { const events = createEventStream() const sessionID = "session-tool-call-id" diff --git a/packages/tui/test/mini/footer.view.test.tsx b/packages/tui/test/mini/footer.view.test.tsx index 2dc583d539..b77f3a6e4c 100644 --- a/packages/tui/test/mini/footer.view.test.tsx +++ b/packages/tui/test/mini/footer.view.test.tsx @@ -21,6 +21,7 @@ import { RunFooterView } from "../../src/mini/footer.view" import { RunEntryContent } from "../../src/mini/scrollback.writer" import { RUN_THEME_FALLBACK, type RunTheme } from "../../src/mini/theme" import type { + FooterQueuedPrompt, FooterState, FooterSubagentState, FooterSubagentTab, @@ -120,13 +121,16 @@ async function renderFooter( height?: number state?: Partial onCycle?: () => void - onSubmit?: (prompt: RunPrompt) => boolean + onSubmit?: (prompt: RunPrompt) => boolean | Promise view?: FooterView onFormReply?: (input: unknown) => void miniSettings?: MiniSettings mono?: boolean onStatus?: (status: string) => void onMiniSettingChange?: (change: MiniSettingChange) => void + queuedPrompts?: FooterQueuedPrompt[] + onQueuedPromptSteer?: (inputID: string) => Promise + onQueuedPromptCancel?: (inputID: string) => Promise } = {}, ) { const [view, setView] = createSignal(input.view ?? { type: "prompt" }) @@ -164,6 +168,7 @@ async function renderFooter( state={state} view={view} subagent={subagents} + queuedPrompts={() => input.queuedPrompts ?? []} theme={input.theme ?? (() => RUN_THEME_FALLBACK)} mono={input.mono ?? false} miniSettings={miniSettings} @@ -173,6 +178,8 @@ async function renderFooter( onFormCancel={() => {}} onCycle={input.onCycle ?? (() => {})} onInterrupt={() => false} + onQueuedPromptSteer={input.onQueuedPromptSteer} + onQueuedPromptCancel={input.onQueuedPromptCancel} onEditorOpen={async () => undefined} onInputClear={() => {}} onExit={() => {}} @@ -913,7 +920,7 @@ test("direct subagent panel closes when moving up from the first item", async () } }) -test("direct pending panel shows durable delivery without edit actions", async () => { +test("direct queued panel steers and deletes selected prompts", async () => { const [prompts] = createSignal([ { messageID: "m-1", @@ -921,16 +928,22 @@ test("direct pending panel shows durable delivery without edit actions", async ( delivery: "queue" as const, }, ]) + const steered: string[] = [] + const deleted: string[] = [] const app = await testRender( () => ( - - RUN_THEME_FALLBACK.footer} - prompts={prompts} - onClose={() => {}} - /> - + + + RUN_THEME_FALLBACK.footer} + prompts={prompts} + onClose={() => {}} + onSteer={(prompt) => steered.push(prompt.messageID)} + onDelete={(prompt) => deleted.push(prompt.messageID)} + /> + + ), { width: 100, height: RUN_SUBAGENT_PANEL_ROWS }, ) @@ -940,19 +953,75 @@ test("direct pending panel shows durable delivery without edit actions", async ( const frame = app.captureCharFrame() const list = panelMenu(app.renderer.root) - expect(frame).toContain("Pending work") + expect(frame).toContain("Queued prompts") expect(frame).toContain("fix the auth test") - expect(frame).toContain("queue") + expect(frame).toContain("queued") + expect(frame).toContain("enter steer · ctrl+d delete") expect(frame).not.toContain("┌") expect(frame).not.toContain("┃") expectPaletteList(list, 0) - expect(frame).not.toContain("edit") - expect(frame).not.toContain("remove") + app.mockInput.pressEnter() + app.mockInput.pressKey("d", { ctrl: true }) + expect(steered).toEqual(["m-1"]) + expect(deleted).toEqual(["m-1"]) } finally { app.renderer.destroy() } }) +test("direct footer steers the oldest queued prompt from an empty composer", async () => { + const steered: string[] = [] + const app = await renderFooter({ + queuedPrompts: [ + { messageID: "m-1", prompt: { text: "first", parts: [] }, delivery: "queue" }, + { messageID: "m-2", prompt: { text: "second", parts: [] }, delivery: "queue" }, + ], + onQueuedPromptSteer: async (inputID) => { + steered.push(inputID) + }, + }) + + try { + await app.renderOnce() + app.mockInput.pressEnter({ meta: true }) + await Bun.sleep(0) + expect(steered).toEqual([]) + app.mockInput.pressEnter() + await Bun.sleep(0) + expect(steered).toEqual(["m-1"]) + } finally { + app.cleanup() + } +}) + +test("direct footer does not steer queued work on a double submit", async () => { + const submitted: RunPrompt[] = [] + const steered: string[] = [] + const app = await renderFooter({ + queuedPrompts: [{ messageID: "m-1", prompt: { text: "queued", parts: [] }, delivery: "queue" }], + onSubmit: async (prompt) => { + submitted.push(prompt) + await Bun.sleep(10) + return true + }, + onQueuedPromptSteer: async (inputID) => { + steered.push(inputID) + }, + }) + + try { + await app.renderOnce() + await app.mockInput.typeText("send once") + app.mockInput.pressEnter() + app.mockInput.pressEnter() + await Bun.sleep(20) + expect(submitted).toHaveLength(1) + expect(steered).toEqual([]) + } finally { + app.cleanup() + } +}) + // OpenTUI currently crashes Bun in the full `test/cli/run` directory run here. // Re-enable after the upstream OpenTUI fix lands in this repo. test.skip("direct footer recreates the frame across command panel transitions", async () => { @@ -1245,7 +1314,7 @@ test.skip("direct footer clears the synthetic skills draft when the panel closes } }) -test("direct footer shows authoritative pending work while running", async () => { +test("direct footer shows authoritative queued work while running", async () => { const [state] = createSignal({ phase: "running", status: "", @@ -1349,9 +1418,9 @@ test("direct footer shows authoritative pending work while running", async () => const hint = statusItems.at(-1)! expect(spinner).toBeDefined() - expect(frame).toContain("1 pending") + expect(frame).toContain("1 queued") expect(frame).toContain("ctrl+b background") - expect(frame).toContain("ctrl+x q 1 pending") + expect(frame).toContain("ctrl+x q 1 queued") expect(frame).toContain("↓ subagents") expect(frame).toContain("ctrl+p cmd") expect(frame).toContain("subagents · ctrl+p cmd") diff --git a/packages/tui/test/mini/stream-v2.transport.test.ts b/packages/tui/test/mini/stream-v2.transport.test.ts index c5b997adef..71dfe221d5 100644 --- a/packages/tui/test/mini/stream-v2.transport.test.ts +++ b/packages/tui/test/mini/stream-v2.transport.test.ts @@ -669,6 +669,14 @@ describe("V2 mini transport", () => { data: { text: "follow up" }, delivery: "queue", }, + { + id: "msg_cancelled", + sessionID: "ses_1", + timeCreated: 2, + type: "user", + data: { text: "remove me" }, + delivery: "queue", + }, ], }, }) @@ -684,11 +692,14 @@ describe("V2 mini transport", () => { .findLast((item) => item.type === "queued.prompts") ?.prompts.map((item) => [item.messageID, item.delivery]) - expect(pending()).toEqual([["msg_queued", "queue"]]) + expect(pending()).toEqual([ + ["msg_queued", "queue"], + ["msg_cancelled", "queue"], + ]) events.push({ - id: "evt_promoted", - created: 2, - type: "session.input.promoted", + id: "evt_steered", + created: 3, + type: "session.input.steered", durable: durable("ses_1", 2), data: { sessionID: "ses_1", inputID: "msg_queued" }, }) @@ -697,7 +708,28 @@ describe("V2 mini transport", () => { expect(ui.commits).toContainEqual( expect.objectContaining({ kind: "user", messageID: "msg_queued", text: "follow up" }), ) - expect(pending()).toEqual([]) + expect(pending()).toEqual([ + ["msg_queued", "steer"], + ["msg_cancelled", "queue"], + ]) + events.push({ + id: "evt_cancelled", + created: 4, + type: "session.input.cancelled", + durable: durable("ses_1", 3), + data: { sessionID: "ses_1", inputID: "msg_cancelled" }, + }) + while (pending()?.length !== 1) await Bun.sleep(0) + expect(pending()).toEqual([["msg_queued", "steer"]]) + events.push({ + id: "evt_promoted", + created: 5, + type: "session.input.promoted", + durable: durable("ses_1", 4), + data: { sessionID: "ses_1", inputID: "msg_queued" }, + }) + while (pending()?.length !== 0) await Bun.sleep(0) + expect(ui.commits.filter((item) => item.messageID === "msg_queued")).toHaveLength(1) const prompt = spyOn(client.session, "prompt").mockImplementation( (request) => ok(promptAdmission(request)) as never, )