Compare commits

...

1 Commits

Author SHA1 Message Date
Dax Raad 0c34e6669f feat(tui): continue pending work after interrupt 2026-08-10 13:52:13 +00:00
19 changed files with 70 additions and 25 deletions
+1 -1
View File
@@ -907,7 +907,7 @@ export type Endpoint5_31Output =
| EventLog.Synced | EventLog.Synced
export type SessionLogOperation<E = never> = (input: Endpoint5_31Input) => Stream.Stream<Endpoint5_31Output, E> export type SessionLogOperation<E = never> = (input: Endpoint5_31Input) => Stream.Stream<Endpoint5_31Output, E>
export type Endpoint5_32Input = { readonly sessionID: Session.ID } export type Endpoint5_32Input = { readonly sessionID: Session.ID; readonly continue?: boolean | undefined }
export type Endpoint5_32Output = void export type Endpoint5_32Output = void
export type SessionInterruptOperation<E = never> = (input: Endpoint5_32Input) => Effect.Effect<Endpoint5_32Output, E> export type SessionInterruptOperation<E = never> = (input: Endpoint5_32Input) => Effect.Effect<Endpoint5_32Output, E>
@@ -596,7 +596,10 @@ const Endpoint5_31 = (raw: RawClient["server.session"]) => (input: Endpoint5_31I
const Endpoint5_32 = (raw: RawClient["server.session"]) => (input: Endpoint5_32Input) => const Endpoint5_32 = (raw: RawClient["server.session"]) => (input: Endpoint5_32Input) =>
preserveEffect<Endpoint5_32Output>()( preserveEffect<Endpoint5_32Output>()(
raw["session.interrupt"]({ params: { sessionID: input["sessionID"] } }).pipe(Effect.mapError(mapClientError)), raw["session.interrupt"]({
params: { sessionID: input["sessionID"] },
query: { continue: input["continue"] },
}).pipe(Effect.mapError(mapClientError)),
) )
const Endpoint5_33 = (raw: RawClient["server.session"]) => (input: Endpoint5_33Input) => const Endpoint5_33 = (raw: RawClient["server.session"]) => (input: Endpoint5_33Input) =>
@@ -875,6 +875,7 @@ export function make(options: ClientOptions) {
{ {
method: "POST", method: "POST",
path: `/api/session/${encodeURIComponent(input.sessionID)}/interrupt`, path: `/api/session/${encodeURIComponent(input.sessionID)}/interrupt`,
query: { continue: input["continue"] },
successStatus: 204, successStatus: 204,
declaredStatuses: [404, 400, 401], declaredStatuses: [404, 400, 401],
empty: true, empty: true,
@@ -3888,7 +3888,10 @@ export type SessionLogInput = {
export type SessionLogOutput = SessionLogItem export type SessionLogOutput = SessionLogItem
export type SessionInterruptInput = { readonly sessionID: { readonly sessionID: string }["sessionID"] } export type SessionInterruptInput = {
readonly sessionID: { readonly sessionID: string }["sessionID"]
readonly continue?: { readonly continue?: boolean | undefined }["continue"]
}
export type SessionInterruptOutput = void export type SessionInterruptOutput = void
+1 -1
View File
@@ -199,7 +199,7 @@ test("session methods retain decoded Effect inputs and outputs", async () => {
const log = yield* client.session const log = yield* client.session
.log({ sessionID: Session.ID.make("ses_test"), after: Event.Seq.make(0) }) .log({ sessionID: Session.ID.make("ses_test"), after: Event.Seq.make(0) })
.pipe(Stream.runCollect) .pipe(Stream.runCollect)
yield* client.session.interrupt({ sessionID: Session.ID.make("ses_test") }) yield* client.session.interrupt({ sessionID: Session.ID.make("ses_test"), continue: true })
const message = yield* client.session.message({ const message = yield* client.session.message({
sessionID: Session.ID.make("ses_test"), sessionID: Session.ID.make("ses_test"),
messageID: SessionMessage.ID.make("msg_model"), messageID: SessionMessage.ID.make("msg_model"),
+2 -2
View File
@@ -543,7 +543,7 @@ test("session methods use the public HTTP contract", async () => {
const context = await client.session.context({ sessionID: "ses_test" }) const context = await client.session.context({ sessionID: "ses_test" })
const log = [] const log = []
for await (const item of client.session.log({ sessionID: "ses_test", after: 0 })) log.push(item) for await (const item of client.session.log({ sessionID: "ses_test", after: 0 })) log.push(item)
await client.session.interrupt({ sessionID: "ses_test" }) await client.session.interrupt({ sessionID: "ses_test", continue: true })
const message = await client.session.message({ sessionID: "ses_test", messageID: "msg_model" }) const message = await client.session.message({ sessionID: "ses_test", messageID: "msg_model" })
expect(page.cursor.next).toBe("next") expect(page.cursor.next).toBe("next")
@@ -568,7 +568,7 @@ test("session methods use the public HTTP contract", async () => {
["POST", "http://localhost:3000/api/session/ses_test/wait"], ["POST", "http://localhost:3000/api/session/ses_test/wait"],
["GET", "http://localhost:3000/api/session/ses_test/context"], ["GET", "http://localhost:3000/api/session/ses_test/context"],
["GET", "http://localhost:3000/api/experimental/session/ses_test/log?after=0"], ["GET", "http://localhost:3000/api/experimental/session/ses_test/log?after=0"],
["POST", "http://localhost:3000/api/session/ses_test/interrupt"], ["POST", "http://localhost:3000/api/session/ses_test/interrupt?continue=true"],
["GET", "http://localhost:3000/api/session/ses_test/message/msg_model"], ["GET", "http://localhost:3000/api/session/ses_test/message/msg_model"],
]) ])
const body = requests.find((request) => request.url.endsWith("/api/session/ses_test/prompt"))?.init?.body const body = requests.find((request) => request.url.endsWith("/api/session/ses_test/prompt"))?.init?.body
+4 -2
View File
@@ -267,7 +267,7 @@ export interface Interface {
readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>> readonly active: Effect.Effect<ReadonlySet<SessionSchema.ID>>
readonly background: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError> readonly background: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError>
readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | SessionRunner.RunError> readonly resume: (sessionID: SessionSchema.ID) => Effect.Effect<void, NotFoundError | SessionRunner.RunError>
readonly interrupt: (sessionID: SessionSchema.ID) => Effect.Effect<void> readonly interrupt: (sessionID: SessionSchema.ID, options?: { continue?: boolean }) => Effect.Effect<void>
readonly synthetic: (input: { readonly synthetic: (input: {
id?: SessionMessage.ID id?: SessionMessage.ID
sessionID: SessionSchema.ID sessionID: SessionSchema.ID
@@ -848,7 +848,9 @@ const layer = Layer.effect(
}), }),
), ),
), ),
interrupt: Effect.fn("Session.interrupt")((sessionID) => Effect.uninterruptible(execution.interrupt(sessionID))), interrupt: Effect.fn("Session.interrupt")((sessionID, options) =>
Effect.uninterruptible(execution.interrupt(sessionID, options)),
),
revert: { revert: {
stage: Effect.fn("Session.revert.stage")(function* (input) { stage: Effect.fn("Session.revert.stage")(function* (input) {
const session = yield* result.get(input.sessionID) const session = yield* result.get(input.sessionID)
+2 -2
View File
@@ -20,7 +20,7 @@ export interface Interface {
/** Registers newly recorded work. Repeated wakeups may coalesce. */ /** Registers newly recorded work. Repeated wakeups may coalesce. */
readonly wake: (sessionID: SessionSchema.ID) => Effect.Effect<void> readonly wake: (sessionID: SessionSchema.ID) => Effect.Effect<void>
/** Interrupt active work owned by this process. Idle interruption is a no-op. */ /** Interrupt active work owned by this process. Idle interruption is a no-op. */
readonly interrupt: (sessionID: SessionSchema.ID) => Effect.Effect<void> readonly interrupt: (sessionID: SessionSchema.ID, options?: { continue?: boolean }) => Effect.Effect<void>
/** Resolves once this process owns no active execution for the Session. Returns immediately when idle and never starts work. */ /** Resolves once this process owns no active execution for the Session. Returns immediately when idle and never starts work. */
readonly awaitIdle: (sessionID: SessionSchema.ID) => Effect.Effect<void> readonly awaitIdle: (sessionID: SessionSchema.ID) => Effect.Effect<void>
} }
@@ -107,7 +107,7 @@ export const layer = Layer.effect(
return Service.of({ return Service.of({
active: coordinator.active, active: coordinator.active,
interrupt: (sessionID) => coordinator.interrupt(sessionID, "user"), interrupt: (sessionID, options) => coordinator.interrupt(sessionID, "user", { preserveWake: options?.continue }),
resume: coordinator.run, resume: coordinator.run,
wake: coordinator.wake, wake: coordinator.wake,
awaitIdle: coordinator.awaitIdle, awaitIdle: coordinator.awaitIdle,
+4 -4
View File
@@ -10,8 +10,8 @@ export interface Coordinator<Key, E, Reason = never> {
readonly run: (key: Key) => Effect.Effect<void, E> readonly run: (key: Key) => Effect.Effect<void, E>
/** Rings the doorbell: an idle key starts an execution; an active one drains again before settling. */ /** Rings the doorbell: an idle key starts an execution; an active one drains again before settling. */
readonly wake: (key: Key) => Effect.Effect<void> readonly wake: (key: Key) => Effect.Effect<void>
/** Stops the active execution, clears its doorbell, and waits for cleanup. No-op when idle. */ /** Stops the active execution and waits for cleanup. Clears its doorbell unless preservation is requested. */
readonly interrupt: (key: Key, reason?: Reason) => Effect.Effect<void> readonly interrupt: (key: Key, reason?: Reason, options?: { preserveWake?: boolean }) => Effect.Effect<void>
/** Resolves once no execution is active for the key. Returns immediately when already idle and never starts work. */ /** Resolves once no execution is active for the key. Returns immediately when already idle and never starts work. */
readonly awaitIdle: (key: Key) => Effect.Effect<void> readonly awaitIdle: (key: Key) => Effect.Effect<void>
} }
@@ -124,12 +124,12 @@ export const make = <Key, E, Reason = never>(options: {
start(key, false) start(key, false)
}) })
const interrupt = (key: Key, reason?: Reason): Effect.Effect<void> => const interrupt = (key: Key, reason?: Reason, options?: { preserveWake?: boolean }): Effect.Effect<void> =>
Effect.suspend(() => { Effect.suspend(() => {
const execution = executions.get(key) const execution = executions.get(key)
if (execution?.owner === undefined || execution.stopping) return Effect.void if (execution?.owner === undefined || execution.stopping) return Effect.void
execution.stopping = true execution.stopping = true
execution.pendingWake = false if (!options?.preserveWake) execution.pendingWake = false
execution.interruptionReason = reason execution.interruptionReason = reason
return Fiber.interrupt(execution.owner) return Fiber.interrupt(execution.owner)
}) })
@@ -301,6 +301,37 @@ describe("SessionRunCoordinator", () => {
), ),
) )
it.effect("continues pending work after interruption when preserving the wake", () =>
Effect.scoped(
Effect.gen(function* () {
const firstStarted = yield* Deferred.make<void>()
const secondStarted = yield* Deferred.make<void>()
let runs = 0
const coordinator = yield* SessionRunCoordinator.make<string, never, string>({
drain: () =>
Effect.sync(() => ++runs).pipe(
Effect.flatMap((run) =>
run === 1
? Deferred.succeed(firstStarted, undefined).pipe(Effect.andThen(Effect.never))
: Deferred.succeed(secondStarted, undefined),
),
),
})
const resumed = yield* coordinator.run("session").pipe(Effect.forkChild)
yield* Deferred.await(firstStarted)
yield* coordinator.wake("session")
yield* coordinator.interrupt("session", "user", { preserveWake: true })
yield* Deferred.await(secondStarted)
const exit = yield* Fiber.await(resumed)
expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue()
yield* coordinator.awaitIdle("session")
expect(runs).toBe(2)
}),
),
)
it.effect("runs a wake registered during interruption cleanup", () => it.effect("runs a wake registered during interruption cleanup", () =>
Effect.scoped( Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
@@ -132,7 +132,8 @@ const execution = (llmClient: Layer.Layer<typeof LLMClient.Service>) =>
active: coordinator.active, active: coordinator.active,
resume: coordinator.run, resume: coordinator.run,
wake: coordinator.wake, wake: coordinator.wake,
interrupt: coordinator.interrupt, interrupt: (sessionID, options) =>
coordinator.interrupt(sessionID, undefined, { preserveWake: options?.continue }),
awaitIdle: coordinator.awaitIdle, awaitIdle: coordinator.awaitIdle,
}) })
}), }),
+2 -1
View File
@@ -391,7 +391,8 @@ const execution = Layer.effect(
active: coordinator.active, active: coordinator.active,
resume: coordinator.run, resume: coordinator.run,
wake: coordinator.wake, wake: coordinator.wake,
interrupt: coordinator.interrupt, interrupt: (sessionID, options) =>
coordinator.interrupt(sessionID, undefined, { preserveWake: options?.continue }),
awaitIdle: coordinator.awaitIdle, awaitIdle: coordinator.awaitIdle,
}) })
}), }),
+3 -1
View File
@@ -647,6 +647,7 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
.add( .add(
HttpApiEndpoint.post("session.interrupt", "/api/session/:sessionID/interrupt", { HttpApiEndpoint.post("session.interrupt", "/api/session/:sessionID/interrupt", {
params: { sessionID: Session.ID }, params: { sessionID: Session.ID },
query: { continue: BooleanFromString.pipe(Schema.optional) },
success: HttpApiSchema.NoContent, success: HttpApiSchema.NoContent,
error: SessionNotFoundError, error: SessionNotFoundError,
}) })
@@ -655,7 +656,8 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
OpenApi.annotations({ OpenApi.annotations({
identifier: "v2.session.interrupt", identifier: "v2.session.interrupt",
summary: "Interrupt session execution", summary: "Interrupt session execution",
description: "Interrupt active execution owned by this OpenCode process. Idle interruption is a no-op.", description:
"Interrupt active execution owned by this OpenCode process. When continue=true, pending work starts after interruption. Idle interruption is a no-op.",
}), }),
), ),
) )
+1 -1
View File
@@ -772,7 +772,7 @@ export const SessionHandler = HttpApiBuilder.group(Api, "server.session", (handl
.handle( .handle(
"session.interrupt", "session.interrupt",
Effect.fn(function* (ctx) { Effect.fn(function* (ctx) {
yield* session.interrupt(ctx.params.sessionID) yield* session.interrupt(ctx.params.sessionID, { continue: ctx.query.continue })
return HttpApiSchema.NoContent.make() return HttpApiSchema.NoContent.make()
}), }),
) )
@@ -432,6 +432,7 @@ export function Prompt(props: PromptProps) {
if (store.interrupt >= 2) { if (store.interrupt >= 2) {
void client.api.session.interrupt({ void client.api.session.interrupt({
sessionID: props.sessionID, sessionID: props.sessionID,
continue: true,
}) })
setStore("interrupt", 0) setStore("interrupt", 0)
} }
+2 -2
View File
@@ -374,7 +374,7 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep
void ( void (
state.stream state.stream
? state.stream.then((item) => item.handle.interruptActiveTurn()) ? state.stream.then((item) => item.handle.interruptActiveTurn())
: state.sdk.session.interrupt({ sessionID: state.sessionID }) : state.sdk.session.interrupt({ sessionID: state.sessionID, continue: true })
) )
.catch(() => {}) .catch(() => {})
.finally(() => { .finally(() => {
@@ -401,7 +401,7 @@ async function runInteractiveRuntime(input: RunRuntimeInput, deps: RunRuntimeDep
}, },
onSubagentInterrupt: (sessionID) => { onSubagentInterrupt: (sessionID) => {
log?.write("send.subagent.interrupt", { sessionID }) log?.write("send.subagent.interrupt", { sessionID })
void state.sdk.session.interrupt({ sessionID }).catch(() => {}) void state.sdk.session.interrupt({ sessionID, continue: true }).catch(() => {})
}, },
onSubagentSelect: (sessionID) => { onSubagentSelect: (sessionID) => {
state.selectSubagent?.(sessionID) state.selectSubagent?.(sessionID)
+2 -2
View File
@@ -1525,7 +1525,7 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
state.wait = active state.wait = active
const interrupt = () => { const interrupt = () => {
active.interrupted = true active.interrupted = true
void sdk.session.interrupt({ sessionID: input.sessionID }).catch(() => {}) void sdk.session.interrupt({ sessionID: input.sessionID, continue: true }).catch(() => {})
} }
next.signal?.addEventListener("abort", interrupt, { once: true }) next.signal?.addEventListener("abort", interrupt, { once: true })
try { try {
@@ -1783,7 +1783,7 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
return return
} }
if (state.wait) state.wait.interrupted = true if (state.wait) state.wait.interrupted = true
await sdk.session.interrupt({ sessionID: input.sessionID }).catch(() => {}) await sdk.session.interrupt({ sessionID: input.sessionID, continue: true }).catch(() => {})
}, },
selectSubagent(sessionID) { selectSubagent(sessionID) {
subagents.select(sdk, sessionID) subagents.select(sdk, sessionID)
@@ -217,7 +217,7 @@ export function SubagentsTab(props: { sessionID: string }) {
run() { run() {
const entry = selectedEntry() const entry = selectedEntry()
if (!entry || entry.status !== "running") return if (!entry || entry.status !== "running") return
void client.api.session.interrupt({ sessionID: entry.sessionID }) void client.api.session.interrupt({ sessionID: entry.sessionID, continue: true })
}, },
}, },
], ],
@@ -1496,7 +1496,7 @@ describe("V2 mini transport", () => {
await transport.interruptActiveTurn() await transport.interruptActiveTurn()
expect(prompt).toHaveBeenCalled() expect(prompt).toHaveBeenCalled()
expect(interrupt).toHaveBeenCalledWith({ sessionID: "ses_1" }) expect(interrupt).toHaveBeenCalledWith({ sessionID: "ses_1", continue: true })
expect(firstPrompt).not.toHaveBeenCalled() expect(firstPrompt).not.toHaveBeenCalled()
expect(firstInterrupt).not.toHaveBeenCalled() expect(firstInterrupt).not.toHaveBeenCalled()
await transport.close() await transport.close()
@@ -2315,7 +2315,7 @@ describe("V2 mini transport", () => {
idle.resolve() idle.resolve()
await turn await turn
expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1" }) expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1", continue: true })
await transport.close() await transport.close()
}) })