Skip to content
122 changes: 122 additions & 0 deletions packages/app/src/context/global-sdk-event-queue.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
import type { Event, SessionStatus } from "@opencode-ai/sdk/v2/client"
import { describe, expect, test } from "bun:test"
import { coalesceQueuedEvents, type QueuedGlobalEvent } from "./global-sdk-event-queue"

const directory = "/repo"

const delta = (partID: string, value: string, messageID = "msg_1", field = "text"): Event => ({
type: "message.part.delta",
properties: { sessionID: "ses_1", messageID, partID, field, delta: value },
})

const updated = (partID: string, messageID = "msg_1"): Event =>
({
type: "message.part.updated",
properties: {
sessionID: "ses_1",
time: 1,
part: { id: partID, sessionID: "ses_1", messageID, type: "text", text: "full" },
},
}) as Event

const status = (sessionID = "ses_1", type: "busy" | "idle" | "retry" = "busy"): Event => ({
type: "session.status",
properties: {
sessionID,
status: type === "retry" ? ({ type, attempt: 1, message: "retry", next: 1 } satisfies SessionStatus) : { type },
},
})

const queued = (...events: Event[]): QueuedGlobalEvent[] => events.map((payload) => ({ directory, payload }))
const queuedIn = (directory: string, ...events: Event[]): QueuedGlobalEvent[] =>
events.map((payload) => ({ directory, payload }))
const eventTypes = (events: QueuedGlobalEvent[]) => events.map((event) => event.payload.type)
const deltas = (events: QueuedGlobalEvent[]) =>
events
.filter((event) => event.payload.type === "message.part.delta")
.map((event) => (event.payload as Extract<Event, { type: "message.part.delta" }>).properties.delta)

describe("global SDK event queue coalescing", () => {
test("combines only contiguous deltas for the same part", () => {
const events = coalesceQueuedEvents(queued(delta("prt_1", "a"), delta("prt_1", "b")))

expect(events).toHaveLength(1)
expect(deltas(events)).toEqual(["ab"])
})

test("does not merge same-part deltas across another part delta", () => {
const events = coalesceQueuedEvents(queued(delta("prt_1", "a"), delta("prt_2", "x"), delta("prt_1", "b")))

expect(deltas(events)).toEqual(["a", "x", "b"])
})

test("does not merge deltas across non-delta barriers", () => {
const events = coalesceQueuedEvents(queued(delta("prt_1", "a"), status(), delta("prt_1", "b")))

expect(eventTypes(events)).toEqual(["message.part.delta", "session.status", "message.part.delta"])
expect(deltas(events)).toEqual(["a", "b"])
})

test("does not merge same-part deltas for different fields", () => {
const events = coalesceQueuedEvents(
queued(delta("prt_1", "a", "msg_1", "text"), delta("prt_1", "b", "msg_1", "metadata")),
)

expect(deltas(events)).toEqual(["a", "b"])
})

test("drops stale deltas before a full part update but keeps later deltas", () => {
const events = coalesceQueuedEvents(queued(delta("prt_1", "stale"), updated("prt_1"), delta("prt_1", "fresh")))

expect(eventTypes(events)).toEqual(["message.part.updated", "message.part.delta"])
expect(deltas(events)).toEqual(["fresh"])
})

test("keeps only the full update and later delta for delta-update-delta ordering", () => {
const events = coalesceQueuedEvents(queued(delta("prt_1", "before"), updated("prt_1"), delta("prt_1", "after")))

expect(eventTypes(events)).toEqual(["message.part.updated", "message.part.delta"])
expect(deltas(events)).toEqual(["after"])
})

test("handles multiple full updates for the same part without resurrecting stale deltas", () => {
const events = coalesceQueuedEvents(queued(delta("prt_1", "stale"), updated("prt_1"), updated("prt_1")))

expect(eventTypes(events)).toEqual(["message.part.updated", "message.part.updated"])
expect(deltas(events)).toEqual([])
})

test("keeps replaceable event indexes correct after stale delta removal", () => {
const events = coalesceQueuedEvents(
queued(delta("prt_1", "stale"), status("ses_1", "busy"), updated("prt_1"), status("ses_1", "idle")),
)

expect(eventTypes(events)).toEqual(["session.status", "message.part.updated"])
expect(events[0].payload).toEqual(status("ses_1", "idle"))
})

test("keeps delta merging and stale pruning isolated by directory", () => {
const events = coalesceQueuedEvents([
...queuedIn("/repo-a", delta("prt_1", "a")),
...queuedIn("/repo-b", delta("prt_1", "b")),
...queuedIn("/repo-a", delta("prt_1", "c")),
...queuedIn("/repo-a", updated("prt_1")),
...queuedIn("/repo-b", delta("prt_1", "d")),
...queuedIn("/repo-b", updated("prt_1")),
])

expect(eventTypes(events)).toEqual(["message.part.updated", "message.part.updated"])
expect(deltas(events)).toEqual([])
expect(events.map((event) => event.directory)).toEqual(["/repo-a", "/repo-b"])
})

test("does not collide composite keys when directories contain separators", () => {
const events = coalesceQueuedEvents([
...queuedIn("/repo:msg_1", delta("prt_1", "a", "prt_2")),
...queuedIn("/repo", delta("msg_1:prt_1", "b", "prt_2")),
])

expect(deltas(events)).toEqual(["a", "b"])
expect(events.map((event) => event.directory)).toEqual(["/repo:msg_1", "/repo"])
})
})
96 changes: 96 additions & 0 deletions packages/app/src/context/global-sdk-event-queue.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
import type { Event } from "@opencode-ai/sdk/v2/client"

export type QueuedGlobalEvent = { directory: string; payload: Event }

const deltaKey = (event: QueuedGlobalEvent) => {
if (event.payload.type !== "message.part.delta") return
const props = event.payload.properties
return JSON.stringify([event.directory, props.messageID, props.partID, props.field])
}

const partKey = (event: QueuedGlobalEvent) => {
if (event.payload.type === "message.part.delta") {
const props = event.payload.properties
return JSON.stringify([event.directory, props.messageID, props.partID])
}
if (event.payload.type === "message.part.updated") {
const part = event.payload.properties.part
return JSON.stringify([event.directory, part.messageID, part.id])
}
}

const replaceableKey = (event: QueuedGlobalEvent) => {
if (event.payload.type === "session.status")
return JSON.stringify(["session.status", event.directory, event.payload.properties.sessionID])
if (event.payload.type === "lsp.updated") return JSON.stringify(["lsp.updated", event.directory])
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

const appendDelta = (event: QueuedGlobalEvent, delta: string): QueuedGlobalEvent => {
if (event.payload.type !== "message.part.delta") return event
return {
...event,
payload: {
...event.payload,
properties: {
...event.payload.properties,
delta: event.payload.properties.delta + delta,
},
},
}
}

export function coalesceQueuedEvents(events: QueuedGlobalEvent[]): QueuedGlobalEvent[] {
const result: (QueuedGlobalEvent | undefined)[] = []
const replaceable = new Map<string, number>()
const deltaIndexesByPart = new Map<string, Set<number>>()
const pushEvent = (event: QueuedGlobalEvent) => {
const index = result.length
result.push(event)

const key = partKey(event)
if (event.payload.type === "message.part.delta" && key) {
const indexes = deltaIndexesByPart.get(key) ?? new Set<number>()
indexes.add(index)
deltaIndexesByPart.set(key, indexes)
}

return index
}

for (const event of events) {
const replaceKey = replaceableKey(event)
if (replaceKey) {
const index = replaceable.get(replaceKey)
if (index !== undefined) {
result[index] = event
continue
}
replaceable.set(replaceKey, pushEvent(event))
continue
}

if (event.payload.type === "message.part.updated") {
const updatedKey = partKey(event)
if (updatedKey) {
for (const index of deltaIndexesByPart.get(updatedKey) ?? []) {
result[index] = undefined
}
deltaIndexesByPart.delete(updatedKey)
}
pushEvent(event)
continue
Comment thread
Astro-Han marked this conversation as resolved.
}

if (event.payload.type === "message.part.delta") {
const last = result[result.length - 1]
if (last?.payload.type === "message.part.delta" && deltaKey(last) === deltaKey(event)) {
result[result.length - 1] = appendDelta(last, event.payload.properties.delta)
continue
}
}

pushEvent(event)
}

return result.filter((event): event is QueuedGlobalEvent => event !== undefined)
}
41 changes: 4 additions & 37 deletions packages/app/src/context/global-sdk.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { makeEventListener } from "@solid-primitives/event-listener"
import { batch, onCleanup, onMount } from "solid-js"
import z from "zod"
import { createSdkForServer } from "@/utils/server"
import { coalesceQueuedEvents, type QueuedGlobalEvent } from "./global-sdk-event-queue"
import { useLanguage } from "./language"
import { usePlatform } from "./platform"
import { useServer } from "./server"
Expand Down Expand Up @@ -44,50 +45,29 @@ export const { use: useGlobalSDK, provider: GlobalSDKProvider } = createSimpleCo
[key: string]: Event
}>()

type Queued = { directory: string; payload: Event }
const FLUSH_FRAME_MS = 16
const STREAM_YIELD_MS = 8
const RECONNECT_DELAY_MS = 250

let queue: Queued[] = []
let buffer: Queued[] = []
const coalesced = new Map<string, number>()
const staleDeltas = new Set<string>()
let queue: QueuedGlobalEvent[] = []
let buffer: QueuedGlobalEvent[] = []
let timer: ReturnType<typeof setTimeout> | undefined
let last = 0

const deltaKey = (directory: string, messageID: string, partID: string) => `${directory}:${messageID}:${partID}`

const key = (directory: string, payload: Event) => {
if (payload.type === "session.status") return `session.status:${directory}:${payload.properties.sessionID}`
if (payload.type === "lsp.updated") return `lsp.updated:${directory}`
if (payload.type === "message.part.updated") {
const part = payload.properties.part
return `message.part.updated:${directory}:${part.messageID}:${part.id}`
}
}

const flush = () => {
if (timer) clearTimeout(timer)
timer = undefined

if (queue.length === 0) return

const events = queue
const skip = staleDeltas.size > 0 ? new Set(staleDeltas) : undefined
const events = coalesceQueuedEvents(queue)
queue = buffer
buffer = events
queue.length = 0
coalesced.clear()
staleDeltas.clear()

last = Date.now()
batch(() => {
for (const event of events) {
if (skip && event.payload.type === "message.part.delta") {
const props = event.payload.properties
if (skip.has(deltaKey(event.directory, props.messageID, props.partID))) continue
}
emitter.emit(event.directory, event.payload)
}
})
Expand Down Expand Up @@ -160,19 +140,6 @@ export const { use: useGlobalSDK, provider: GlobalSDKProvider } = createSimpleCo
if (payload.type === "sync") {
continue
}
const k = key(directory, payload)
if (k) {
const i = coalesced.get(k)
if (i !== undefined) {
queue[i] = { directory, payload }
if (payload.type === "message.part.updated") {
const part = payload.properties.part
staleDeltas.add(deltaKey(directory, part.messageID, part.id))
}
continue
}
coalesced.set(k, queue.length)
}
queue.push({ directory, payload })
schedule()

Expand Down
16 changes: 16 additions & 0 deletions packages/app/src/pages/session/use-session-hash-scroll.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,3 +14,19 @@ describe("messageIdFromHash", () => {
expect(messageIdFromHash("#review-panel")).toBeUndefined()
})
})

describe("useSessionHashScroll", () => {
test("clearing a message hash notifies the timeline to leave hash history mode", async () => {
const source = await Bun.file(new URL("./use-session-hash-scroll.ts", import.meta.url)).text()

expect(source).toContain("onMessageHashCleared")
expect(source).toContain("input.onMessageHashCleared?.()")
})

test("timeline wires hash clearing to guarded latest-window recovery", async () => {
const source = await Bun.file(new URL("./use-session-timeline-interaction.ts", import.meta.url)).text()

expect(source).toContain("onMessageHashCleared")
expect(source).toContain("historyWindow.clearHashTarget()")
})
})
7 changes: 4 additions & 3 deletions packages/app/src/pages/session/use-session-hash-scroll.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,13 @@ export const useSessionHashScroll = (input: {
pendingMessage: () => string | undefined
setPendingMessage: (value: string | undefined) => void
setActiveMessage: (message: UserMessage | undefined) => void
setTurnStart: (value: number) => void
markHashTarget: (index: number) => void
autoScroll: { pause: () => void; forceScrollToBottom: () => void }
scroller: () => HTMLDivElement | undefined
anchor: (id: string) => string
scheduleScrollState: (el: HTMLDivElement) => void
consumePendingMessage: (key: string) => string | undefined
onMessageHashCleared?: () => void
}) => {
const visibleUserMessages = createMemo(() => input.visibleUserMessages())
const messageById = createMemo(() => new Map(visibleUserMessages().map((m) => [m.id, m])))
Expand Down Expand Up @@ -52,6 +53,7 @@ export const useSessionHashScroll = (input: {
if (!location.hash) return
clearingHash = location.hash
navigate(location.pathname + location.search, { replace: true })
input.onMessageHashCleared?.()
}

const updateHash = (id: string) => {
Expand Down Expand Up @@ -91,9 +93,8 @@ export const useSessionHashScroll = (input: {
if (input.currentMessageId() !== message.id) input.setActiveMessage(message)

const index = messageIndex().get(message.id) ?? -1
if (index !== -1) input.markHashTarget(index)
if (index !== -1 && index < input.turnStart()) {
input.setTurnStart(index)

queue(() => {
seek(message.id, behavior)
})
Expand Down
Loading
Loading