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
61 changes: 57 additions & 4 deletions packages/opencode/src/session/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -972,11 +972,64 @@ export const layer: Layer.Layer<Service, never, Bus.Service | Storage.Service> =
.pipe(Effect.orElseSucceed((): Snapshot.FileDiff[] => []))
})

const terminalizeStaleExternalResultQuestions = Effect.fn("Session.terminalizeStaleExternalResultQuestions")(
function* (messages: MessageV2.WithParts[]) {
const now = Date.now()
const next: MessageV2.WithParts[] = []
for (const message of messages) {
let changed = false
const parts: MessageV2.Part[] = []
for (const part of message.parts) {
if (
part.type !== "tool" ||
part.tool !== "question" ||
part.state.status !== "running" ||
part.state.metadata?.externalResultReady !== true
) {
parts.push(part)
continue
}

const lookup = ExternalResult.lookup({
sessionID: part.sessionID,
messageID: part.messageID,
callID: part.callID,
})
if (lookup.state !== "not_found") {
parts.push(part)
continue
}

changed = true
parts.push(
yield* updatePart({
...part,
state: {
status: "error",
input: part.state.input,
error: "Question cancelled before the user answered it.",
reason: "shutdown",
metadata: {
...part.state.metadata,
interrupted: true,
stale_external_result: true,
},
time: { start: part.state.time.start, end: now },
},
}),
)
}
next.push(changed ? { ...message, parts } : message)
}
return next
},
Comment thread
Astro-Han marked this conversation as resolved.
)

const messages = Effect.fn("Session.messages")(function* (input: { sessionID: SessionID; limit?: number }) {
if (input.limit) {
return MessageV2.page({ sessionID: input.sessionID, limit: input.limit }).items
}
return Array.from(MessageV2.stream(input.sessionID)).reverse()
const items = input.limit
? MessageV2.page({ sessionID: input.sessionID, limit: input.limit }).items
: Array.from(MessageV2.stream(input.sessionID)).reverse()
return yield* terminalizeStaleExternalResultQuestions(items)
})

const removeMessage = Effect.fn("Session.removeMessage")(function* (input: {
Expand Down
155 changes: 155 additions & 0 deletions packages/opencode/test/session/session.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { Log } from "@opencode-ai/core/util/log"
import { Instance } from "../../src/project/instance"
import { MessageV2 } from "../../src/session/message-v2"
import { MessageID, PartID } from "../../src/session/schema"
import { ExternalResult } from "../../src/tool/external-result"
import { tmpdir } from "../fixture/fixture"
import { Database, eq } from "../../src/storage/db"
import { MessageTable, SessionTable } from "../../src/session/session.sql"
Expand Down Expand Up @@ -734,6 +735,160 @@ describe("step-finish token propagation via Bus event", () => {
})

describe("Session", () => {
async function createRunningQuestionSession(directory: string, input?: { externalResultReady?: boolean }) {
const session = await SessionNs.create({})
const userID = MessageID.ascending()
await SessionNs.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 SessionNs.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 SessionNs.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
}

test("messages terminalizes stale running external-result questions", async () => {
await using tmp = await tmpdir({ git: true })
ExternalResult.__resetForTests()

try {
await Instance.provide({
directory: tmp.path,
fn: async () => {
const { session, assistantID, partID } = await createRunningQuestionSession(tmp.path)

const messages = await SessionNs.messages({ sessionID: session.id })
const part = expectToolPart(messages.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 SessionNs.remove(session.id)
},
})
} finally {
ExternalResult.__resetForTests()
}
})

test("messages preserves live pending external-result questions", async () => {
await using tmp = await tmpdir({ git: true })
ExternalResult.__resetForTests()

try {
await 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 messages = await SessionNs.messages({ sessionID: session.id })
const part = expectToolPart(messages.flatMap((msg) => msg.parts).find((item) => item.id === partID))
expect(part.state.status).toBe("running")

await SessionNs.remove(session.id)
},
})
} finally {
ExternalResult.__resetForTests()
}
})

test("messages preserves unready external-result questions", async () => {
await using tmp = await tmpdir({ git: true })
ExternalResult.__resetForTests()

try {
await Instance.provide({
directory: tmp.path,
fn: async () => {
const { session, partID } = await createRunningQuestionSession(tmp.path, { externalResultReady: false })

const messages = await SessionNs.messages({ sessionID: session.id })
const part = expectToolPart(messages.flatMap((msg) => msg.parts).find((item) => item.id === partID))
expect(part.state.status).toBe("running")

await SessionNs.remove(session.id)
},
})
} finally {
ExternalResult.__resetForTests()
}
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.

test("remove works without an instance", async () => {
await using tmp = await tmpdir({ git: true })

Expand Down