Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion packages/opencode/src/server/instance/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
26 changes: 26 additions & 0 deletions packages/opencode/src/session/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -524,6 +529,11 @@ export interface Interface {
readonly findActiveWorktreeBinding: (directory: string) => Effect.Effect<Info | undefined>
readonly diff: (sessionID: SessionID) => Effect.Effect<Snapshot.FileDiff[]>
readonly messages: (input: { sessionID: SessionID; limit?: number }) => Effect.Effect<MessageV2.WithParts[]>
readonly messagesPage: (input: {
sessionID: SessionID
limit: number
before?: string
}) => Effect.Effect<ReturnType<typeof MessageV2.page>>
readonly children: (parentID: SessionID) => Effect.Effect<Info[]>
readonly remove: (sessionID: SessionID) => Effect.Effect<void>
readonly updateMessage: <T extends MessageV2.Info>(msg: T) => Effect.Effect<T>
Expand Down Expand Up @@ -1032,6 +1042,20 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
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
Expand Down Expand Up @@ -1096,6 +1120,7 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
findActiveWorktreeBinding,
diff,
messages,
messagesPage,
children,
remove,
updateMessage,
Expand Down Expand Up @@ -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)))
Expand Down
198 changes: 198 additions & 0 deletions packages/opencode/test/server/session-messages.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand All @@ -31,6 +32,7 @@ const svc = {
}

afterEach(async () => {
ExternalResult.__resetForTests()
await Instance.disposeAll()
})

Expand Down Expand Up @@ -72,7 +74,203 @@ async function fill(sessionID: SessionID, count: number, time = (i: number) => D
return ids
}

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: time },
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: time + 1 },
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: time + 2 },
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("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()
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(() =>
Expand Down