mirror of
https://github.com/anomalyco/opencode.git
synced 2026-09-28 13:33:35 +08:00
fix: mini reconnect (#49301)
A stalled catalog refresh kept reconnection from completing and prevented new prompts. Allow the stream to reconnect independently of catalog state.
This commit is contained in:
@@ -411,6 +411,7 @@ export function RunFooterView(props: RunFooterViewProps) {
|
||||
if (notice()) return notice()
|
||||
if (!footerDetails()) return shell() ? "Shell" : ""
|
||||
if (busy()) {
|
||||
if (stateStatus() === "reconnecting") return "reconnecting"
|
||||
return interruptLabel() ? `${interruptLabel()} stop` : "Running"
|
||||
}
|
||||
return stateStatus() || (shell() ? "Shell" : "")
|
||||
|
||||
@@ -520,6 +520,14 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
const abortReady = () => readyReject(new Error("Mini closed before the event stream connected"))
|
||||
controller.signal.addEventListener("abort", abortReady, { once: true })
|
||||
const offFooterClose = input.footer.onClose(() => controller.abort())
|
||||
const waitUntilConnected = async (signal?: AbortSignal) => {
|
||||
const abort = signal ? AbortSignal.any([signal, controller.signal]) : controller.signal
|
||||
while (!state.connected) {
|
||||
if (state.closed || controller.signal.aborted || input.footer.isClosed || signal?.aborted)
|
||||
throw new Error("Event stream aborted")
|
||||
await wait(25, abort)
|
||||
}
|
||||
}
|
||||
const current = (attempt: Attempt) =>
|
||||
!state.closed &&
|
||||
!controller.signal.aborted &&
|
||||
@@ -991,7 +999,7 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
phase: state.rootActive ? "running" : "idle",
|
||||
status: state.rootActive ? "assistant responding" : blockerStatus(state.view),
|
||||
})
|
||||
if (!state.rootActive) await input.footer.idle()
|
||||
if (!state.rootActive && !next.reconnect) await input.footer.idle()
|
||||
if (!current(attempt)) return
|
||||
}
|
||||
|
||||
@@ -1462,13 +1470,6 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
})
|
||||
return task
|
||||
}
|
||||
const settleCatalog = async (attempt: Attempt) => {
|
||||
while (current(attempt)) {
|
||||
const refreshes = catalogRefreshes.get(attempt.generation)
|
||||
if (!refreshes || refreshes.size === 0) return
|
||||
await Promise.all(refreshes)
|
||||
}
|
||||
}
|
||||
const settleCatalogRefreshes = async () => {
|
||||
while (catalogRefreshes.size > 0)
|
||||
await Promise.all([...catalogRefreshes.values()].flatMap((refreshes) => [...refreshes]))
|
||||
@@ -1506,18 +1507,14 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
),
|
||||
consume,
|
||||
])
|
||||
await Promise.race([refreshCatalog(attempt), consume])
|
||||
if (!current(attempt)) throw new Error("Event stream disconnected")
|
||||
state.initial = false
|
||||
do {
|
||||
for (const event of buffered.splice(0)) apply(attempt, event)
|
||||
await Promise.race([subagents.ready(), consume])
|
||||
await Promise.race([settleCatalog(attempt), consume])
|
||||
} while (buffered.length > 0)
|
||||
for (const event of buffered.splice(0)) apply(attempt, event)
|
||||
if (!current(attempt)) throw new Error("Event stream disconnected")
|
||||
booting = false
|
||||
state.connected = true
|
||||
readyResolve()
|
||||
void refreshCatalog(attempt)
|
||||
await consume
|
||||
} finally {
|
||||
connection.abort()
|
||||
@@ -1563,7 +1560,7 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
|
||||
const runShellTurn = async (next: SessionTurnInput) => {
|
||||
if (state.wait || state.shellWait) throw new Error("prompt already running")
|
||||
if (!state.connected) throw new Error("Event stream is reconnecting")
|
||||
await waitUntilConnected(next.signal)
|
||||
const client = sdk
|
||||
const abort = new AbortController()
|
||||
const onAbort = () => abort.abort()
|
||||
@@ -1796,7 +1793,7 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
return {
|
||||
async admitPromptTurn(next, delivery) {
|
||||
if (next.prompt.mode === "shell") throw new Error("This prompt cannot be queued")
|
||||
if (!state.connected) throw new Error("Event stream is reconnecting")
|
||||
await waitUntilConnected(next.signal)
|
||||
const client = sdk
|
||||
if (!next.prompt.command && next.agent)
|
||||
await client.session.switchAgent({ sessionID: input.sessionID, agent: next.agent }, { signal: next.signal })
|
||||
@@ -1821,7 +1818,7 @@ export async function createSessionTransport(input: StreamInput): Promise<Sessio
|
||||
return
|
||||
}
|
||||
if (state.wait || state.shellWait) throw new Error("prompt already running")
|
||||
if (!state.connected) throw new Error("Event stream is reconnecting")
|
||||
await waitUntilConnected(next.signal)
|
||||
const client = sdk
|
||||
const messageID = next.prompt.messageID
|
||||
if (!messageID) throw new Error("Prompt message ID is required")
|
||||
|
||||
@@ -1847,70 +1847,6 @@ describe("V2 mini transport", () => {
|
||||
|
||||
firstEvents.close()
|
||||
while (!replacementHydrating) await Bun.sleep(0)
|
||||
await expect(
|
||||
transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_blocked", text: "blocked", parts: [] },
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
}),
|
||||
).rejects.toThrow("Event stream is reconnecting")
|
||||
secondEvents.push({
|
||||
id: "evt_buffered_text",
|
||||
created: 2,
|
||||
type: "session.text.delta",
|
||||
data: {
|
||||
sessionID: "ses_1",
|
||||
assistantMessageID: "msg_assistant",
|
||||
ordinal: 0,
|
||||
delta: " replacement",
|
||||
},
|
||||
})
|
||||
let resized = false
|
||||
const resize = transport.replayOnResize({
|
||||
localRows: () => [],
|
||||
reset: async () => {
|
||||
resized = true
|
||||
},
|
||||
})
|
||||
releaseHydration()
|
||||
while (
|
||||
!ui.events.some(
|
||||
(event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === "frm_child",
|
||||
)
|
||||
)
|
||||
await Bun.sleep(0)
|
||||
while (refreshes < 2) await Bun.sleep(0)
|
||||
await resize
|
||||
expect(resized).toBe(false)
|
||||
await expect(
|
||||
transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_catalog_blocked", text: "blocked", parts: [] },
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
}),
|
||||
).rejects.toThrow("Event stream is reconnecting")
|
||||
releaseCatalog()
|
||||
await Bun.sleep(0)
|
||||
|
||||
expect(current).toEqual([second])
|
||||
expect(first.event.subscribe).toHaveBeenCalledTimes(1)
|
||||
expect(second.event.subscribe).toHaveBeenCalledTimes(1)
|
||||
expect(second.session.list).toHaveBeenCalled()
|
||||
expect(second.session.form.list).toHaveBeenCalledWith(
|
||||
{ sessionID: "ses_child" },
|
||||
{ signal: expect.any(AbortSignal) },
|
||||
)
|
||||
expect(ui.commits.filter((commit) => commit.messageID === "msg_assistant").map((commit) => commit.text)).toEqual([
|
||||
"partial",
|
||||
" replacement",
|
||||
])
|
||||
|
||||
const prompt = spyOn(second.session, "prompt").mockImplementation((request) => {
|
||||
queueMicrotask(() => {
|
||||
secondEvents.push({
|
||||
@@ -1930,7 +1866,7 @@ describe("V2 mini transport", () => {
|
||||
})
|
||||
return ok({ data: promptAdmission(request) }) as never
|
||||
})
|
||||
await transport.runPromptTurn({
|
||||
const queued = transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
@@ -1938,6 +1874,42 @@ describe("V2 mini transport", () => {
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
})
|
||||
await Bun.sleep(0)
|
||||
expect(prompt).not.toHaveBeenCalled()
|
||||
secondEvents.push({
|
||||
id: "evt_buffered_text",
|
||||
created: 2,
|
||||
type: "session.text.delta",
|
||||
data: {
|
||||
sessionID: "ses_1",
|
||||
assistantMessageID: "msg_assistant",
|
||||
ordinal: 0,
|
||||
delta: " replacement",
|
||||
},
|
||||
})
|
||||
releaseHydration()
|
||||
while (
|
||||
!ui.events.some(
|
||||
(event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === "frm_child",
|
||||
)
|
||||
)
|
||||
await Bun.sleep(0)
|
||||
await queued
|
||||
while (refreshes < 2) await Bun.sleep(0)
|
||||
releaseCatalog()
|
||||
|
||||
expect(current).toEqual([second])
|
||||
expect(first.event.subscribe).toHaveBeenCalledTimes(1)
|
||||
expect(second.event.subscribe).toHaveBeenCalledTimes(1)
|
||||
expect(second.session.list).toHaveBeenCalled()
|
||||
expect(second.session.form.list).toHaveBeenCalledWith(
|
||||
{ sessionID: "ses_child" },
|
||||
{ signal: expect.any(AbortSignal) },
|
||||
)
|
||||
expect(ui.commits.filter((commit) => commit.messageID === "msg_assistant").map((commit) => commit.text)).toEqual([
|
||||
"partial",
|
||||
" replacement",
|
||||
])
|
||||
const interrupt = spyOn(second.session, "interrupt").mockImplementation(() => ok({ interrupted: true }))
|
||||
await transport.interruptActiveTurn()
|
||||
|
||||
@@ -1948,6 +1920,73 @@ describe("V2 mini transport", () => {
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("sends a prompt after the event stream reconnects", async () => {
|
||||
const first = feed()
|
||||
const second = feed()
|
||||
first.push(connected("evt_connected_1"))
|
||||
second.push(connected("evt_connected_2"))
|
||||
const client = sdk({ streams: [first, second] })
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
first.close()
|
||||
while (!ui.events.some((event) => event.type === "stream.patch" && event.patch.status === "reconnecting"))
|
||||
await Bun.sleep(0)
|
||||
const prompt = spyOn(client.session, "prompt").mockImplementation(
|
||||
(request) => ok({ data: promptAdmission(request) }) as never,
|
||||
)
|
||||
await transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_after_reconnect", text: "hello", parts: [] },
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
})
|
||||
expect(prompt).toHaveBeenCalled()
|
||||
expect(client.event.subscribe).toHaveBeenCalledTimes(2)
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("reconnects even when catalog refresh hangs", async () => {
|
||||
const first = feed()
|
||||
const second = feed()
|
||||
first.push(connected("evt_connected_1"))
|
||||
second.push(connected("evt_connected_2"))
|
||||
const client = sdk({ streams: [first, second] })
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
onCatalogRefresh: (signal) =>
|
||||
new Promise<void>((_resolve, reject) => {
|
||||
signal?.addEventListener("abort", () => reject(new Error("aborted")), { once: true })
|
||||
}),
|
||||
})
|
||||
first.close()
|
||||
while (!ui.events.some((event) => event.type === "stream.patch" && event.patch.status === "reconnecting"))
|
||||
await Bun.sleep(0)
|
||||
const prompt = spyOn(client.session, "prompt").mockImplementation(
|
||||
(request) => ok({ data: promptAdmission(request) }) as never,
|
||||
)
|
||||
await transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_after_hanging_catalog", text: "hello", parts: [] },
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
})
|
||||
expect(prompt).toHaveBeenCalled()
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("reconciles buffered deltas already present in a resize snapshot", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
|
||||
Reference in New Issue
Block a user