feat(tui): manage queued prompts

This commit is contained in:
Kit Langton
2026-08-06 21:28:32 -04:00
parent 9b021f5879
commit d81ae0f0d4
31 changed files with 1024 additions and 103 deletions
@@ -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" }])
})
})
@@ -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))
+59 -27
View File
@@ -251,38 +251,48 @@ export type Endpoint5_21Input = { readonly sessionID: Session.ID }
export type Endpoint5_21Output = ReadonlyArray<SessionPending.Info>
export type SessionPendingListOperation<E = never> = (input: Endpoint5_21Input) => Effect.Effect<Endpoint5_21Output, E>
export type Endpoint5_22Input = { readonly sessionID: Session.ID }
export type Endpoint5_22Output = ReadonlyArray<InstructionEntry.Info>
export type SessionInstructionsEntryListOperation<E = never> = (
export type Endpoint5_22Input = { readonly sessionID: Session.ID; readonly inputID: SessionMessage.ID }
export type Endpoint5_22Output = void
export type SessionPendingCancelOperation<E = never> = (
input: Endpoint5_22Input,
) => Effect.Effect<Endpoint5_22Output, E>
export type Endpoint5_23Input = {
export type Endpoint5_23Input = { readonly sessionID: Session.ID; readonly inputID: SessionMessage.ID }
export type Endpoint5_23Output = void
export type SessionPendingSteerOperation<E = never> = (input: Endpoint5_23Input) => Effect.Effect<Endpoint5_23Output, E>
export type Endpoint5_24Input = { readonly sessionID: Session.ID }
export type Endpoint5_24Output = ReadonlyArray<InstructionEntry.Info>
export type SessionInstructionsEntryListOperation<E = never> = (
input: Endpoint5_24Input,
) => Effect.Effect<Endpoint5_24Output, E>
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<E = never> = (
input: Endpoint5_23Input,
) => Effect.Effect<Endpoint5_23Output, E>
input: Endpoint5_25Input,
) => Effect.Effect<Endpoint5_25Output, E>
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<E = never> = (
input: Endpoint5_24Input,
) => Effect.Effect<Endpoint5_24Output, E>
input: Endpoint5_26Input,
) => Effect.Effect<Endpoint5_26Output, E>
export type Endpoint5_25Input = { readonly sessionID: Session.ID; readonly prompt: string }
export type Endpoint5_25Output = { readonly text: string }
export type SessionGenerateOperation<E = never> = (input: Endpoint5_25Input) => Effect.Effect<Endpoint5_25Output, E>
export type Endpoint5_27Input = { readonly sessionID: Session.ID; readonly prompt: string }
export type Endpoint5_27Output = { readonly text: string }
export type SessionGenerateOperation<E = never> = (input: Endpoint5_27Input) => Effect.Effect<Endpoint5_27Output, E>
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<E = never> = (input: Endpoint5_26Input) => Stream.Stream<Endpoint5_26Output, E>
export type SessionLogOperation<E = never> = (input: Endpoint5_28Input) => Stream.Stream<Endpoint5_28Output, E>
export type Endpoint5_27Input = { readonly sessionID: Session.ID }
export type Endpoint5_27Output = void
export type SessionInterruptOperation<E = never> = (input: Endpoint5_27Input) => Effect.Effect<Endpoint5_27Output, E>
export type Endpoint5_29Input = { readonly sessionID: Session.ID }
export type Endpoint5_29Output = void
export type SessionInterruptOperation<E = never> = (input: Endpoint5_29Input) => Effect.Effect<Endpoint5_29Output, E>
export type Endpoint5_28Input = { readonly sessionID: Session.ID }
export type Endpoint5_28Output = void
export type SessionBackgroundOperation<E = never> = (input: Endpoint5_28Input) => Effect.Effect<Endpoint5_28Output, E>
export type Endpoint5_30Input = { readonly sessionID: Session.ID }
export type Endpoint5_30Output = void
export type SessionBackgroundOperation<E = never> = (input: Endpoint5_30Input) => Effect.Effect<Endpoint5_30Output, E>
export type Endpoint5_29Input = { readonly sessionID: Session.ID; readonly messageID: SessionMessage.ID }
export type Endpoint5_29Output = SessionMessage.Info
export type SessionMessageOperation<E = never> = (input: Endpoint5_29Input) => Effect.Effect<Endpoint5_29Output, E>
export type Endpoint5_31Input = { readonly sessionID: Session.ID; readonly messageID: SessionMessage.ID }
export type Endpoint5_31Output = SessionMessage.Info
export type SessionMessageOperation<E = never> = (input: Endpoint5_31Input) => Effect.Effect<Endpoint5_31Output, E>
export interface SessionApi<E = never> {
readonly list: SessionListOperation<E>
@@ -888,7 +916,11 @@ export interface SessionApi<E = never> {
readonly commit: SessionRevertCommitOperation<E>
}
readonly context: SessionContextOperation<E>
readonly pending: { readonly list: SessionPendingListOperation<E> }
readonly pending: {
readonly list: SessionPendingListOperation<E>
readonly cancel: SessionPendingCancelOperation<E>
readonly steer: SessionPendingSteerOperation<E>
}
readonly instructions: {
readonly entry: {
readonly list: SessionInstructionsEntryListOperation<E>
+39 -21
View File
@@ -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<Endpoint5_22Output>()(
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<Endpoint5_23Output>()(
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<Endpoint5_24Output>()(
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<Endpoint5_23Output>()(
const Endpoint5_25 = (raw: RawClient["server.session"]) => (input: Endpoint5_25Input) =>
preserveEffect<Endpoint5_25Output>()(
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<Endpoint5_24Output>()(
const Endpoint5_26 = (raw: RawClient["server.session"]) => (input: Endpoint5_26Input) =>
preserveEffect<Endpoint5_26Output>()(
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<Endpoint5_25Output>()(
const Endpoint5_27 = (raw: RawClient["server.session"]) => (input: Endpoint5_27Input) =>
preserveEffect<Endpoint5_27Output>()(
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<Endpoint5_26Output>()(
const Endpoint5_28 = (raw: RawClient["server.session"]) => (input: Endpoint5_28Input) =>
preserveStream<Endpoint5_28Output>()(
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<Endpoint5_27Output>()(
const Endpoint5_29 = (raw: RawClient["server.session"]) => (input: Endpoint5_29Input) =>
preserveEffect<Endpoint5_29Output>()(
raw["session.interrupt"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError)),
)
const Endpoint5_28 = (raw: RawClient["server.session"]) => (input: Endpoint5_28Input) =>
preserveEffect<Endpoint5_28Output>()(
const Endpoint5_30 = (raw: RawClient["server.session"]) => (input: Endpoint5_30Input) =>
preserveEffect<Endpoint5_30Output>()(
raw["session.background"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError)),
)
const Endpoint5_29 = (raw: RawClient["server.session"]) => (input: Endpoint5_29Input) =>
preserveEffect<Endpoint5_29Output>()(
const Endpoint5_31 = (raw: RawClient["server.session"]) => (input: Endpoint5_31Input) =>
preserveEffect<Endpoint5_31Output>()(
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) =>
@@ -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<SessionPendingCancelOutput>(
{
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<SessionPendingSteerOutput>(
{
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: {
@@ -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<SessionPendingInfo> }["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<InstructionEntryInfo> }["data"]
+21
View File
@@ -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",
+44
View File
@@ -133,6 +133,13 @@ export class CompactionConflictError extends Schema.TaggedErrorClass<CompactionC
export class BusyError extends Schema.TaggedErrorClass<BusyError>()("Session.BusyError", {
sessionID: SessionSchema.ID,
}) {}
export class PendingInputConflictError extends Schema.TaggedErrorClass<PendingInputConflictError>()(
"Session.PendingInputConflictError",
{
sessionID: SessionSchema.ID,
inputID: SessionMessage.ID,
},
) {}
export class SkillNotFoundError extends Schema.TaggedErrorClass<SkillNotFoundError>()("Session.SkillNotFoundError", {
skill: Skill.ID,
}) {}
@@ -181,6 +188,14 @@ export interface Interface {
* unhandled compaction barriers.
*/
readonly pending: (sessionID: SessionSchema.ID) => Effect.Effect<SessionPending.Info[], NotFoundError>
readonly cancelPending: (input: {
sessionID: SessionSchema.ID
inputID: SessionMessage.ID
}) => Effect.Effect<void, NotFoundError | PendingInputConflictError>
readonly steerPending: (input: {
sessionID: SessionSchema.ID
inputID: SessionMessage.ID
}) => Effect.Effect<void, NotFoundError | PendingInputConflictError>
/**
* 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
@@ -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,
+69
View File
@@ -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,
+12
View File
@@ -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)
+68
View File
@@ -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)
}),
)
})
+30
View File
@@ -491,6 +491,36 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(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 },
+26 -1
View File
@@ -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
@@ -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",
+52
View File
@@ -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) {
+6 -2
View File
@@ -59,6 +59,7 @@ export type PromptProps = {
visible?: boolean
disabled?: boolean
onSubmit?: () => void
onEmptySubmit?: () => boolean | Promise<boolean>
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
}
+2
View File
@@ -104,6 +104,7 @@ export const Definitions = {
session_background: keybind("ctrl+b", "Background blocking session tools"),
session_compact: keybind("<leader>c", "Compact the session"),
session_queued_prompts: keybind("<leader>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",
+54 -8
View File
@@ -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<string, number>) => 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<string, number>, 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":
+32 -5
View File
@@ -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<RunFooterTheme>
prompts: Accessor<FooterQueuedPrompt[]>
onClose: () => void
onSteer: (prompt: FooterQueuedPrompt) => void
onDelete: (prompt: FooterQueuedPrompt) => void
onRows?: (rows: number) => void
mono?: boolean
}) {
const entries = createMemo(() =>
const entries = createMemo<QueuedPromptEntry[]>(() =>
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 (
<PanelShell
title="Pending work"
title="Queued prompts"
query={controller.query()}
count={controller.items().length}
total={entries().length}
@@ -866,6 +892,7 @@ export function RunQueuedPromptSelectBody(props: {
theme={props.theme}
inputRef={controller.inputRef}
onQuery={controller.setQuery}
hint={["enter steer", deleteShortcut() ? `${deleteShortcut()} delete` : undefined].filter(Boolean).join(" · ")}
mono={props.mono}
>
<RunFooterMenu
@@ -875,7 +902,7 @@ export function RunQueuedPromptSelectBody(props: {
offset={controller.menu.offset}
rows={controller.menu.rows}
limit={SUBAGENT_LIST_ROWS}
empty="No pending work"
empty="No queued prompts"
border={false}
paddingLeft={panelPad(props.mono)}
paddingRight={panelPad(props.mono)}
+34 -9
View File
@@ -32,7 +32,15 @@ import { realignEditorPromptParts, resolveEditorSlashValue } from "./prompt.edit
import { monoTruncateMiddle } from "./mono"
import { FOOTER_MENU_ROWS, createFooterMenuState, type RunFooterMenuItem } from "./footer.menu"
import type { RunFooterTheme } from "./theme"
import type { FooterState, RunAgent, RunCommand, RunPrompt, RunPromptPart, RunReference } from "./types"
import type {
FooterQueuedPrompt,
FooterState,
RunAgent,
RunCommand,
RunPrompt,
RunPromptPart,
RunReference,
} from "./types"
const AUTOCOMPLETE_ROWS = FOOTER_MENU_ROWS
const AUTOCOMPLETE_BOTTOM_ROWS = 1
@@ -73,6 +81,8 @@ type PromptInput = {
theme: Accessor<RunFooterTheme>
mono: Accessor<boolean>
history?: Accessor<RunPrompt[]>
queuedPrompts: Accessor<FooterQueuedPrompt[]>
onQueuedPromptSteer: (inputID: string) => Promise<boolean>
onSubmit: (input: RunPrompt) => boolean | Promise<boolean>
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)
})
}
+4
View File
@@ -96,6 +96,8 @@ type RunFooterOptions = {
onVariantSelect?: (variant: string | undefined) => CycleResult | void | Promise<CycleResult | void>
onInterrupt?: () => void
onBackground?: () => void
onQueuedPromptSteer?: (inputID: string) => Promise<void>
onQueuedPromptCancel?: (inputID: string) => Promise<void>
onEditorOpen: (input: { value: string }) => Promise<string | undefined>
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,
+33 -8
View File
@@ -92,13 +92,15 @@ type RunFooterViewProps = {
mono: boolean
miniSettings: () => MiniSettings
history?: () => RunPrompt[]
onSubmit: (input: RunPrompt) => boolean
onSubmit: (input: RunPrompt) => boolean | Promise<boolean>
onPermissionReply: (input: PermissionReply) => void | Promise<void>
onFormReply: (input: FormReply) => void | Promise<void>
onFormCancel: (input: FormCancel) => void | Promise<void>
onCycle: () => void
onInterrupt: () => boolean
onBackground?: () => void
onQueuedPromptSteer?: (inputID: string) => Promise<void>
onQueuedPromptCancel?: (inputID: string) => Promise<void>
onEditorOpen: (input: { value: string }) => Promise<string | undefined>
onInputClear: () => void
onExitRequest?: () => boolean
@@ -132,6 +134,7 @@ export function RunFooterView(props: RunFooterViewProps) {
const [route, setRoute] = createSignal<FooterPromptRoute>({ 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) {
<Match when={selectingQueued()}>
<RunQueuedPromptSelectBody
theme={theme}
prompts={queuedPrompts}
prompts={queue}
onClose={closePanel}
onSteer={(item) => {
void queuedPromptAction("steer", item.messageID).then((steered) => {
if (steered) closePanel()
})
}}
onDelete={(item) => {
void queuedPromptAction("delete", item.messageID)
}}
onRows={setSubagentMenuRows}
mono={props.mono}
/>
@@ -70,6 +70,8 @@ export type LifecycleInput = {
onVariantSelect?: (variant: string | undefined) => CycleResult | void | Promise<CycleResult | void>
onInterrupt?: () => void
onBackground?: () => void
onQueuedPromptSteer?: (inputID: string) => Promise<void>
onQueuedPromptCancel?: (inputID: string) => Promise<void>
onSubagentSelect?: (sessionID: string | undefined) => void
onSubagentInterrupt?: (sessionID: string) => void
}
@@ -243,6 +245,8 @@ export async function createRuntimeLifecycle(input: LifecycleInput): Promise<Lif
onVariantSelect: input.onVariantSelect,
onInterrupt: input.onInterrupt,
onBackground: input.onBackground,
onQueuedPromptSteer: input.onQueuedPromptSteer,
onQueuedPromptCancel: input.onQueuedPromptCancel,
onEditorOpen: async ({ value }) => {
if (closed || renderer.isDestroyed) {
return
+10
View File
@@ -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(() => {})
@@ -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)
@@ -934,6 +934,29 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
write([], { phase: "running", status: "waiting for assistant" })
return
}
if (event.type === "session.input.steered") {
const pending = state.pending.get(event.data.inputID)
if (!pending) return
state.pending.set(event.data.inputID, { ...pending, delivery: "steer" })
syncPending()
if (state.messageIDs.has(event.data.inputID)) return
state.messageIDs.add(event.data.inputID)
write([
{
kind: "user",
source: "system",
text: pending.prompt.text,
phase: "start",
messageID: event.data.inputID,
},
])
return
}
if (event.type === "session.input.cancelled") {
state.admitted.delete(event.data.inputID)
if (state.pending.delete(event.data.inputID)) syncPending()
return
}
if (event.type === "session.step.started") {
state.stepModel = { providerID: event.data.model.providerID, modelID: event.data.model.id }
write([], { phase: "running", status: "assistant responding" })
+41 -1
View File
@@ -375,6 +375,24 @@ export function Session() {
})
const dialog = useDialog()
const renderer = useRenderer()
const steerQueuedPrompt = async (inputID: string) => {
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(() => (
<DialogSelect
@@ -384,6 +402,23 @@ export function Session() {
value: prompt.id,
footer: `${index + 1} of ${queuedPrompts().length}`,
}))}
onSelect={(option) => {
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}
/>
</Match>
+86
View File
@@ -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<typeof useData>
let rows!: ReturnType<typeof createSessionRows>
let client!: ReturnType<typeof useClient>
function Probe() {
client = useClient()
data = useData()
rows = createSessionRows(() => sessionID)
return <box />
}
const app = await testRender(() => (
<TestTuiContexts>
<ClientProvider api={createApi(calls.fetch)}>
<ProjectProvider>
<DataProvider>
<Probe />
</DataProvider>
</ProjectProvider>
</ClientProvider>
</TestTuiContexts>
))
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"
+85 -16
View File
@@ -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<FooterState>
onCycle?: () => void
onSubmit?: (prompt: RunPrompt) => boolean
onSubmit?: (prompt: RunPrompt) => boolean | Promise<boolean>
view?: FooterView
onFormReply?: (input: unknown) => void
miniSettings?: MiniSettings
mono?: boolean
onStatus?: (status: string) => void
onMiniSettingChange?: (change: MiniSettingChange) => void
queuedPrompts?: FooterQueuedPrompt[]
onQueuedPromptSteer?: (inputID: string) => Promise<void>
onQueuedPromptCancel?: (inputID: string) => Promise<void>
} = {},
) {
const [view, setView] = createSignal<FooterView>(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(
() => (
<box width={100} height={RUN_SUBAGENT_PANEL_ROWS}>
<RunQueuedPromptSelectBody
theme={() => RUN_THEME_FALLBACK.footer}
prompts={prompts}
onClose={() => {}}
/>
</box>
<Keymap.Provider config={tuiConfig}>
<box width={100} height={RUN_SUBAGENT_PANEL_ROWS}>
<RunQueuedPromptSelectBody
theme={() => RUN_THEME_FALLBACK.footer}
prompts={prompts}
onClose={() => {}}
onSteer={(prompt) => steered.push(prompt.messageID)}
onDelete={(prompt) => deleted.push(prompt.messageID)}
/>
</box>
</Keymap.Provider>
),
{ 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<FooterState>({
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")
@@ -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,
)