From 7c78e44d0175ee480bf0d02ef6286fe67bad770e Mon Sep 17 00:00:00 2001 From: Yuhan Lei Date: Thu, 28 May 2026 09:37:35 +0800 Subject: [PATCH 1/2] fix(session): repair stale paginated question blockers --- .../opencode/src/server/instance/session.ts | 2 +- packages/opencode/src/session/session.ts | 26 +++ .../test/server/session-messages.test.ts | 156 ++++++++++++++++++ 3 files changed, 183 insertions(+), 1 deletion(-) diff --git a/packages/opencode/src/server/instance/session.ts b/packages/opencode/src/server/instance/session.ts index 7c8df281f..8db48b73a 100644 --- a/packages/opencode/src/server/instance/session.ts +++ b/packages/opencode/src/server/instance/session.ts @@ -1130,7 +1130,7 @@ export const SessionRoutes = lazy(() => return c.json(messages) } - const page = await MessageV2.page({ + const page = await Session.messagesPage({ sessionID, limit: query.limit, before: query.before, diff --git a/packages/opencode/src/session/session.ts b/packages/opencode/src/session/session.ts index e4f267f93..d818dc99e 100644 --- a/packages/opencode/src/session/session.ts +++ b/packages/opencode/src/session/session.ts @@ -338,6 +338,11 @@ export const SetRevertInput = z.object({ summary: Info.shape.summary, }) export const MessagesInput = z.object({ sessionID: SessionID.zod, limit: z.number().optional() }) +export const MessagesPageInput = z.object({ + sessionID: SessionID.zod, + limit: z.number(), + before: z.string().optional(), +}) const AbsoluteDirectory = z .string() .min(1, "Expected an absolute directory path") @@ -524,6 +529,11 @@ export interface Interface { readonly findActiveWorktreeBinding: (directory: string) => Effect.Effect readonly diff: (sessionID: SessionID) => Effect.Effect readonly messages: (input: { sessionID: SessionID; limit?: number }) => Effect.Effect + readonly messagesPage: (input: { + sessionID: SessionID + limit: number + before?: string + }) => Effect.Effect> readonly children: (parentID: SessionID) => Effect.Effect readonly remove: (sessionID: SessionID) => Effect.Effect readonly updateMessage: (msg: T) => Effect.Effect @@ -1032,6 +1042,20 @@ export const layer: Layer.Layer = return yield* terminalizeStaleExternalResultQuestions(items) }) + const messagesPage = Effect.fn("Session.messagesPage")(function* (input: { + sessionID: SessionID + limit: number + before?: string + }) { + const page = MessageV2.page({ + sessionID: input.sessionID, + limit: input.limit, + before: input.before, + }) + const items = yield* terminalizeStaleExternalResultQuestions(page.items) + return { ...page, items } + }) + const removeMessage = Effect.fn("Session.removeMessage")(function* (input: { sessionID: SessionID messageID: MessageID @@ -1096,6 +1120,7 @@ export const layer: Layer.Layer = findActiveWorktreeBinding, diff, messages, + messagesPage, children, remove, updateMessage, @@ -1125,6 +1150,7 @@ export const setTitle = fn(SetTitleInput, (input) => runPromise((svc) => svc.set export const setArchived = fn(SetArchivedInput, (input) => runPromise((svc) => svc.setArchived(input))) export const setPermission = fn(SetPermissionInput, (input) => runPromise((svc) => svc.setPermission(input))) export const messages = fn(MessagesInput, (input) => runPromise((svc) => svc.messages(input))) +export const messagesPage = fn(MessagesPageInput, (input) => runPromise((svc) => svc.messagesPage(input))) export const removePart = fn(RemovePartInput, (input) => runPromise((svc) => svc.removePart(input))) export const updateMessage = fn(MessageV2.Info, (input) => runPromise((svc) => svc.updateMessage(input))) export const updatePart = fn(MessageV2.Part, (input) => runPromise((svc) => svc.updatePart(input))) diff --git a/packages/opencode/test/server/session-messages.test.ts b/packages/opencode/test/server/session-messages.test.ts index 1901d2c8e..17b84ad27 100644 --- a/packages/opencode/test/server/session-messages.test.ts +++ b/packages/opencode/test/server/session-messages.test.ts @@ -5,6 +5,7 @@ import { Server } from "../../src/server/server" import { Session as SessionNs } from "../../src/session" import { MessageV2 } from "../../src/session/message-v2" import { MessageID, PartID, type SessionID } from "../../src/session/schema" +import { ExternalResult } from "../../src/tool/external-result" import { Log } from "../../src/util" import { tmpdir } from "../fixture/fixture" @@ -31,6 +32,7 @@ const svc = { } afterEach(async () => { + ExternalResult.__resetForTests() await Instance.disposeAll() }) @@ -72,7 +74,161 @@ async function fill(sessionID: SessionID, count: number, time = (i: number) => D return ids } +async function createRunningQuestionSession(directory: string, input?: { externalResultReady?: boolean }) { + const session = await svc.create({}) + const userID = MessageID.ascending() + await svc.updateMessage({ + id: userID, + sessionID: session.id, + role: "user", + time: { created: Date.now() }, + agent: "user", + model: { providerID: "test", modelID: "test" }, + tools: {}, + mode: "", + } as unknown as MessageV2.Info) + + const assistantID = MessageID.ascending() + await svc.updateMessage({ + id: assistantID, + sessionID: session.id, + role: "assistant", + parentID: userID, + time: { created: Date.now() }, + agent: "build", + mode: "build", + path: { cwd: directory, root: directory }, + cost: 0, + tokens: { + total: 0, + input: 0, + output: 0, + reasoning: 0, + cache: { read: 0, write: 0 }, + }, + modelID: "test", + providerID: "test", + } as unknown as MessageV2.Info) + + const partID = PartID.ascending() + const callID = "call_stale_question" + await svc.updatePart({ + id: partID, + sessionID: session.id, + messageID: assistantID, + type: "tool", + tool: "question", + callID, + state: { + status: "running", + input: { + questions: [ + { + question: "Continue?", + options: [{ label: "Yes" }, { label: "No" }], + }, + ], + }, + raw: "", + time: { start: Date.now() }, + metadata: { externalResultReady: input?.externalResultReady ?? true }, + }, + } as unknown as MessageV2.Part) + + return { session, assistantID, partID, callID } +} + +function expectToolPart(part: MessageV2.Part | undefined) { + expect(part?.type).toBe("tool") + if (part?.type !== "tool") throw new Error("expected tool part") + return part +} + describe("session messages endpoint", () => { + test("terminalizes stale running external-result questions in paginated responses", async () => { + await using tmp = await tmpdir({ git: true }) + ExternalResult.__resetForTests() + await withoutWatcher(() => + Instance.provide({ + directory: tmp.path, + fn: async () => { + const { session, assistantID, partID } = await createRunningQuestionSession(tmp.path) + const app = Server.Default().app + + const res = await app.request(`/session/${session.id}/message?limit=2`) + expect(res.status).toBe(200) + const body = (await res.json()) as MessageV2.WithParts[] + const part = expectToolPart(body.flatMap((msg) => msg.parts).find((item) => item.id === partID)) + expect(part.state.status).toBe("error") + if (part.state.status !== "error") throw new Error("expected error state") + expect(part.state.metadata?.interrupted).toBe(true) + expect(part.state.metadata?.stale_external_result).toBe(true) + + const persisted = expectToolPart( + MessageV2.get({ sessionID: session.id, messageID: assistantID }).parts.find( + (item) => item.id === partID, + ), + ) + expect(persisted.state.status).toBe("error") + + await svc.remove(session.id) + }, + }), + ) + }) + + test("preserves live pending external-result questions in paginated responses", async () => { + await using tmp = await tmpdir({ git: true }) + ExternalResult.__resetForTests() + await withoutWatcher(() => + Instance.provide({ + directory: tmp.path, + fn: async () => { + const { session, assistantID, partID, callID } = await createRunningQuestionSession(tmp.path) + await Effect.runPromise( + ExternalResult.register({ + sessionID: session.id, + messageID: assistantID, + callID, + inputSnapshot: { questions: ["q1"] }, + }), + ) + const app = Server.Default().app + + const res = await app.request(`/session/${session.id}/message?limit=2`) + expect(res.status).toBe(200) + const body = (await res.json()) as MessageV2.WithParts[] + const part = expectToolPart(body.flatMap((msg) => msg.parts).find((item) => item.id === partID)) + expect(part.state.status).toBe("running") + + await svc.remove(session.id) + }, + }), + ) + }) + + test("preserves unready external-result questions in paginated responses", async () => { + await using tmp = await tmpdir({ git: true }) + ExternalResult.__resetForTests() + await withoutWatcher(() => + Instance.provide({ + directory: tmp.path, + fn: async () => { + const { session, partID } = await createRunningQuestionSession(tmp.path, { externalResultReady: false }) + const app = Server.Default().app + + const res = await app.request(`/session/${session.id}/message?limit=2`) + expect(res.status).toBe(200) + const body = (await res.json()) as MessageV2.WithParts[] + const part = expectToolPart(body.flatMap((msg) => msg.parts).find((item) => item.id === partID)) + expect(part.state.status).toBe("running") + + await svc.remove(session.id) + }, + }), + ) + }) + test("returns cursor headers for older pages", async () => { await using tmp = await tmpdir({ git: true }) await withoutWatcher(() => From fdc7337bb37a96de675d10a8603a3f9357d18d7f Mon Sep 17 00:00:00 2001 From: Yuhan Lei Date: Thu, 28 May 2026 12:59:56 +0800 Subject: [PATCH 2/2] test(session): cover stale question cursor pagination --- .../test/server/session-messages.test.ts | 50 +++++++++++++++++-- 1 file changed, 46 insertions(+), 4 deletions(-) diff --git a/packages/opencode/test/server/session-messages.test.ts b/packages/opencode/test/server/session-messages.test.ts index 17b84ad27..4456442a4 100644 --- a/packages/opencode/test/server/session-messages.test.ts +++ b/packages/opencode/test/server/session-messages.test.ts @@ -74,14 +74,15 @@ async function fill(sessionID: SessionID, count: number, time = (i: number) => D return ids } -async function createRunningQuestionSession(directory: string, input?: { externalResultReady?: boolean }) { +async function createRunningQuestionSession(directory: string, input?: { externalResultReady?: boolean; time?: number }) { const session = await svc.create({}) const userID = MessageID.ascending() + const time = input?.time ?? Date.now() await svc.updateMessage({ id: userID, sessionID: session.id, role: "user", - time: { created: Date.now() }, + time: { created: time }, agent: "user", model: { providerID: "test", modelID: "test" }, tools: {}, @@ -94,7 +95,7 @@ async function createRunningQuestionSession(directory: string, input?: { externa sessionID: session.id, role: "assistant", parentID: userID, - time: { created: Date.now() }, + time: { created: time + 1 }, agent: "build", mode: "build", path: { cwd: directory, root: directory }, @@ -130,7 +131,7 @@ async function createRunningQuestionSession(directory: string, input?: { externa ], }, raw: "", - time: { start: Date.now() }, + time: { start: time + 2 }, metadata: { externalResultReady: input?.externalResultReady ?? true }, }, } as unknown as MessageV2.Part) @@ -177,6 +178,47 @@ describe("session messages endpoint", () => { ) }) + test("terminalizes stale running external-result questions when fetched from an older cursor page", async () => { + await using tmp = await tmpdir({ git: true }) + ExternalResult.__resetForTests() + await withoutWatcher(() => + Instance.provide({ + directory: tmp.path, + fn: async () => { + const time = Date.now() + const { session, assistantID, partID } = await createRunningQuestionSession(tmp.path, { time }) + await fill(session.id, 3, (i) => time + 10 + i) + const app = Server.Default().app + + const latest = await app.request(`/session/${session.id}/message?limit=2`) + expect(latest.status).toBe(200) + const cursor = latest.headers.get("x-next-cursor") + expect(cursor).toBeTruthy() + + const older = await app.request( + `/session/${session.id}/message?limit=2&before=${encodeURIComponent(cursor!)}`, + ) + expect(older.status).toBe(200) + const body = (await older.json()) as MessageV2.WithParts[] + const part = expectToolPart(body.flatMap((msg) => msg.parts).find((item) => item.id === partID)) + expect(part.state.status).toBe("error") + if (part.state.status !== "error") throw new Error("expected error state") + expect(part.state.metadata?.interrupted).toBe(true) + expect(part.state.metadata?.stale_external_result).toBe(true) + + const persisted = expectToolPart( + MessageV2.get({ sessionID: session.id, messageID: assistantID }).parts.find( + (item) => item.id === partID, + ), + ) + expect(persisted.state.status).toBe("error") + + await svc.remove(session.id) + }, + }), + ) + }) + test("preserves live pending external-result questions in paginated responses", async () => { await using tmp = await tmpdir({ git: true }) ExternalResult.__resetForTests()