fix(sessions): resync session status on SSE reconnect to clear stuck 'processing' after network drop (closes #123) (#124)
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
57
src/lib/session-status-reconcile.test.ts
Normal file
57
src/lib/session-status-reconcile.test.ts
Normal file
@@ -0,0 +1,57 @@
|
|||||||
|
import { test } from "node:test"
|
||||||
|
import assert from "node:assert/strict"
|
||||||
|
import { isSessionActuallyIdle } from "./session-status-reconcile.ts"
|
||||||
|
import type { Message } from "./sdk.ts"
|
||||||
|
|
||||||
|
const userMsg = (id: string, overrides: Partial<Message> = {}): Message => ({
|
||||||
|
id,
|
||||||
|
sessionID: "s1",
|
||||||
|
role: "user",
|
||||||
|
time: { created: 1 },
|
||||||
|
...overrides,
|
||||||
|
})
|
||||||
|
|
||||||
|
const assistantMsg = (id: string, overrides: Partial<Message> = {}): Message => ({
|
||||||
|
id,
|
||||||
|
sessionID: "s1",
|
||||||
|
role: "assistant",
|
||||||
|
time: { created: 1 },
|
||||||
|
...overrides,
|
||||||
|
})
|
||||||
|
|
||||||
|
test("no messages -> still busy (nothing to reconcile from)", () => {
|
||||||
|
assert.equal(isSessionActuallyIdle(undefined), false)
|
||||||
|
assert.equal(isSessionActuallyIdle(null), false)
|
||||||
|
assert.equal(isSessionActuallyIdle([]), false)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("last message is a completed assistant reply -> idle (the missed-event case)", () => {
|
||||||
|
const messages = [userMsg("u1"), assistantMsg("a1", { time: { created: 1, completed: 5 } })]
|
||||||
|
assert.equal(isSessionActuallyIdle(messages), true)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("last message is an assistant reply that errored out -> idle", () => {
|
||||||
|
const messages = [userMsg("u1"), assistantMsg("a1", { error: { message: "boom" } })]
|
||||||
|
assert.equal(isSessionActuallyIdle(messages), true)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("last message is a user prompt awaiting a reply -> still busy", () => {
|
||||||
|
const messages = [assistantMsg("a1", { time: { created: 1, completed: 5 } }), userMsg("u2")]
|
||||||
|
assert.equal(isSessionActuallyIdle(messages), false)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("last message is an assistant reply still streaming (no completed, no error) -> still busy", () => {
|
||||||
|
const messages = [userMsg("u1"), assistantMsg("a1")]
|
||||||
|
assert.equal(isSessionActuallyIdle(messages), false)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("a follow-up user prompt after a completed assistant reply -> still busy again", () => {
|
||||||
|
// Server queued a second turn: the previous reply completed, but a new user
|
||||||
|
// message was appended after it, so the run may be in progress again.
|
||||||
|
const messages = [
|
||||||
|
userMsg("u1"),
|
||||||
|
assistantMsg("a1", { time: { created: 1, completed: 5 } }),
|
||||||
|
userMsg("u2"),
|
||||||
|
]
|
||||||
|
assert.equal(isSessionActuallyIdle(messages), false)
|
||||||
|
})
|
||||||
31
src/lib/session-status-reconcile.ts
Normal file
31
src/lib/session-status-reconcile.ts
Normal file
@@ -0,0 +1,31 @@
|
|||||||
|
import type { Message } from "./sdk"
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Decide whether a session the client believes is "busy" has actually
|
||||||
|
* finished, based on the tail of its message history.
|
||||||
|
*
|
||||||
|
* Why this exists (issue #123): `sessionStatus`/`sending` are SSE-driven —
|
||||||
|
* the server's busy -> idle `session.status` event is the only thing that
|
||||||
|
* normally clears them. If the network drops while a session is busy and
|
||||||
|
* that busy -> idle event fires DURING the outage, it is lost: SSE reconnect
|
||||||
|
* resumes the stream from "now", it does not replay missed events. Without a
|
||||||
|
* resync, the UI shows a stuck 'processing' spinner forever even though the
|
||||||
|
* server finished long ago.
|
||||||
|
*
|
||||||
|
* Heuristic: idle iff the most recent message is an assistant message that
|
||||||
|
* has terminated — either it completed normally (`time.completed` set) or it
|
||||||
|
* ended in error (`error` set; the server still finalizes the message on
|
||||||
|
* error, it just never gets a successful `time.completed`). Anything else —
|
||||||
|
* the last message is a user prompt still awaiting a reply, or an assistant
|
||||||
|
* message that hasn't finished streaming — means the run may still be in
|
||||||
|
* progress server-side. Callers MUST treat that as "still busy" and leave the
|
||||||
|
* local state alone: this heuristic only ever clears a stale busy flag, it
|
||||||
|
* never forces a session busy that the server hasn't reported as such, so a
|
||||||
|
* genuinely still-busy session is never clobbered.
|
||||||
|
*/
|
||||||
|
export function isSessionActuallyIdle(messages: Message[] | null | undefined): boolean {
|
||||||
|
if (!messages || messages.length === 0) return false
|
||||||
|
const last = messages[messages.length - 1]
|
||||||
|
if (last.role !== "assistant") return false
|
||||||
|
return Boolean(last.time?.completed) || Boolean(last.error)
|
||||||
|
}
|
||||||
@@ -8,6 +8,7 @@ import { addBreadcrumb } from "../lib/sentry"
|
|||||||
import { AnalyticsEvent, track } from "../lib/analytics"
|
import { AnalyticsEvent, track } from "../lib/analytics"
|
||||||
import { recordSuccessfulSession } from "../lib/store-review"
|
import { recordSuccessfulSession } from "../lib/store-review"
|
||||||
import { isAuthError } from "../lib/api-error"
|
import { isAuthError } from "../lib/api-error"
|
||||||
|
import { isSessionActuallyIdle } from "../lib/session-status-reconcile"
|
||||||
import type { Client, Part, Session, Message } from "../lib/sdk"
|
import type { Client, Part, Session, Message } from "../lib/sdk"
|
||||||
|
|
||||||
// Session status from the server
|
// Session status from the server
|
||||||
@@ -88,6 +89,62 @@ export async function refreshPending(client: Client, sessionID: string) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Re-sync any session currently marked "busy" against the server after an
|
||||||
|
// SSE reconnect. sessionStatus/sending are SSE-driven and there is normally
|
||||||
|
// no other path to idle — if the server's busy -> idle `session.status`
|
||||||
|
// event fired while the network was down, SSE reconnect resumes the stream
|
||||||
|
// from "now" (it does not replay missed events), so without this the busy
|
||||||
|
// flag would never clear and the UI would show a stuck 'processing' spinner
|
||||||
|
// forever (issue #123).
|
||||||
|
//
|
||||||
|
// Only ever CLEARS a busy flag the server confirms is stale via
|
||||||
|
// isSessionActuallyIdle — it never marks a session busy, so it can't
|
||||||
|
// clobber a genuinely still-busy session. Also re-checks sessionStatus right
|
||||||
|
// before writing, so a real session.status event that lands while the fetch
|
||||||
|
// is in flight (e.g. the session went busy again) wins over this resync.
|
||||||
|
async function resyncBusySessions() {
|
||||||
|
const busySessionIDs = Object.entries(useEvents.getState().sessionStatus)
|
||||||
|
.filter(([, status]) => status.type === "busy")
|
||||||
|
.map(([sessionID]) => sessionID)
|
||||||
|
if (busySessionIDs.length === 0) return
|
||||||
|
|
||||||
|
await Promise.all(
|
||||||
|
busySessionIDs.map(async (sessionID) => {
|
||||||
|
try {
|
||||||
|
const sessionsState = useSessions.getState()
|
||||||
|
const session =
|
||||||
|
sessionsState.sessions.find((s) => s.id === sessionID) ??
|
||||||
|
(sessionsState.currentSession?.id === sessionID ? sessionsState.currentSession : undefined)
|
||||||
|
const connState = useConnections.getState()
|
||||||
|
const client = session?.directory
|
||||||
|
? connState.clientForDirectory(session.directory) ?? connState.client
|
||||||
|
: connState.client
|
||||||
|
if (!client) return
|
||||||
|
|
||||||
|
const response = await client.session.messages(sessionID)
|
||||||
|
const messages = (response || []).map((m) => m.info)
|
||||||
|
if (!isSessionActuallyIdle(messages)) return // server says still busy - leave it alone
|
||||||
|
|
||||||
|
// A fresh session.status event may have landed on the SSE stream
|
||||||
|
// while this fetch was in flight — that's authoritative, don't
|
||||||
|
// stomp on it.
|
||||||
|
if (useEvents.getState().sessionStatus[sessionID]?.type !== "busy") return
|
||||||
|
|
||||||
|
useEvents.setState((state) => ({
|
||||||
|
sessionStatus: { ...state.sessionStatus, [sessionID]: { type: "idle" } },
|
||||||
|
statusText: { ...state.statusText, [sessionID]: "" },
|
||||||
|
}))
|
||||||
|
useSessions.setState((state) => ({ sending: { ...state.sending, [sessionID]: false } }))
|
||||||
|
if (useSessions.getState().currentSession?.id === sessionID) {
|
||||||
|
useSessions.getState().refreshMessages()
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
console.warn("[Events] Failed to resync session status for", sessionID, err)
|
||||||
|
}
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
export const useEvents = create<EventsState>((set, get) => ({
|
export const useEvents = create<EventsState>((set, get) => ({
|
||||||
connected: false,
|
connected: false,
|
||||||
authError: false,
|
authError: false,
|
||||||
@@ -118,6 +175,12 @@ export const useEvents = create<EventsState>((set, get) => ({
|
|||||||
// Run in background
|
// Run in background
|
||||||
;(async () => {
|
;(async () => {
|
||||||
let reconnectScheduled = false
|
let reconnectScheduled = false
|
||||||
|
// True if this connect() call is resuming after a prior disconnect —
|
||||||
|
// gates the one-time busy-session resync below so a cold app start
|
||||||
|
// (sessionStatus is always empty then) never triggers it, and a run of
|
||||||
|
// failed retries can't re-arm the check on every attempt.
|
||||||
|
const isReconnect = get().reconnectAttempts > 0
|
||||||
|
let resyncedAfterReconnect = false
|
||||||
const stableTimer = setTimeout(() => {
|
const stableTimer = setTimeout(() => {
|
||||||
if (!currentController.signal.aborted) {
|
if (!currentController.signal.aborted) {
|
||||||
set({ reconnectAttempts: 0, lastDisconnectAt: null })
|
set({ reconnectAttempts: 0, lastDisconnectAt: null })
|
||||||
@@ -163,6 +226,14 @@ export const useEvents = create<EventsState>((set, get) => ({
|
|||||||
for await (const event of client.global.events(currentController.signal)) {
|
for await (const event of client.global.events(currentController.signal)) {
|
||||||
if (currentController.signal.aborted) break
|
if (currentController.signal.aborted) break
|
||||||
|
|
||||||
|
// The stream is genuinely live again (we're actually receiving
|
||||||
|
// data, not just optimistically marked "connected") — resync once
|
||||||
|
// per reconnect, not on every event.
|
||||||
|
if (isReconnect && !resyncedAfterReconnect) {
|
||||||
|
resyncedAfterReconnect = true
|
||||||
|
void resyncBusySessions()
|
||||||
|
}
|
||||||
|
|
||||||
const payload = (event as any).payload || event
|
const payload = (event as any).payload || event
|
||||||
const type = payload.type as string
|
const type = payload.type as string
|
||||||
const props = payload.properties || {}
|
const props = payload.properties || {}
|
||||||
|
|||||||
Reference in New Issue
Block a user