From e4095af8d33a70648fd72f7ea49488848ba1349d Mon Sep 17 00:00:00 2001 From: marius-kilocode Date: Fri, 29 May 2026 12:15:48 +0200 Subject: [PATCH] refactor(cli): remove Question compatibility facade --- .../src/kilo-sessions/kilo-sessions.ts | 4 +- .../src/kilo-sessions/remote-sender.ts | 25 +- .../opencode/src/kilocode/plan-followup.ts | 3 + packages/opencode/src/question/index.ts | 22 +- .../instance/httpapi/handlers/question.ts | 18 +- .../src/server/routes/instance/question.ts | 28 +- packages/opencode/src/session/prompt.ts | 4 +- .../test/kilocode/plan-exit-detection.test.ts | 19 +- .../kilocode/prompt-dismiss-contract.test.ts | 4 +- .../kilocode/question-dismiss-all.test.ts | 239 +++++++----------- .../kilocode/session-compaction-cap.test.ts | 1 + .../session-prompt-compaction-safety.test.ts | 1 + .../kilocode/session-prompt-queue.test.ts | 120 --------- .../kilocode/sessions/remote-sender.test.ts | 121 +++++++-- .../opencode/test/question/question.test.ts | 62 +++-- packages/opencode/test/session/prompt.test.ts | 46 ++++ .../test/session/snapshot-tool-race.test.ts | 1 + script/check-opencode-promise-facades.ts | 1 - 18 files changed, 367 insertions(+), 352 deletions(-) diff --git a/packages/opencode/src/kilo-sessions/kilo-sessions.ts b/packages/opencode/src/kilo-sessions/kilo-sessions.ts index 7d2aecf2c4f..75da60c7c24 100644 --- a/packages/opencode/src/kilo-sessions/kilo-sessions.ts +++ b/packages/opencode/src/kilo-sessions/kilo-sessions.ts @@ -196,7 +196,9 @@ export namespace KiloSessions { ).filter((p) => p.sessionID === sessionID) if (permissions.length > 0) return "permission" - const questions = (await Question.list()).filter((q) => q.sessionID === sessionID) + const questions = ( + await AppRuntime.runPromise(Question.Service.use((svc) => svc.list())) + ).filter((q) => q.sessionID === sessionID) if (questions.length > 0) return "question" const status = await AppRuntime.runPromise(SessionStatus.Service.use((svc) => svc.get(SessionID.make(sessionID)))) diff --git a/packages/opencode/src/kilo-sessions/remote-sender.ts b/packages/opencode/src/kilo-sessions/remote-sender.ts index 38a70975b3f..5b6b2a3d081 100644 --- a/packages/opencode/src/kilo-sessions/remote-sender.ts +++ b/packages/opencode/src/kilo-sessions/remote-sender.ts @@ -76,6 +76,11 @@ export namespace RemoteSender { readonly list: () => Promise> readonly reply: (input: Permission.ReplyInput) => Promise } + question?: { + readonly list: () => Promise> + readonly reply: (input: Parameters[0]) => Promise + readonly reject: (requestID: QuestionID) => Promise + } } export type Sender = { @@ -97,6 +102,20 @@ export namespace RemoteSender { return AppRuntime.runPromise(Permission.Service.use((svc) => svc.reply(input))) }, } + const question = options.question ?? { + list: async () => { + const { AppRuntime } = await import("@/effect/app-runtime") + return AppRuntime.runPromise(Question.Service.use((svc) => svc.list())) + }, + reply: async (input: Parameters[0]) => { + const { AppRuntime } = await import("@/effect/app-runtime") + return AppRuntime.runPromise(Question.Service.use((svc) => svc.reply(input))) + }, + reject: async (requestID: QuestionID) => { + const { AppRuntime } = await import("@/effect/app-runtime") + return AppRuntime.runPromise(Question.Service.use((svc) => svc.reject(requestID))) + }, + } const sub = options.subscribe ?? @@ -146,7 +165,7 @@ export namespace RemoteSender { async function replay(sessionId: string) { const [suggestions, questions, permissions] = await Promise.all([ Suggestion.list(), - Question.list(), + question.list(), permission.list(), ]) for (const suggestion of suggestions) { @@ -310,7 +329,7 @@ export namespace RemoteSender { } const dir = msg.sessionId ? directoryFor(msg.sessionId) : Promise.resolve(options.directory) dispatchQuick(msg, dir, () => - Question.reply({ ...parsed.data, requestID: QuestionID.make(parsed.data.requestID) }), + question.reply({ ...parsed.data, requestID: QuestionID.make(parsed.data.requestID) }), ) return } @@ -325,7 +344,7 @@ export namespace RemoteSender { return } const dir = msg.sessionId ? directoryFor(msg.sessionId) : Promise.resolve(options.directory) - dispatchQuick(msg, dir, () => Question.reject(QuestionID.make(parsed.data.requestID))) + dispatchQuick(msg, dir, () => question.reject(QuestionID.make(parsed.data.requestID))) return } if (msg.command === "suggestion_accept") { diff --git a/packages/opencode/src/kilocode/plan-followup.ts b/packages/opencode/src/kilocode/plan-followup.ts index 50f1953af15..343f0533f29 100644 --- a/packages/opencode/src/kilocode/plan-followup.ts +++ b/packages/opencode/src/kilocode/plan-followup.ts @@ -48,6 +48,9 @@ export const PlanFollowupRuntime = { reject(requestID: Parameters[0]) { return questions().runPromise((svc) => svc.reject(requestID)) }, + reply(input: Parameters[0]) { + return questions().runPromise((svc) => svc.reply(input)) + }, }, todo: { get(sessionID: SessionID) { diff --git a/packages/opencode/src/question/index.ts b/packages/opencode/src/question/index.ts index 70d60171117..4a4901ebc4e 100644 --- a/packages/opencode/src/question/index.ts +++ b/packages/opencode/src/question/index.ts @@ -7,7 +7,6 @@ import { zod } from "@/util/effect-zod" import * as Log from "@opencode-ai/core/util/log" import { withStatics } from "@/util/schema" import { QuestionID } from "./schema" -import { makeRuntime } from "@/effect/run-service" // kilocode_change import { KiloQuestion } from "@/kilocode/question" // kilocode_change const log = Log.create({ service: "question" }) @@ -134,6 +133,10 @@ export class RejectedError extends Schema.TaggedErrorClass()("Que } } +export class NotFoundError extends Schema.TaggedErrorClass()("Question.NotFoundError", { + requestID: QuestionID, +}) {} + interface PendingEntry { info: Request deferred: Deferred.Deferred, RejectedError> @@ -152,8 +155,8 @@ export interface Interface { blocking?: boolean // kilocode_change tool?: Tool }) => Effect.Effect, RejectedError> - readonly reply: (input: { requestID: QuestionID; answers: ReadonlyArray }) => Effect.Effect - readonly reject: (requestID: QuestionID) => Effect.Effect + readonly reply: (input: { requestID: QuestionID; answers: ReadonlyArray }) => Effect.Effect + readonly reject: (requestID: QuestionID) => Effect.Effect readonly list: () => Effect.Effect> readonly dismissAll: (sessionID: SessionID) => Effect.Effect // kilocode_change } @@ -225,7 +228,7 @@ export const layer = Layer.effect( const existing = pending.get(input.requestID) if (!existing) { log.warn("reply for unknown request", { requestID: input.requestID }) - return + return yield* new NotFoundError({ requestID: input.requestID }) } pending.delete(input.requestID) log.info("replied", { requestID: input.requestID, answers: input.answers }) @@ -242,7 +245,7 @@ export const layer = Layer.effect( const existing = pending.get(requestID) if (!existing) { log.warn("reject for unknown request", { requestID }) - return + return yield* new NotFoundError({ requestID }) } pending.delete(requestID) log.info("rejected", { requestID }) @@ -273,13 +276,4 @@ export const layer = Layer.effect( export const defaultLayer = layer.pipe(Layer.provide(Bus.layer)) -// kilocode_change start - legacy promise helpers for Kilo callsites -const { runPromise } = makeRuntime(Service, defaultLayer) -export const list = () => runPromise((svc) => svc.list()) -export const ask = (input: Parameters[0]) => runPromise((svc) => svc.ask(input)) -export const reply = (input: Parameters[0]) => runPromise((svc) => svc.reply(input)) -export const reject = (requestID: QuestionID) => runPromise((svc) => svc.reject(requestID)) -export const dismissAll = (sessionID: string) => runPromise((svc) => svc.dismissAll(SessionID.make(sessionID))) -// kilocode_change end - export * as Question from "." diff --git a/packages/opencode/src/server/routes/instance/httpapi/handlers/question.ts b/packages/opencode/src/server/routes/instance/httpapi/handlers/question.ts index 3a4d316179c..8aac0cea8de 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/handlers/question.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/handlers/question.ts @@ -1,7 +1,7 @@ import { Question } from "@/question" import { QuestionID } from "@/question/schema" import { Effect } from "effect" -import { HttpApiBuilder } from "effect/unstable/httpapi" +import { HttpApiBuilder, HttpApiError } from "effect/unstable/httpapi" // kilocode_change - map Question missing requests to declared 404 errors import { InstanceHttpApi } from "../api" export const questionHandlers = HttpApiBuilder.group(InstanceHttpApi, "question", (handlers) => @@ -16,15 +16,21 @@ export const questionHandlers = HttpApiBuilder.group(InstanceHttpApi, "question" params: { requestID: QuestionID } payload: Question.Reply }) { - yield* svc.reply({ - requestID: ctx.params.requestID, - answers: ctx.payload.answers, - }) + // kilocode_change start - map missing Question requests to the declared transport error + yield* svc + .reply({ + requestID: ctx.params.requestID, + answers: ctx.payload.answers, + }) + .pipe(Effect.mapError(() => new HttpApiError.NotFound({}))) + // kilocode_change end return true }) const reject = Effect.fn("QuestionHttpApi.reject")(function* (ctx: { params: { requestID: QuestionID } }) { - yield* svc.reject(ctx.params.requestID) + // kilocode_change start - map missing Question requests to the declared transport error + yield* svc.reject(ctx.params.requestID).pipe(Effect.mapError(() => new HttpApiError.NotFound({}))) + // kilocode_change end return true }) diff --git a/packages/opencode/src/server/routes/instance/question.ts b/packages/opencode/src/server/routes/instance/question.ts index 51ecb48ccdd..2b574a470dd 100644 --- a/packages/opencode/src/server/routes/instance/question.ts +++ b/packages/opencode/src/server/routes/instance/question.ts @@ -1,8 +1,10 @@ +import { Effect } from "effect" // kilocode_change - translate Question not-found failures for the legacy route import { Hono } from "hono" import { describeRoute, validator } from "hono-openapi" import { resolver } from "hono-openapi" import { QuestionID } from "@/question/schema" import { Question } from "@/question" +import { NotFoundError } from "@/storage/storage" // kilocode_change - expose upstream Question not-found behavior on legacy routes import z from "zod" import { errors } from "../../error" import { lazy } from "@/util/lazy" @@ -69,10 +71,18 @@ export const QuestionRoutes = lazy(() => const params = c.req.valid("param") const json = c.req.valid("json") const svc = yield* Question.Service - yield* svc.reply({ - requestID: params.requestID, - answers: json.answers, - }) + // kilocode_change start - preserve documented 404 for unknown requests + yield* svc + .reply({ + requestID: params.requestID, + answers: json.answers, + }) + .pipe( + Effect.mapError( + () => new NotFoundError({ message: `Question request not found: ${params.requestID}` }), + ), + ) + // kilocode_change end return true }), ) @@ -104,7 +114,15 @@ export const QuestionRoutes = lazy(() => jsonRequest("QuestionRoutes.reject", c, function* () { const params = c.req.valid("param") const svc = yield* Question.Service - yield* svc.reject(params.requestID) + // kilocode_change start - preserve documented 404 for unknown requests + yield* svc + .reject(params.requestID) + .pipe( + Effect.mapError( + () => new NotFoundError({ message: `Question request not found: ${params.requestID}` }), + ), + ) + // kilocode_change end return true }), ), diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index 17748fa93f7..aff791c03ad 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -118,6 +118,7 @@ export const layer = Layer.effect( const commands = yield* Command.Service const config = yield* Config.Service const permission = yield* Permission.Service + const question = yield* Question.Service // kilocode_change - dismiss superseded pending questions through the shared service const fsys = yield* AppFileSystem.Service const mcp = yield* MCP.Service const lsp = yield* LSP.Service @@ -1441,7 +1442,7 @@ NOTE: At any point in time through this workflow you should feel free to ask the // runLoop checks hasFollowup between steps to break out once it has been // enqueued during the turn. yield* Effect.promise(() => Suggestion.dismissAll(input.sessionID)) - yield* Effect.promise(() => Question.dismissAll(input.sessionID)) + yield* question.dismissAll(input.sessionID) if (input.noReply === true) return message return yield* KiloSessionPromptQueue.enqueue( input.sessionID, @@ -2005,6 +2006,7 @@ export const defaultLayer = Layer.suspend(() => Layer.provide(SessionProcessor.defaultLayer), Layer.provide(Command.defaultLayer), Layer.provide(Permission.defaultLayer), + Layer.provide(Question.defaultLayer), // kilocode_change - provide pending question dismissal dependency Layer.provide(MCP.defaultLayer), Layer.provide(LSP.defaultLayer), Layer.provide(ToolRegistry.defaultLayer), diff --git a/packages/opencode/test/kilocode/plan-exit-detection.test.ts b/packages/opencode/test/kilocode/plan-exit-detection.test.ts index 5e6371b71b2..162cc3c93ef 100644 --- a/packages/opencode/test/kilocode/plan-exit-detection.test.ts +++ b/packages/opencode/test/kilocode/plan-exit-detection.test.ts @@ -6,8 +6,7 @@ import { SessionID, MessageID, PartID } from "../../src/session/schema" import { ModelID, ProviderID } from "../../src/provider/schema" import { Instance } from "../../src/project/instance" import { WithInstance } from "../../src/project/with-instance" -import { PlanFollowup } from "../../src/kilocode/plan-followup" -import { Question } from "../../src/question" +import { PlanFollowup, PlanFollowupRuntime } from "../../src/kilocode/plan-followup" import { Session } from "../../src/session/session" import { MessageV2 } from "../../src/session/message-v2" import { SessionPrompt } from "../../src/session/prompt" @@ -109,7 +108,7 @@ async function seed(input: { async function waitQuestion(sessionID: string) { for (let i = 0; i < 50; i++) { - const list = await Question.list() + const list = await PlanFollowupRuntime.question.list() const question = list.find((item) => item.sessionID === sessionID) if (question) return question await Bun.sleep(10) @@ -141,7 +140,7 @@ describe("plan_exit detection", () => { expect(question).toBeDefined() if (!question) return expect(question.questions[0].header).toBe("Implement") - await Question.reject(question.id) + await PlanFollowupRuntime.question.reject(question.id) await expect(pending).resolves.toBe("break") })) @@ -180,7 +179,7 @@ describe("plan_exit detection", () => { PlanFollowup.ANSWER_CONTINUE, ]) expect(question.questions[0].options.find((item) => item.label === PlanFollowup.ANSWER_CONTINUE)?.mode).toBe("code") - await Question.reject(question.id) + await PlanFollowupRuntime.question.reject(question.id) await expect(pending).resolves.toBe("break") } finally { if (prev === undefined) delete process.env.KILO_CLIENT @@ -210,7 +209,7 @@ describe("plan_exit detection", () => { const question = await waitQuestion(seeded.sessionID) expect(question).toBeDefined() if (!question) return - await Question.reply({ + await PlanFollowupRuntime.question.reply({ requestID: question.id, answers: [[PlanFollowup.ANSWER_CONTINUE]], }) @@ -233,7 +232,7 @@ describe("plan_exit detection", () => { text: "Here is a partial plan, I have questions", }) expect(SessionPrompt.shouldAskPlanFollowup({ messages: seeded.messages, abort: AbortSignal.any([]) })).toBe(false) - const list = await Question.list() + const list = await PlanFollowupRuntime.question.list() expect(list).toHaveLength(0) })) @@ -312,7 +311,7 @@ describe("plan_exit detection", () => { expect(SessionPrompt.shouldAskPlanFollowup({ messages, abort: AbortSignal.any([]) })).toBe(false) // Confirm no questions were posted - const list = await Question.list() + const list = await PlanFollowupRuntime.question.list() expect(list).toHaveLength(0) })) @@ -426,7 +425,7 @@ describe("plan_exit detection", () => { const question = await waitQuestion(seeded.sessionID) expect(question).toBeDefined() if (!question) return - await Question.reply({ + await PlanFollowupRuntime.question.reply({ requestID: question.id, answers: [[PlanFollowup.ANSWER_CONTINUE]], }) @@ -529,7 +528,7 @@ describe("plan_exit detection", () => { expect(question).toBeDefined() if (!question) return expect(question.questions[0].header).toBe("Implement") - await Question.reply({ + await PlanFollowupRuntime.question.reply({ requestID: question.id, answers: [[PlanFollowup.ANSWER_CONTINUE]], }) diff --git a/packages/opencode/test/kilocode/prompt-dismiss-contract.test.ts b/packages/opencode/test/kilocode/prompt-dismiss-contract.test.ts index c5981d9f6e7..abcad3f51ad 100644 --- a/packages/opencode/test/kilocode/prompt-dismiss-contract.test.ts +++ b/packages/opencode/test/kilocode/prompt-dismiss-contract.test.ts @@ -36,9 +36,9 @@ describe("prompt.ts Kilo-specific invariants", () => { // an in-flight handle.process blocked on a pending tool prompt can return. // Critically, the block must NOT call state.cancel or KiloSessionPromptQueue.reserve — // either of those would abort the running streamText mid-tokens, which was - // the #9332 regression. Order: dismissAll(Suggestion) → dismissAll(Question) → enqueue. + // the #9332 regression. Order: dismissAll(Suggestion), question.dismissAll, enqueue. const block = content.match( - /kilocode_change start[^\n]*unblock tools[\s\S]*?Suggestion\.dismissAll[\s\S]*?Question\.dismissAll[\s\S]*?KiloSessionPromptQueue\.enqueue/, + /kilocode_change start[^\n]*unblock tools[\s\S]*?Suggestion\.dismissAll[\s\S]*?question\.dismissAll[\s\S]*?KiloSessionPromptQueue\.enqueue/, ) expect(block).not.toBeNull() expect(content).not.toMatch(/state\.cancel\(input\.sessionID\)/) diff --git a/packages/opencode/test/kilocode/question-dismiss-all.test.ts b/packages/opencode/test/kilocode/question-dismiss-all.test.ts index b82ae133713..8b41c6bb2fb 100644 --- a/packages/opencode/test/kilocode/question-dismiss-all.test.ts +++ b/packages/opencode/test/kilocode/question-dismiss-all.test.ts @@ -1,174 +1,119 @@ -import { describe, expect, test } from "bun:test" -import { Effect } from "effect" +import { describe, expect } from "bun:test" +import { Cause, Effect, Exit, Fiber, Layer } from "effect" +import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" import { KiloSessionPromptQueue } from "../../src/kilocode/session/prompt-queue" -import { WithInstance } from "../../src/project/with-instance" import { Question } from "../../src/question" import { MessageID, SessionID } from "../../src/session/schema" -import { tmpdir } from "../fixture/fixture" +import { testEffect } from "../lib/effect" + +const it = testEffect(Layer.mergeAll(Question.defaultLayer, CrossSpawnSpawner.defaultLayer)) + +const prompt = [ + { + header: "Continue?", + question: "Should I continue?", + options: [ + { label: "Yes", description: "Go" }, + { label: "No", description: "Stop" }, + ], + }, +] + +const waitFor = (question: Question.Interface, count: number) => + Effect.gen(function* () { + for (let i = 0; i < 50; i++) { + const pending = yield* question.list() + if (pending.length >= count) return pending + yield* Effect.sleep("10 millis") + } + return yield* Effect.fail(new Error(`timed out waiting for ${count} pending question request(s)`)) + }) describe("Question.dismissAll", () => { - test("rejects pending asks for the target session and clears them", async () => { - await using tmp = await tmpdir({ git: true }) - await WithInstance.provide({ - directory: tmp.path, - fn: async () => { + it.instance( + "rejects pending asks for the target session and clears them", + () => + Effect.gen(function* () { + const question = yield* Question.Service const sesA = SessionID.make("ses_a") const sesB = SessionID.make("ses_b") + const a1 = yield* question.ask({ sessionID: sesA, questions: prompt }).pipe(Effect.forkScoped) + const a2 = yield* question.ask({ sessionID: sesA, questions: prompt }).pipe(Effect.forkScoped) + const b1 = yield* question.ask({ sessionID: sesB, questions: prompt }).pipe(Effect.forkScoped) - const a1 = Question.ask({ - sessionID: sesA, - questions: [ - { - header: "Continue?", - question: "Should I continue?", - options: [ - { label: "Yes", description: "Go" }, - { label: "No", description: "Stop" }, - ], - }, - ], - }).catch((err) => { - if (err instanceof Question.RejectedError) return "rejected" - throw err - }) + expect(yield* waitFor(question, 3)).toHaveLength(3) + yield* question.dismissAll(sesA) - const a2 = Question.ask({ - sessionID: sesA, - questions: [ - { - header: "Retry?", - question: "Try again?", - options: [ - { label: "Retry", description: "Retry" }, - { label: "Cancel", description: "Cancel" }, - ], - }, - ], - }).catch((err) => { - if (err instanceof Question.RejectedError) return "rejected" - throw err - }) - - const b1 = Question.ask({ - sessionID: sesB, - questions: [ - { - header: "Deploy?", - question: "Deploy now?", - options: [ - { label: "Ship", description: "Ship" }, - { label: "Wait", description: "Wait" }, - ], - }, - ], - }).catch((err) => { - if (err instanceof Question.RejectedError) return "rejected-b" - throw err - }) - - // Wait for all three asks to register so we can dismiss them. - for (let i = 0; i < 50; i++) { - if ((await Question.list()).length >= 3) break - await Bun.sleep(10) + for (const fiber of [a1, a2]) { + const exit = yield* Fiber.await(fiber) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Question.RejectedError) } - expect(await Question.list()).toHaveLength(3) - // Track whether B's promise settles. - let settled = false - b1.then(() => { - settled = true - }) + yield* Effect.sleep("10 millis") - await Question.dismissAll("ses_a") - - expect(await a1).toBe("rejected") - expect(await a2).toBe("rejected") - - await new Promise((r) => setTimeout(r, 10)) - expect(settled).toBe(false) - - const remaining = await Question.list() + const remaining = yield* question.list() expect(remaining).toHaveLength(1) expect(remaining[0]?.sessionID).toBe(sesB) - await Question.reject(remaining[0]!.id) - expect(await b1).toBe("rejected-b") - }, - }) - }) + yield* question.reject(remaining[0]!.id) + const exit = yield* Fiber.await(b1) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Question.RejectedError) + }), + { git: true }, + ) - test("is a no-op when no questions exist", async () => { - await using tmp = await tmpdir({ git: true }) - await WithInstance.provide({ - directory: tmp.path, - fn: async () => { - await Question.dismissAll("ses_missing") - expect(await Question.list()).toEqual([]) - }, - }) - }) + it.instance( + "is a no-op when no questions exist", + () => + Effect.gen(function* () { + const question = yield* Question.Service + yield* question.dismissAll(SessionID.make("ses_missing")) + expect(yield* question.list()).toEqual([]) + }), + { git: true }, + ) - test("ask rejects immediately when a followup is queued on the session", async () => { - // When a newer prompt has already been enqueued on the session, a tool - // that subsequently calls Question.ask would otherwise block the run until - // the user manually dismisses it. Verify the pre-emptive hasFollowup check - // rejects with RejectedError before any pending entry is registered. - await using tmp = await tmpdir({ git: true }) - await WithInstance.provide({ - directory: tmp.path, - fn: async () => { + it.instance( + "ask rejects immediately when a followup is queued on the session", + () => + Effect.gen(function* () { + const question = yield* Question.Service const sessionID = SessionID.make("ses_auto_ask") const started = Promise.withResolvers() const release = Promise.withResolvers() - // Slot 1 stays running so activeSince is pinned to its seq. - const first = Effect.runPromise( - KiloSessionPromptQueue.enqueue( - sessionID, - MessageID.make("message_ask_1"), - Effect.gen(function* () { - started.resolve() - yield* Effect.promise(() => release.promise) - return "first" as const - }), - Effect.succeed("first-cancelled" as const), - ), - ) - await started.promise + const first = yield* KiloSessionPromptQueue.enqueue( + sessionID, + MessageID.make("message_ask_1"), + Effect.gen(function* () { + started.resolve() + yield* Effect.promise(() => release.promise) + return "first" as const + }), + Effect.succeed("first-cancelled" as const), + ).pipe(Effect.forkScoped) + yield* Effect.promise(() => started.promise) - // Slot 2 arrives while slot 1 is active — latest > activeSince. - const second = Effect.runPromise( - KiloSessionPromptQueue.enqueue( - sessionID, - MessageID.make("message_ask_2"), - Effect.succeed("second" as const), - Effect.succeed("second-cancelled" as const), - ), - ) - await Bun.sleep(10) + const second = yield* KiloSessionPromptQueue.enqueue( + sessionID, + MessageID.make("message_ask_2"), + Effect.succeed("second" as const), + Effect.succeed("second-cancelled" as const), + ).pipe(Effect.forkScoped) + yield* Effect.sleep("10 millis") expect(KiloSessionPromptQueue.hasFollowup(sessionID)).toBe(true) - await expect( - Question.ask({ - sessionID, - questions: [ - { - header: "Continue?", - question: "Should I continue?", - options: [ - { label: "Yes", description: "Go" }, - { label: "No", description: "Stop" }, - ], - }, - ], - }), - ).rejects.toBeInstanceOf(Question.RejectedError) - expect(await Question.list()).toEqual([]) + const exit = yield* question.ask({ sessionID, questions: prompt }).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Question.RejectedError) + expect(yield* question.list()).toEqual([]) release.resolve() - expect(await first).toBe("first") - expect(await second).toBe("second") - }, - }) - }) + expect(yield* Fiber.join(first)).toBe("first") + expect(yield* Fiber.join(second)).toBe("second") + }), + { git: true }, + ) }) diff --git a/packages/opencode/test/kilocode/session-compaction-cap.test.ts b/packages/opencode/test/kilocode/session-compaction-cap.test.ts index cbd4fc6d7fd..239e63029a2 100644 --- a/packages/opencode/test/kilocode/session-compaction-cap.test.ts +++ b/packages/opencode/test/kilocode/session-compaction-cap.test.ts @@ -165,6 +165,7 @@ function makeHttp() { Layer.provideMerge(proc), Layer.provideMerge(registry), Layer.provideMerge(trunc), + Layer.provideMerge(question), // kilocode_change - SessionPrompt now dismisses questions via its service dependency Layer.provide(Instruction.defaultLayer), Layer.provide(SystemPrompt.defaultLayer), Layer.provideMerge(deps), diff --git a/packages/opencode/test/kilocode/session-prompt-compaction-safety.test.ts b/packages/opencode/test/kilocode/session-prompt-compaction-safety.test.ts index a472954f378..20c1c3d03c6 100644 --- a/packages/opencode/test/kilocode/session-prompt-compaction-safety.test.ts +++ b/packages/opencode/test/kilocode/session-prompt-compaction-safety.test.ts @@ -158,6 +158,7 @@ function makeHttp() { Layer.provideMerge(proc), Layer.provideMerge(registry), Layer.provideMerge(trunc), + Layer.provideMerge(question), // kilocode_change - SessionPrompt now dismisses questions via its service dependency Layer.provide(Instruction.defaultLayer), Layer.provide(SystemPrompt.defaultLayer), Layer.provideMerge(deps), diff --git a/packages/opencode/test/kilocode/session-prompt-queue.test.ts b/packages/opencode/test/kilocode/session-prompt-queue.test.ts index c07663aff97..267901a9e56 100644 --- a/packages/opencode/test/kilocode/session-prompt-queue.test.ts +++ b/packages/opencode/test/kilocode/session-prompt-queue.test.ts @@ -4,7 +4,6 @@ import { Effect } from "effect" import { Bus } from "../../src/bus" import { KiloSessionPromptQueue } from "@/kilocode/session/prompt-queue" import { Suggestion } from "../../src/kilocode/suggestion" -import { Question } from "../../src/question" import { ModelID, ProviderID } from "../../src/provider/schema" import { WithInstance } from "../../src/project/with-instance" import { Session } from "../../src/session/session" @@ -628,59 +627,6 @@ describe("session prompt queue", () => { }) }) - test("new prompt dismisses a pending question", async () => { - const asked = Promise.withResolvers() - const rejected = Promise.withResolvers() - await using tmp = await tmpdir({ git: true }) - - await WithInstance.provide({ - directory: tmp.path, - fn: async () => { - const session = await Session.create({ title: "Question unblock regression" }) - const offAsked = Bus.subscribe(Question.Event.Asked, (event) => { - if (event.properties.sessionID === session.id) asked.resolve() - }) - const offRejected = Bus.subscribe(Question.Event.Rejected, (event) => { - if (event.properties.sessionID === session.id) rejected.resolve() - }) - - try { - const pending = Question.ask({ - sessionID: session.id, - questions: [ - { - header: "Continue?", - question: "Should I continue?", - options: [ - { label: "Yes", description: "Go ahead" }, - { label: "No", description: "Stop" }, - ], - }, - ], - }).catch((err) => { - if (err instanceof Question.RejectedError) return "rejected" - throw err - }) - - await asked.promise - await SessionPrompt.prompt({ - sessionID: session.id, - agent: "code", - parts: [{ type: "text", text: "replacement prompt" }], - noReply: true, - }) - await rejected.promise - - expect(await pending).toBe("rejected") - expect(await Question.list()).toEqual([]) - } finally { - offAsked() - offRejected() - } - }, - }) - }) - test("auto-dismisses a suggestion shown after a queued prompt", async () => { // Reverse ordering of the "new prompt dismisses a pending suggestion" test: // queue the follow-up first, then open the blocker. Suggestion.show must see @@ -746,70 +692,4 @@ describe("session prompt queue", () => { }) }) - test("auto-dismisses a question shown after a queued prompt", async () => { - await using tmp = await tmpdir({ git: true }) - await WithInstance.provide({ - directory: tmp.path, - fn: async () => { - const sessionID = SessionID.make("ses_auto_question") - const started = Promise.withResolvers() - const release = Promise.withResolvers() - - const first = Effect.runPromise( - KiloSessionPromptQueue.enqueue( - sessionID, - MessageID.make("message_auto_q_1"), - Effect.gen(function* () { - started.resolve() - yield* Effect.promise(() => release.promise) - return "first" as const - }), - Effect.succeed("first-cancelled" as const), - ), - ) - await started.promise - - const second = Effect.runPromise( - KiloSessionPromptQueue.enqueue( - sessionID, - MessageID.make("message_auto_q_2"), - Effect.succeed("second" as const), - Effect.succeed("second-cancelled" as const), - ), - ) - await Bun.sleep(10) - expect(KiloSessionPromptQueue.hasFollowup(sessionID)).toBe(true) - - let asked = 0 - const offAsked = Bus.subscribe(Question.Event.Asked, (event) => { - if (event.properties.sessionID === sessionID) asked++ - }) - try { - await expect( - Question.ask({ - sessionID, - questions: [ - { - header: "Continue?", - question: "Should I continue?", - options: [ - { label: "Yes", description: "Go ahead" }, - { label: "No", description: "Stop" }, - ], - }, - ], - }), - ).rejects.toBeInstanceOf(Question.RejectedError) - } finally { - offAsked() - } - expect(asked).toBe(0) - expect(await Question.list()).toEqual([]) - - release.resolve() - expect(await first).toBe("first") - expect(await second).toBe("second") - }, - }) - }) }) diff --git a/packages/opencode/test/kilocode/sessions/remote-sender.test.ts b/packages/opencode/test/kilocode/sessions/remote-sender.test.ts index 734d974dc5a..a5d5d254390 100644 --- a/packages/opencode/test/kilocode/sessions/remote-sender.test.ts +++ b/packages/opencode/test/kilocode/sessions/remote-sender.test.ts @@ -6,6 +6,7 @@ import type { RemoteWS } from "../../../src/kilo-sessions/remote-ws" import type { RemoteProtocol } from "../../../src/kilo-sessions/remote-protocol" import { SessionPrompt } from "../../../src/session/prompt" import { Question } from "../../../src/question" +import { QuestionID } from "../../../src/question/schema" import { Permission } from "../../../src/permission" import { PermissionID } from "../../../src/permission/schema" import { Suggestion } from "../../../src/kilocode/suggestion" // kilocode_change @@ -55,6 +56,14 @@ function permissions(items: Permission.Request[] = []) { } } +function questions(items: Question.Request[] = []) { + return { + list: async () => items, + reply: async (_input: Parameters[0]) => {}, + reject: async (_requestID: QuestionID) => {}, + } +} + // kilocode_change start afterEach(() => { mock.restore() @@ -444,15 +453,18 @@ describe("RemoteSender", () => { test("question_reply sends response after work completes", async () => { const { conn, sent } = fakeConn() - let provideCalled = false + const calls: Parameters[0][] = [] const sender = RemoteSender.create({ conn, directory: "/tmp/test", log: nolog, subscribe: fakeBus().subscribe, - provide: async () => { - provideCalled = true - return {} as any + provide: async (input: { directory: string; init?: Effect.Effect; fn: () => R }) => input.fn(), + question: { + ...questions(), + reply: async (input) => { + calls.push(input) + }, }, }) @@ -463,12 +475,12 @@ describe("RemoteSender", () => { data: { requestID: "r1", answers: [["yes"]] }, }) - // Response not sent synchronously — waits for provide to finish + // Response not sent synchronously - waits for provide to finish. expect(sent).toHaveLength(0) await new Promise((r) => setTimeout(r, 10)) - expect(provideCalled).toBe(true) + expect(calls).toEqual([{ requestID: QuestionID.make("r1"), answers: [["yes"]] }]) expect(sent).toHaveLength(1) expect(sent[0]).toEqual({ type: "response", id: "req_q", result: {} }) }) @@ -511,8 +523,12 @@ describe("RemoteSender", () => { directory: "/tmp/test", log: nolog, subscribe: fakeBus().subscribe, - provide: async () => { - throw new Error("boom") + provide: async (input: { directory: string; init?: Effect.Effect; fn: () => R }) => input.fn(), + question: { + ...questions(), + reply: async () => { + throw new Error("boom") + }, }, }) @@ -531,6 +547,37 @@ describe("RemoteSender", () => { expect(sent[0].error).toContain("boom") }) + test("question_reply reports unknown request errors", async () => { + const { conn, sent } = fakeConn() + const sender = RemoteSender.create({ + conn, + directory: "/tmp/test", + log: nolog, + subscribe: fakeBus().subscribe, + provide: async (input: { directory: string; init?: Effect.Effect; fn: () => R }) => input.fn(), + question: { + ...questions(), + reply: async (input) => { + throw new Question.NotFoundError({ requestID: input.requestID }) + }, + }, + }) + + sender.handle({ + type: "command", + id: "req_q_missing", + command: "question_reply", + data: { requestID: "missing", answers: [["yes"]] }, + }) + + await new Promise((r) => setTimeout(r, 10)) + + expect(sent).toHaveLength(1) + expect(sent[0].type).toBe("response") + expect(sent[0].id).toBe("req_q_missing") + expect(sent[0].error).toContain("Question.NotFoundError") + }) + test("suggestion_accept sends response after work completes", async () => { const { conn, sent } = fakeConn() const accept = spyOn(Suggestion, "accept").mockResolvedValue(true) @@ -578,15 +625,18 @@ describe("RemoteSender", () => { test("question_reject sends response after work completes", async () => { const { conn, sent } = fakeConn() - let provideCalled = false + const calls: QuestionID[] = [] const sender = RemoteSender.create({ conn, directory: "/tmp/test", log: nolog, subscribe: fakeBus().subscribe, - provide: async () => { - provideCalled = true - return {} as any + provide: async (input: { directory: string; init?: Effect.Effect; fn: () => R }) => input.fn(), + question: { + ...questions(), + reject: async (requestID) => { + calls.push(requestID) + }, }, }) @@ -599,11 +649,42 @@ describe("RemoteSender", () => { await new Promise((r) => setTimeout(r, 10)) - expect(provideCalled).toBe(true) + expect(calls).toEqual([QuestionID.make("r1")]) expect(sent).toHaveLength(1) expect(sent[0]).toEqual({ type: "response", id: "req_qr", result: {} }) }) + test("question_reject reports unknown request errors", async () => { + const { conn, sent } = fakeConn() + const sender = RemoteSender.create({ + conn, + directory: "/tmp/test", + log: nolog, + subscribe: fakeBus().subscribe, + provide: async (input: { directory: string; init?: Effect.Effect; fn: () => R }) => input.fn(), + question: { + ...questions(), + reject: async (requestID) => { + throw new Question.NotFoundError({ requestID }) + }, + }, + }) + + sender.handle({ + type: "command", + id: "req_qr_missing", + command: "question_reject", + data: { requestID: "missing" }, + }) + + await new Promise((r) => setTimeout(r, 10)) + + expect(sent).toHaveLength(1) + expect(sent[0].type).toBe("response") + expect(sent[0].id).toBe("req_qr_missing") + expect(sent[0].error).toContain("Question.NotFoundError") + }) + test("question_reject with invalid data sends error response", () => { const { conn, sent } = fakeConn() const sender = RemoteSender.create({ @@ -900,10 +981,6 @@ describe("RemoteSender", () => { const bus = fakeBus() spyOn(Suggestion, "list").mockResolvedValue([]) - spyOn(Question, "list").mockResolvedValue([ - { id: "question_1", sessionID: "ses_target", questions: [{ type: "text", text: "Continue?" }] } as any, - { id: "question_2", sessionID: "ses_other", questions: [{ type: "text", text: "Unrelated?" }] } as any, - ]) const sender = RemoteSender.create({ conn, @@ -912,6 +989,10 @@ describe("RemoteSender", () => { subscribe: bus.subscribe, provide: async (input: any) => input.fn(), permission: permissions(), + question: questions([ + { id: "question_1", sessionID: "ses_target", questions: [{ type: "text", text: "Continue?" }] } as any, + { id: "question_2", sessionID: "ses_other", questions: [{ type: "text", text: "Unrelated?" }] } as any, + ]), }) sender.handle({ type: "subscribe", sessionId: "ses_target" }) @@ -932,7 +1013,6 @@ describe("RemoteSender", () => { const bus = fakeBus() spyOn(Suggestion, "list").mockResolvedValue([]) - spyOn(Question, "list").mockResolvedValue([]) const sender = RemoteSender.create({ conn, @@ -940,6 +1020,7 @@ describe("RemoteSender", () => { log: nolog, subscribe: bus.subscribe, provide: async (input: any) => input.fn(), + question: questions(), permission: permissions([ { id: "permission_1", @@ -987,7 +1068,6 @@ describe("RemoteSender", () => { spyOn(Suggestion, "list").mockResolvedValue([ { id: "sug_1", sessionID: "ses_other", text: "Review?", actions: [] } as any, ]) - spyOn(Question, "list").mockResolvedValue([{ id: "question_1", sessionID: "ses_other", questions: [] } as any]) const sender = RemoteSender.create({ conn, @@ -995,6 +1075,7 @@ describe("RemoteSender", () => { log: nolog, subscribe: bus.subscribe, provide: async (input: any) => input.fn(), + question: questions([{ id: "question_1", sessionID: "ses_other", questions: [] } as any]), permission: permissions([ { id: "permission_1", @@ -1032,7 +1113,6 @@ describe("RemoteSender", () => { actions: [{ label: "Skip", prompt: "skip" }], } as any, ]) - spyOn(Question, "list").mockResolvedValue([]) const sender = RemoteSender.create({ conn, @@ -1041,6 +1121,7 @@ describe("RemoteSender", () => { subscribe: bus.subscribe, provide: async (input: any) => input.fn(), permission: permissions(), + question: questions(), }) sender.handle({ type: "subscribe", sessionId: "ses_target" }) diff --git a/packages/opencode/test/question/question.test.ts b/packages/opencode/test/question/question.test.ts index eb099277619..d0ecca62006 100644 --- a/packages/opencode/test/question/question.test.ts +++ b/packages/opencode/test/question/question.test.ts @@ -1,11 +1,11 @@ -import { afterEach, expect, test } from "bun:test" +import { afterEach, expect } from "bun:test" import { Cause, Effect, Exit, Fiber, Layer } from "effect" import { Question } from "../../src/question" import { Instance } from "../../src/project/instance" import { WithInstance } from "../../src/project/with-instance" import { InstanceRuntime } from "../../src/project/instance-runtime" import { QuestionID } from "../../src/question/schema" -import { disposeAllInstances, provideInstance, reloadTestInstance, tmpdir, tmpdirScoped } from "../fixture/fixture" +import { disposeAllInstances, provideInstance, reloadTestInstance, tmpdirScoped } from "../fixture/fixture" import { SessionID } from "../../src/session/schema" import { testEffect } from "../lib/effect" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" @@ -15,6 +15,7 @@ const it = testEffect(Layer.mergeAll(Question.defaultLayer, CrossSpawnSpawner.de const askEffect = Effect.fn("QuestionTest.ask")(function* (input: { sessionID: SessionID questions: ReadonlyArray + blocking?: boolean // kilocode_change tool?: Question.Tool }) { const question = yield* Question.Service @@ -110,12 +111,11 @@ it.instance( ) // kilocode_change start - review follow-up uses non-blocking question prompts -test("ask - preserves blocking flag", async () => { - await using tmp = await tmpdir({ git: true }) - await WithInstance.provide({ - directory: tmp.path, - fn: async () => { - const askPromise = Question.ask({ +it.instance( + "ask - preserves blocking flag", + () => + Effect.gen(function* () { + const fiber = yield* askEffect({ sessionID: SessionID.make("ses_test"), blocking: false, questions: [ @@ -125,16 +125,18 @@ test("ask - preserves blocking flag", async () => { options: [{ label: "Start", description: "Run review" }], }, ], - }) + }).pipe(Effect.forkScoped) - const pending = await Question.list() + const pending = yield* waitForPending(1) expect(pending[0]?.blocking).toBe(false) - await Question.reject(pending[0].id) - await expect(askPromise).rejects.toBeInstanceOf(Question.RejectedError) - }, - }) -}) + yield* rejectEffect(pending[0].id) + const exit = yield* Fiber.await(fiber) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Question.RejectedError) + }), + { git: true }, +) // kilocode_change end // reply tests @@ -206,11 +208,16 @@ it.instance( ) it.instance( - "reply - does nothing for unknown requestID", + "reply - fails for unknown requestID", () => - replyEffect({ - requestID: QuestionID.make("que_unknown"), - answers: [["Option 1"]], + Effect.gen(function* () { + const id = QuestionID.make("que_unknown") + const exit = yield* replyEffect({ requestID: id, answers: [["Option 1"]] }).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) + if (!Exit.isFailure(exit)) return + const err = Cause.squash(exit.cause) + expect(err).toBeInstanceOf(Question.NotFoundError) + if (err instanceof Question.NotFoundError) expect(err.requestID).toBe(id) }), { git: true }, ) @@ -275,9 +282,20 @@ it.instance( { git: true }, ) -it.instance("reject - does nothing for unknown requestID", () => rejectEffect(QuestionID.make("que_unknown")), { - git: true, -}) +it.instance( + "reject - fails for unknown requestID", + () => + Effect.gen(function* () { + const id = QuestionID.make("que_unknown") + const exit = yield* rejectEffect(id).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) + if (!Exit.isFailure(exit)) return + const err = Cause.squash(exit.cause) + expect(err).toBeInstanceOf(Question.NotFoundError) + if (err instanceof Question.NotFoundError) expect(err.requestID).toBe(id) + }), + { git: true }, +) // multiple questions tests diff --git a/packages/opencode/test/session/prompt.test.ts b/packages/opencode/test/session/prompt.test.ts index ea5aa9255ce..fd6655dd545 100644 --- a/packages/opencode/test/session/prompt.test.ts +++ b/packages/opencode/test/session/prompt.test.ts @@ -213,6 +213,7 @@ function makeHttp() { Layer.provideMerge(proc), Layer.provideMerge(registry), Layer.provideMerge(trunc), + Layer.provideMerge(question), // kilocode_change - SessionPrompt now dismisses questions via its service dependency Layer.provide(Instruction.defaultLayer), Layer.provide(SystemPrompt.defaultLayer), Layer.provideMerge(deps), @@ -396,6 +397,51 @@ it.live("loop calls LLM and returns assistant message", () => ), ) +// kilocode_change start - replacement prompts unblock pending Question service requests +it.live("new prompt dismisses a pending question", () => + provideTmpdirServer( + Effect.fnUntraced(function* () { + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + const question = yield* Question.Service + const chat = yield* sessions.create({ title: "Question unblock regression" }) + const pending = yield* question + .ask({ + sessionID: chat.id, + questions: [ + { + header: "Continue?", + question: "Should I continue?", + options: [ + { label: "Yes", description: "Go ahead" }, + { label: "No", description: "Stop" }, + ], + }, + ], + }) + .pipe(Effect.forkScoped) + yield* waitFor( + "pending question", + question.list().pipe(Effect.map((items) => items.find((item) => item.sessionID === chat.id))), + ) + + yield* prompt.prompt({ + sessionID: chat.id, + agent: "build", + parts: [{ type: "text", text: "replacement prompt" }], + noReply: true, + }) + + const exit = yield* Fiber.await(pending) + expect(Exit.isFailure(exit)).toBe(true) + if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Question.RejectedError) + expect(yield* question.list()).toEqual([]) + }), + { git: true, config: providerCfg }, + ), +) +// kilocode_change end + it.live("prompt emits v2 prompted and synthetic events", () => provideTmpdirServer( Effect.fnUntraced(function* () { diff --git a/packages/opencode/test/session/snapshot-tool-race.test.ts b/packages/opencode/test/session/snapshot-tool-race.test.ts index 26cd1111824..2949abd9727 100644 --- a/packages/opencode/test/session/snapshot-tool-race.test.ts +++ b/packages/opencode/test/session/snapshot-tool-race.test.ts @@ -150,6 +150,7 @@ function makeHttp() { Layer.provideMerge(proc), Layer.provideMerge(registry), Layer.provideMerge(trunc), + Layer.provideMerge(question), // kilocode_change - SessionPrompt now dismisses questions via its service dependency Layer.provide(Instruction.defaultLayer), Layer.provide(SystemPrompt.defaultLayer), Layer.provideMerge(deps), diff --git a/script/check-opencode-promise-facades.ts b/script/check-opencode-promise-facades.ts index 2e0b10778ad..748bfdf1134 100644 --- a/script/check-opencode-promise-facades.ts +++ b/script/check-opencode-promise-facades.ts @@ -23,7 +23,6 @@ const allow: Record = { "bus/index.ts": "core bus callback and synchronous runtime boundary", "cli/cmd/tui/config/tui.ts": "separately tracked TUI config facade", "installation/index.ts": "existing installation facade outside #10655", - "question/index.ts": "transitional facade deferred for upstream reconciliation in #10655", "session/compaction.ts": "existing compaction facade outside #10655", "session/prompt.ts": "transitional facade tracked by #10655", "session/session.ts": "transitional facade tracked by #10655",