diff --git a/packages/app/e2e/actions.ts b/packages/app/e2e/actions.ts index 90f71f2c6..f3f0019b4 100644 --- a/packages/app/e2e/actions.ts +++ b/packages/app/e2e/actions.ts @@ -790,7 +790,20 @@ export async function seedSessionQuestion( }, }) - if (!result) throw new Error("Timed out seeding question request") + if (!result) { + const [questions, status] = await Promise.all([ + sdk.question.list().then((x) => x.data ?? []).catch((error) => ({ error: String(error) })), + sdk.session.status().then((x) => x.data ?? {}).catch((error) => ({ error: String(error) })), + ]) + throw new Error( + `Timed out seeding question request: ${JSON.stringify({ + sessionID: input.sessionID, + wantedHeader: first.header, + questions, + status, + })}`, + ) + } return { id: result.id } } diff --git a/packages/app/e2e/backend.ts b/packages/app/e2e/backend.ts index a67a16ea9..3575c4c5c 100644 --- a/packages/app/e2e/backend.ts +++ b/packages/app/e2e/backend.ts @@ -85,6 +85,7 @@ export function createBackendEnv(input: { XDG_STATE_HOME: path.join(input.sandbox, "state"), OPENCODE_CLIENT: "app", OPENCODE_STRICT_CONFIG_DEPS: "true", + OPENCODE_E2E_ENABLED: "true", OPENCODE_E2E_LLM_URL: input.llmUrl, } for (const key of Object.keys(env)) { diff --git a/packages/app/e2e/fixtures.ts b/packages/app/e2e/fixtures.ts index 15ceec054..2626221e4 100644 --- a/packages/app/e2e/fixtures.ts +++ b/packages/app/e2e/fixtures.ts @@ -116,6 +116,7 @@ async function promptSend(page: Page) { } type ProjectHandle = { + url: string directory: string slug: string gotoSession: (sessionID?: string) => Promise @@ -517,6 +518,9 @@ function makeProject( gotoSession, trackSession, trackDirectory, + get url() { + return backend.url + }, get directory() { return need().directory }, diff --git a/packages/app/e2e/session/session-composer-dock.spec.ts b/packages/app/e2e/session/session-composer-dock.spec.ts index 64a7735b6..82968e042 100644 --- a/packages/app/e2e/session/session-composer-dock.spec.ts +++ b/packages/app/e2e/session/session-composer-dock.spec.ts @@ -1,4 +1,6 @@ import { mkdir } from "node:fs/promises" +import type { Page } from "@playwright/test" +import type { QuestionRequest } from "@opencode-ai/sdk/v2/client" import { test, expect } from "../fixtures" import { composerEvent, @@ -21,6 +23,11 @@ import { dict as enDict } from "../../src/i18n/en" type Sdk = Parameters[0] type PermissionRule = { permission: string; pattern: string; action: "allow" | "deny" | "ask" } +type ProjectQuestionSeed = { + url: string + directory: string + sdk: Sdk +} async function withDockSession( sdk: Sdk, @@ -80,6 +87,74 @@ async function withDockSeed(sdk: Sdk, sessionID: string, fn: () => Promise } } +function globalEventStream(page: Page) { + return { + cursor: () => + page.evaluate(() => { + const win = window as Window & { + __opencode_e2e?: { globalEventStream?: { cursor: () => string | undefined } } + } + return win.__opencode_e2e?.globalEventStream?.cursor() + }), + stop: () => + page.evaluate(() => { + const win = window as Window & { + __opencode_e2e?: { globalEventStream?: { stop: () => void } } + } + win.__opencode_e2e?.globalEventStream?.stop() + }), + start: () => + page.evaluate(() => { + const win = window as Window & { + __opencode_e2e?: { globalEventStream?: { start: () => void } } + } + win.__opencode_e2e?.globalEventStream?.start() + }), + } +} + +async function e2eAskQuestion( + project: ProjectQuestionSeed, + input: { sessionID: string; questions: typeof defaultQuestions }, +) { + const response = await fetch(`${project.url}/question/__e2e/ask?directory=${encodeURIComponent(project.directory)}`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(input), + }) + expect(response.status).toBe(204) +} + +async function e2ePublishQuestionAsked(project: ProjectQuestionSeed, request: QuestionRequest) { + const response = await fetch( + `${project.url}/question/__e2e/publish-asked?directory=${encodeURIComponent(project.directory)}`, + { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ request }), + }, + ) + expect(response.status).toBe(204) +} + +async function waitForQuestionSeed(project: ProjectQuestionSeed, sessionID: string) { + let current: QuestionRequest | undefined + await expect + .poll( + async () => { + const questions = await project.sdk.question.list().then((response) => response.data ?? []) + current = questions.find( + (question) => question.sessionID === sessionID && question.questions[0]?.header === defaultQuestions[0]?.header, + ) + return !!current + }, + { timeout: 30_000 }, + ) + .toBe(true) + if (!current) throw new Error("Question seed was not visible after polling") + return current +} + async function clearPermissionDock(page: any, label: RegExp) { const dock = page.locator(permissionDockSelector) await expect(dock).toBeVisible() @@ -383,6 +458,57 @@ test("blocked question flow unblocks after submit", async ({ page, llm, project ) }) +test("question dock recovers after missed question.asked via SSE replay", async ({ page, project }) => { + await project.open() + await withDockSession( + project.sdk, + "e2e composer dock question replay", + async (session) => { + await withDockSeed(project.sdk, session.id, async () => { + await project.gotoSession(session.id) + + const stream = globalEventStream(page) + + await expect.poll(stream.cursor, { timeout: 10_000 }).toMatch(/:/) + await stream.stop() + await e2eAskQuestion(project, { sessionID: session.id, questions: defaultQuestions }) + await waitForQuestionSeed(project, session.id) + + await expect(page.locator(questionDockSelector)).toHaveCount(0, { timeout: 750 }) + await stream.start() + + await expectQuestionBlocked(page) + await expect(page.locator(questionDockSelector)).toHaveCount(1) + }) + }, + { trackSession: project.trackSession }, + ) +}) + +test("stale question.asked does not reopen after question reply", async ({ page, project }) => { + await project.open() + await withDockSession( + project.sdk, + "e2e composer dock stale question", + async (session) => { + await withDockSeed(project.sdk, session.id, async () => { + await project.gotoSession(session.id) + + await e2eAskQuestion(project, { sessionID: session.id, questions: defaultQuestions }) + const request = await waitForQuestionSeed(project, session.id) + + await expectQuestionBlocked(page) + await project.sdk.question.reply({ requestID: request.id, questionReply: { answers: [["Continue"]] } }) + await expectQuestionOpen(page) + + await e2ePublishQuestionAsked(project, request) + await expect(page.locator(questionDockSelector)).toHaveCount(0, { timeout: 1_000 }) + }) + }, + { trackSession: project.trackSession }, + ) +}) + test("blocked question flow supports skipping one question before submit", async ({ page, llm, project }) => { await project.open() await withDockSession( diff --git a/packages/app/src/context/global-sdk.tsx b/packages/app/src/context/global-sdk.tsx index 256f03bdb..7dd638e36 100644 --- a/packages/app/src/context/global-sdk.tsx +++ b/packages/app/src/context/global-sdk.tsx @@ -4,11 +4,13 @@ import { createGlobalEmitter } from "@solid-primitives/event-bus" import { makeEventListener } from "@solid-primitives/event-listener" import { batch, onCleanup, onMount } from "solid-js" import z from "zod" +import type { E2EWindow } from "@/testing/terminal" 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" +import { createSseCursor } from "./global-sdk/sse-cursor" const abortError = z.object({ name: z.literal("AbortError"), @@ -91,6 +93,7 @@ export const { use: useGlobalSDK, provider: GlobalSDKProvider } = createSimpleCo const HEARTBEAT_TIMEOUT_MS = 15_000 let lastEventAt = Date.now() let heartbeat: ReturnType | undefined + const replayCursor = createSseCursor() const resetHeartbeat = () => { lastEventAt = Date.now() if (heartbeat) clearTimeout(heartbeat) @@ -119,6 +122,10 @@ export const { use: useGlobalSDK, provider: GlobalSDKProvider } = createSimpleCo try { const events = await eventSdk.global.event({ signal: attempt.signal, + headers: replayCursor.headers(), + onSseEvent: (event) => { + replayCursor.update(event.id) + }, onSseError: (error) => { if (aborted(error)) return if (streamErrorLogged) return @@ -178,7 +185,21 @@ export const { use: useGlobalSDK, provider: GlobalSDKProvider } = createSimpleCo clearHeartbeat() } + const e2e = () => { + if (typeof window === "undefined") return + const state = (window as E2EWindow).__opencode_e2e + if (!state) return + state.globalEventStream = { + stop, + start: () => { + void start() + }, + cursor: replayCursor.current, + } + } + onMount(() => { + e2e() makeEventListener(document, "visibilitychange", () => { if (document.visibilityState !== "visible") return if (!started) return diff --git a/packages/app/src/context/global-sdk/sse-cursor.test.ts b/packages/app/src/context/global-sdk/sse-cursor.test.ts new file mode 100644 index 000000000..8a9a5de10 --- /dev/null +++ b/packages/app/src/context/global-sdk/sse-cursor.test.ts @@ -0,0 +1,30 @@ +import { describe, expect, test } from "bun:test" +import { createSseCursor } from "./sse-cursor" + +describe("createSseCursor", () => { + test("starts without a cursor", () => { + const cursor = createSseCursor() + expect(cursor.current()).toBeUndefined() + expect(cursor.headers()).toBeUndefined() + }) + + test("stores the latest non-empty event id", () => { + const cursor = createSseCursor() + cursor.update(undefined) + cursor.update("") + cursor.update("boot:1") + cursor.update("boot:2") + + expect(cursor.current()).toBe("boot:2") + }) + + test("builds Last-Event-ID headers when a cursor exists", () => { + const cursor = createSseCursor() + cursor.update("boot:7") + + const headers = cursor.headers() + + expect(headers).toBeInstanceOf(Headers) + expect(headers?.get("Last-Event-ID")).toBe("boot:7") + }) +}) diff --git a/packages/app/src/context/global-sdk/sse-cursor.ts b/packages/app/src/context/global-sdk/sse-cursor.ts new file mode 100644 index 000000000..c33668c51 --- /dev/null +++ b/packages/app/src/context/global-sdk/sse-cursor.ts @@ -0,0 +1,20 @@ +export function createSseCursor() { + let value: string | undefined + + return { + current() { + return value + }, + update(id: string | undefined) { + if (!id) return + // SSE ids come from the server replay layer; keep this helper transport-only. + value = id + }, + headers() { + if (!value) return undefined + const headers = new Headers() + headers.set("Last-Event-ID", value) + return headers + }, + } +} diff --git a/packages/app/src/context/global-sync.tsx b/packages/app/src/context/global-sync.tsx index 2528ffcdf..9396b4f8a 100644 --- a/packages/app/src/context/global-sync.tsx +++ b/packages/app/src/context/global-sync.tsx @@ -16,6 +16,7 @@ import { Persist, persisted } from "@/utils/persist" import type { InitError } from "../pages/error" import { useGlobalSDK } from "./global-sdk" import { bootstrapDirectory, bootstrapGlobal, clearProviderRev } from "./global-sync/bootstrap" +import { createBlockerTerminalCache } from "./global-sync/blocker-terminal-cache" import { createChildStoreManager } from "./global-sync/child-store" import { applyDirectoryEvent, applyGlobalEvent, cleanupDroppedSessionCaches } from "./global-sync/event-reducer" import { createRefreshQueue } from "./global-sync/queue" @@ -57,6 +58,7 @@ function createGlobalSync() { const booting = new Map>() const sessionLoads = new Map>() const sessionMeta = new Map() + const blockerTerminals = createBlockerTerminalCache() const [projectCache, setProjectCache, projectInit] = persisted( Persist.global("globalSync.project", ["globalSync.project.v1"]), @@ -166,6 +168,7 @@ function createGlobalSync() { onDispose: (directory) => { queue.clear(directory) sessionMeta.delete(directory) + blockerTerminals.clearDirectory(directory) sdkCache.delete(directory) clearProviderRev(directory) clearSessionPrefetchDirectory(directory) @@ -330,6 +333,7 @@ function createGlobalSync() { setStore, push: queue.push, setSessionTodo, + blockerTerminals, vcsCache: children.vcsCache.get(directory), loadLsp: () => { void sdkFor(directory) diff --git a/packages/app/src/context/global-sync/blocker-terminal-cache.test.ts b/packages/app/src/context/global-sync/blocker-terminal-cache.test.ts new file mode 100644 index 000000000..91b1fb057 --- /dev/null +++ b/packages/app/src/context/global-sync/blocker-terminal-cache.test.ts @@ -0,0 +1,51 @@ +import { describe, expect, test } from "bun:test" +import { createBlockerTerminalCache } from "./blocker-terminal-cache" + +describe("createBlockerTerminalCache", () => { + test("marks and finds terminal blocker ids by kind, directory, session, and request", () => { + const cache = createBlockerTerminalCache({ now: () => 1000 }) + + cache.mark("question", "/repo", "ses_1", "q1") + + expect(cache.has("question", "/repo", "ses_1", "q1")).toBe(true) + expect(cache.has("permission", "/repo", "ses_1", "q1")).toBe(false) + expect(cache.has("question", "/other", "ses_1", "q1")).toBe(false) + expect(cache.has("question", "/repo", "ses_2", "q1")).toBe(false) + }) + + test("expires old entries by ttl", () => { + let now = 1000 + const cache = createBlockerTerminalCache({ ttlMs: 100, now: () => now }) + + cache.mark("question", "/repo", "ses_1", "q1") + now = 1200 + + expect(cache.has("question", "/repo", "ses_1", "q1")).toBe(false) + }) + + test("prunes oldest entries by max size", () => { + let now = 1000 + const cache = createBlockerTerminalCache({ max: 2, now: () => now }) + + cache.mark("question", "/repo", "ses_1", "q1") + now += 1 + cache.mark("question", "/repo", "ses_1", "q2") + now += 1 + cache.mark("question", "/repo", "ses_1", "q3") + + expect(cache.has("question", "/repo", "ses_1", "q1")).toBe(false) + expect(cache.has("question", "/repo", "ses_1", "q2")).toBe(true) + expect(cache.has("question", "/repo", "ses_1", "q3")).toBe(true) + }) + + test("clears all entries for a directory", () => { + const cache = createBlockerTerminalCache({ now: () => 1000 }) + + cache.mark("question", "/repo", "ses_1", "q1") + cache.mark("question", "/other", "ses_1", "q1") + cache.clearDirectory("/repo") + + expect(cache.has("question", "/repo", "ses_1", "q1")).toBe(false) + expect(cache.has("question", "/other", "ses_1", "q1")).toBe(true) + }) +}) diff --git a/packages/app/src/context/global-sync/blocker-terminal-cache.ts b/packages/app/src/context/global-sync/blocker-terminal-cache.ts new file mode 100644 index 000000000..994185980 --- /dev/null +++ b/packages/app/src/context/global-sync/blocker-terminal-cache.ts @@ -0,0 +1,50 @@ +export type BlockerKind = "question" | "permission" + +type Entry = { + directory: string + createdAt: number +} + +const DEFAULT_MAX = 2048 +const DEFAULT_TTL_MS = 5 * 60 * 1000 + +export function createBlockerTerminalCache(input?: { max?: number; ttlMs?: number; now?: () => number }) { + const max = input?.max ?? DEFAULT_MAX + const ttlMs = input?.ttlMs ?? DEFAULT_TTL_MS + const now = input?.now ?? Date.now + const entries = new Map() + + const keyFor = (kind: BlockerKind, directory: string, sessionID: string, requestID: string) => + JSON.stringify([kind, directory, sessionID, requestID]) + + const prune = () => { + const cutoff = now() - ttlMs + for (const [key, entry] of entries) { + if (entry.createdAt >= cutoff) continue + entries.delete(key) + } + while (entries.size > max) { + const oldest = entries.keys().next().value + if (!oldest) break + entries.delete(oldest) + } + } + + return { + mark(kind: BlockerKind, directory: string, sessionID: string, requestID: string) { + const key = keyFor(kind, directory, sessionID, requestID) + entries.delete(key) + entries.set(key, { directory, createdAt: now() }) + prune() + }, + has(kind: BlockerKind, directory: string, sessionID: string, requestID: string) { + prune() + return entries.has(keyFor(kind, directory, sessionID, requestID)) + }, + clearDirectory(directory: string) { + for (const [key, entry] of entries) { + if (entry.directory === directory) entries.delete(key) + } + }, + } +} diff --git a/packages/app/src/context/global-sync/event-reducer.test.ts b/packages/app/src/context/global-sync/event-reducer.test.ts index cfff37d2b..bacc0146c 100644 --- a/packages/app/src/context/global-sync/event-reducer.test.ts +++ b/packages/app/src/context/global-sync/event-reducer.test.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test" import type { Message, Part, PermissionRequest, Project, QuestionRequest, Session } from "@opencode-ai/sdk/v2/client" import { createStore } from "solid-js/store" import type { State } from "./types" +import { createBlockerTerminalCache } from "./blocker-terminal-cache" import { applyDirectoryEvent, applyGlobalEvent, cleanupDroppedSessionCaches } from "./event-reducer" const rootSession = (input: { id: string; parentID?: string; archived?: number; created?: number; updated?: number }) => @@ -493,6 +494,66 @@ describe("applyDirectoryEvent", () => { expect(store.question[sessionID]?.map((x) => x.id)).toEqual(["q_1", "q_3"]) }) + test("question.replied before question.asked prevents stale ask from reopening", () => { + const blockerTerminals = createBlockerTerminalCache({ now: () => 1000 }) + const [store, setStore] = createStore(baseState()) + + applyDirectoryEvent({ + event: { type: "question.replied", properties: { sessionID: "ses_1", requestID: "q1" } }, + directory: "/repo", + store, + setStore, + push() {}, + loadLsp() {}, + blockerTerminals, + }) + + applyDirectoryEvent({ + event: { + type: "question.asked", + properties: questionRequest("q1", "ses_1"), + }, + directory: "/repo", + store, + setStore, + push() {}, + loadLsp() {}, + blockerTerminals, + }) + + expect(store.question.ses_1).toBeUndefined() + }) + + test("permission.replied before permission.asked prevents stale ask from reopening", () => { + const blockerTerminals = createBlockerTerminalCache({ now: () => 1000 }) + const [store, setStore] = createStore(baseState()) + + applyDirectoryEvent({ + event: { type: "permission.replied", properties: { sessionID: "ses_1", requestID: "perm_1" } }, + directory: "/repo", + store, + setStore, + push() {}, + loadLsp() {}, + blockerTerminals, + }) + + applyDirectoryEvent({ + event: { + type: "permission.asked", + properties: permissionRequest("perm_1", "ses_1"), + }, + directory: "/repo", + store, + setStore, + push() {}, + loadLsp() {}, + blockerTerminals, + }) + + expect(store.permission.ses_1).toBeUndefined() + }) + test("updates vcs branch in store and cache", () => { const [store, setStore] = createStore(baseState({ vcs: { branch: "main", default_branch: "main" } })) const [cacheStore, setCacheStore] = createStore({ diff --git a/packages/app/src/context/global-sync/event-reducer.ts b/packages/app/src/context/global-sync/event-reducer.ts index 351d3cdf0..09dfb9939 100644 --- a/packages/app/src/context/global-sync/event-reducer.ts +++ b/packages/app/src/context/global-sync/event-reducer.ts @@ -15,6 +15,7 @@ import type { State, VcsCache } from "./types" import { trimSessions } from "./session-trim" import { dropSessionCaches } from "./session-cache" import { diffs as list, message as clean } from "@/utils/diffs" +import type { createBlockerTerminalCache } from "./blocker-terminal-cache" const SKIP_PARTS = new Set(["patch", "step-start", "step-finish"]) @@ -95,6 +96,7 @@ export function applyDirectoryEvent(input: { loadLsp: () => void vcsCache?: VcsCache setSessionTodo?: (sessionID: string, todos: Todo[] | undefined) => void + blockerTerminals?: ReturnType }) { const event = input.event switch (event.type) { @@ -277,6 +279,7 @@ export function applyDirectoryEvent(input: { } case "permission.asked": { const permission = event.properties as PermissionRequest + if (input.blockerTerminals?.has("permission", input.directory, permission.sessionID, permission.id)) break const permissions = input.store.permission[permission.sessionID] if (!permissions) { input.setStore("permission", permission.sessionID, [permission]) @@ -298,6 +301,7 @@ export function applyDirectoryEvent(input: { } case "permission.replied": { const props = event.properties as { sessionID: string; requestID: string } + input.blockerTerminals?.mark("permission", input.directory, props.sessionID, props.requestID) const permissions = input.store.permission[props.sessionID] if (!permissions) break const result = Binary.search(permissions, props.requestID, (p) => p.id) @@ -313,6 +317,7 @@ export function applyDirectoryEvent(input: { } case "question.asked": { const question = event.properties as QuestionRequest + if (input.blockerTerminals?.has("question", input.directory, question.sessionID, question.id)) break const questions = input.store.question[question.sessionID] if (!questions) { input.setStore("question", question.sessionID, [question]) @@ -335,6 +340,7 @@ export function applyDirectoryEvent(input: { case "question.replied": case "question.rejected": { const props = event.properties as { sessionID: string; requestID: string } + input.blockerTerminals?.mark("question", input.directory, props.sessionID, props.requestID) const questions = input.store.question[props.sessionID] if (!questions) break const result = Binary.search(questions, props.requestID, (q) => q.id) diff --git a/packages/app/src/testing/terminal.ts b/packages/app/src/testing/terminal.ts index db8001ddf..86855a94a 100644 --- a/packages/app/src/testing/terminal.ts +++ b/packages/app/src/testing/terminal.ts @@ -30,6 +30,11 @@ export type E2EWindow = Window & { terminals?: Record controls?: Record } + globalEventStream?: { + stop: () => void + start: () => void + cursor: () => string | undefined + } } } diff --git a/packages/opencode/src/server/event-replay.ts b/packages/opencode/src/server/event-replay.ts new file mode 100644 index 000000000..58b051fc3 --- /dev/null +++ b/packages/opencode/src/server/event-replay.ts @@ -0,0 +1,169 @@ +import { randomUUID } from "node:crypto" + +export type GlobalEventEnvelope = { + directory?: string + project?: string + workspace?: string + payload: { + type: string + properties: unknown + } +} + +export type ReplayCursor = { + bootID: string + seq: number +} + +export type ReplayRecord = { + id: string + seq: number + createdAt: number + envelope: GlobalEventEnvelope +} + +export type ReplaySnapshot = { + bootID: string + fenceSeq: number + fenceID: string + replay: ReplayRecord[] + gap: boolean + invalidCursor: boolean +} + +const DEFAULT_MAX_RECORDS = 2048 +const DEFAULT_MAX_AGE_MS = 5 * 60 * 1000 + +// Keep this list intentionally small. These events are low-volume and session +// or blocker critical; high-volume streaming events recover through bootstrap. +const REPLAYABLE_EVENT_TYPES = new Set([ + "question.asked", + "question.replied", + "question.rejected", + "permission.asked", + "permission.replied", + "session.created", + "session.updated", + "session.deleted", + "session.status", + "server.instance.disposed", +]) + +export function parseReplayCursor(input: string | undefined): ReplayCursor | undefined { + if (!input) return undefined + const index = input.lastIndexOf(":") + if (index <= 0 || index === input.length - 1) return undefined + + const bootID = input.slice(0, index) + const seq = Number(input.slice(index + 1)) + if (!Number.isSafeInteger(seq) || seq < 0) return undefined + + return { bootID, seq } +} + +export function isReplayableGlobalEvent(envelope: GlobalEventEnvelope): boolean { + return REPLAYABLE_EVENT_TYPES.has(envelope.payload?.type) +} + +export class EventReplayStore { + private bootID: string + private readonly maxRecords: number + private readonly maxAgeMs: number + private readonly now: () => number + private seq = 0 + private records: ReplayRecord[] = [] + + constructor(input?: { bootID?: string; maxRecords?: number; maxAgeMs?: number; now?: () => number }) { + this.bootID = input?.bootID ?? `${Date.now().toString(36)}-${randomUUID().slice(0, 8)}` + const maxRecords = input?.maxRecords ?? DEFAULT_MAX_RECORDS + const maxAgeMs = input?.maxAgeMs ?? DEFAULT_MAX_AGE_MS + if (!Number.isInteger(maxRecords) || maxRecords < 0) { + throw new RangeError("maxRecords must be a non-negative integer") + } + if (!Number.isFinite(maxAgeMs) || maxAgeMs < 0) { + throw new RangeError("maxAgeMs must be a non-negative number") + } + this.maxRecords = maxRecords + this.maxAgeMs = maxAgeMs + this.now = input?.now ?? Date.now + } + + latestID(): string { + return this.formatID(this.seq) + } + + append(envelope: GlobalEventEnvelope): ReplayRecord | undefined { + if (!isReplayableGlobalEvent(envelope)) return undefined + + const seq = ++this.seq + const record: ReplayRecord = { + id: this.formatID(seq), + seq, + createdAt: this.now(), + envelope: structuredClone(envelope), + } + + this.records.push(record) + this.prune() + + return record + } + + snapshot(lastEventID: string | undefined): ReplaySnapshot { + this.prune() + const fenceSeq = this.seq + const fenceID = this.formatID(fenceSeq) + const parsed = parseReplayCursor(lastEventID) + const invalidCursor = !!lastEventID && (!parsed || parsed.bootID !== this.bootID || parsed.seq > fenceSeq) + const cursorSeq = invalidCursor || !parsed ? fenceSeq : parsed.seq + const earliestSeq = this.records[0]?.seq + const gap = + !invalidCursor && + !!parsed && + parsed.seq < fenceSeq && + (earliestSeq === undefined || parsed.seq < earliestSeq - 1) + const replay = + invalidCursor || !parsed + ? [] + : this.records.filter((record) => record.seq > cursorSeq && record.seq <= fenceSeq) + + return { + bootID: this.bootID, + fenceSeq, + fenceID, + replay, + gap, + invalidCursor, + } + } + + reset() { + this.bootID = `${this.now().toString(36)}-${randomUUID().slice(0, 8)}` + this.seq = 0 + this.clear() + } + + clear() { + this.records = [] + } + + clearDirectory(directory: string) { + this.records = this.records.filter((record) => record.envelope.directory !== directory) + } + + recordsForTest(): readonly ReplayRecord[] { + return this.records.slice() + } + + private formatID(seq: number): string { + return `${this.bootID}:${seq}` + } + + private prune() { + const cutoff = this.now() - this.maxAgeMs + while (this.records.length > this.maxRecords) this.records.shift() + while (this.records[0] && this.records[0].createdAt < cutoff) this.records.shift() + } +} + +export * as EventReplay from "./event-replay" diff --git a/packages/opencode/src/server/instance/global.ts b/packages/opencode/src/server/instance/global.ts index 7da22212c..12b3b89d4 100644 --- a/packages/opencode/src/server/instance/global.ts +++ b/packages/opencode/src/server/instance/global.ts @@ -14,11 +14,125 @@ import { Log } from "@opencode-ai/core/util/log" import { lazy } from "../../util/lazy" import { Config } from "../../config/config" import { errors } from "../error" +import { EventReplayStore, type GlobalEventEnvelope, type ReplayRecord } from "../event-replay" const log = Log.create({ service: "server" }) export const GlobalDisposedEvent = BusEvent.define("global.disposed", z.object({})) +type SsePacket = { + id?: string + data: string + replaySeq?: number +} + +export type GlobalEventReplayPacket = SsePacket + +function packetForEnvelope(envelope: GlobalEventEnvelope): SsePacket { + return { + data: JSON.stringify(envelope), + } +} + +function packetForRecord(record: ReplayRecord): SsePacket { + return { + id: record.id, + replaySeq: record.seq, + data: JSON.stringify(record.envelope), + } +} + +export function createGlobalEventReplayBridge(input?: { replayStore?: EventReplayStore }) { + const replayStore = input?.replayStore ?? new EventReplayStore() + const listeners = new Set<(packet: GlobalEventReplayPacket) => void>() + + return { + replayStore, + append(event: GlobalEventEnvelope) { + if (event.payload.type === GlobalDisposedEvent.type) { + // Global dispose starts a new replay generation. Online clients still + // receive the live packet; offline clients refresh on bootID mismatch. + replayStore.reset() + } + if (event.payload.type === "server.instance.disposed" && event.directory) { + replayStore.clearDirectory(event.directory) + } + const record = replayStore.append(event) + const packet = record ? packetForRecord(record) : packetForEnvelope(event) + for (const listener of listeners) listener(packet) + return packet + }, + subscribe(listener: (packet: GlobalEventReplayPacket) => void) { + listeners.add(listener) + return () => listeners.delete(listener) + }, + } +} + +export function openGlobalEventReplayConnection(input: { + bridge: ReturnType + lastEventID?: string + push: (packet: GlobalEventReplayPacket) => void +}) { + const liveBuffer: GlobalEventReplayPacket[] = [] + let replaying = true + const pushLive = (packet: GlobalEventReplayPacket) => { + if (replaying) { + liveBuffer.push(packet) + return + } + input.push(packet) + } + + const unsubscribeLive = input.bridge.subscribe(pushLive) + const opened = input.bridge.replayStore.snapshot(input.lastEventID) + + const pushConnected = (id?: string) => { + input.push({ + id, + data: JSON.stringify({ + payload: { + type: "server.connected", + properties: {}, + }, + }), + }) + } + + // Fresh connect seeds a cursor. Valid reconnect replays only missed records. + // Invalid or gapped reconnect sends one refresh signal and advances the cursor. + if (!input.lastEventID) { + pushConnected(opened.fenceID) + } + + // Do not send partial replay for invalid/gapped cursors. Missing earlier + // events can make retained blocker records stale, so bootstrap owns recovery. + if (opened.invalidCursor || opened.gap) { + if (input.lastEventID) pushConnected(opened.fenceID) + } else { + for (const record of opened.replay) { + input.push(packetForRecord(record)) + } + } + + replaying = false + for (const packet of liveBuffer) { + if (packet.replaySeq !== undefined && packet.replaySeq <= opened.fenceSeq) continue + input.push(packet) + } + liveBuffer.length = 0 + + return () => { + unsubscribeLive() + } +} + +const globalEventReplay = createGlobalEventReplayBridge() + +GlobalBus.on("event", (event) => { + globalEventReplay.append(event as GlobalEventEnvelope) +}) + async function streamEvents(c: Context, subscribe: (q: AsyncQueue) => () => void) { return streamSSE(c, async (stream) => { const q = new AsyncQueue() @@ -69,6 +183,55 @@ async function streamEvents(c: Context, subscribe: (q: AsyncQueue }) } +export async function streamGlobalEvents( + c: Context, + bridge: ReturnType = globalEventReplay, +) { + const lastEventID = c.req.header("Last-Event-ID") ?? c.req.header("last-event-id") ?? undefined + + return streamSSE(c, async (stream) => { + const q = new AsyncQueue() + let done = false + const unsubscribe = openGlobalEventReplayConnection({ + bridge, + lastEventID, + push: (packet) => q.push(packet), + }) + + const heartbeat = setInterval(() => { + q.push({ + data: JSON.stringify({ + payload: { + type: "server.heartbeat", + properties: {}, + }, + }), + }) + }, 10_000) + + const stop = () => { + if (done) return + done = true + clearInterval(heartbeat) + unsubscribe() + q.push(null) + log.info("global event disconnected") + } + + stream.onAbort(stop) + + try { + for await (const packet of q) { + if (packet === null) return + const { replaySeq: _replaySeq, ...sse } = packet + await stream.writeSSE(sse) + } + } finally { + stop() + } + }) +} + export const GlobalRoutes = lazy(() => new Hono() .get( @@ -126,13 +289,7 @@ export const GlobalRoutes = lazy(() => c.header("X-Accel-Buffering", "no") c.header("X-Content-Type-Options", "nosniff") - return streamEvents(c, (q) => { - async function handler(event: any) { - q.push(JSON.stringify(event)) - } - GlobalBus.on("event", handler) - return () => GlobalBus.off("event", handler) - }) + return streamGlobalEvents(c) }, ) .get( diff --git a/packages/opencode/src/server/instance/question.ts b/packages/opencode/src/server/instance/question.ts index bf1dfe284..833ae2c35 100644 --- a/packages/opencode/src/server/instance/question.ts +++ b/packages/opencode/src/server/instance/question.ts @@ -4,12 +4,65 @@ import { resolver } from "hono-openapi" import { QuestionID } from "@/question/schema" import { Question } from "../../question" import { AppRuntime } from "@/effect/app-runtime" +import { Bus } from "@/bus" +import { Env } from "@/env" +import { SessionID } from "@/session/schema" +import { Log } from "@opencode-ai/core/util/log" import z from "zod" import { errors } from "../error" import { lazy } from "../../util/lazy" +const log = Log.create({ service: "server" }) +const e2eQuestionRoutesEnabled = () => + Env.get("OPENCODE_E2E_ENABLED") === "true" && !!Env.get("OPENCODE_E2E_LLM_URL") + export const QuestionRoutes = lazy(() => new Hono() + // E2E-only hooks exercise question transport/UI flows without relying on flaky LLM seeding. + .post( + "/__e2e/ask", + validator( + "json", + z.object({ + sessionID: SessionID.zod, + questions: z.array(Question.Info.zod).min(1).max(4), + }), + ), + async (c) => { + if (!e2eQuestionRoutesEnabled()) return c.notFound() + + const json = c.req.valid("json") + void AppRuntime.runPromise( + Question.Service.use((svc) => + svc.ask({ + sessionID: json.sessionID, + questions: json.questions, + }), + ), + ).catch((error) => { + log.error("e2e question seed failed", { sessionID: json.sessionID, error }) + }) + + return c.body(null, 204) + }, + ) + .post( + "/__e2e/publish-asked", + validator( + "json", + z.object({ + request: Question.Request.zod, + }), + ), + async (c) => { + if (!e2eQuestionRoutesEnabled()) return c.notFound() + + const json = c.req.valid("json") + await AppRuntime.runPromise(Bus.Service.use((bus) => bus.publish(Question.Event.Asked, json.request))) + + return c.body(null, 204) + }, + ) .get( "/", describeRoute({ diff --git a/packages/opencode/test/server/event-replay.test.ts b/packages/opencode/test/server/event-replay.test.ts new file mode 100644 index 000000000..0163b3b9b --- /dev/null +++ b/packages/opencode/test/server/event-replay.test.ts @@ -0,0 +1,211 @@ +import { describe, expect, test } from "bun:test" +import { EventReplayStore, isReplayableGlobalEvent, parseReplayCursor } from "../../src/server/event-replay" + +const question = (id: string, sessionID = "ses_1") => ({ + directory: "/repo", + project: "proj", + workspace: "work", + payload: { + type: "question.asked", + properties: { id, sessionID, questions: [{ header: "H", question: "Q", options: [] }] }, + }, +}) + +const event = (type: string) => ({ + directory: "/repo", + payload: { + type, + properties: {}, + }, +}) + +describe("parseReplayCursor", () => { + test("parses boot id and numeric sequence", () => { + expect(parseReplayCursor("boot-abc:42")).toEqual({ bootID: "boot-abc", seq: 42 }) + }) + + test("rejects malformed cursors", () => { + expect(parseReplayCursor(undefined)).toBeUndefined() + expect(parseReplayCursor("")).toBeUndefined() + expect(parseReplayCursor("boot")).toBeUndefined() + expect(parseReplayCursor("boot:abc")).toBeUndefined() + expect(parseReplayCursor("boot:-1")).toBeUndefined() + }) +}) + +describe("isReplayableGlobalEvent", () => { + test("allows blocker and session state events", () => { + expect(isReplayableGlobalEvent(event("question.asked"))).toBe(true) + expect(isReplayableGlobalEvent(event("question.replied"))).toBe(true) + expect(isReplayableGlobalEvent(event("question.rejected"))).toBe(true) + expect(isReplayableGlobalEvent(event("permission.asked"))).toBe(true) + expect(isReplayableGlobalEvent(event("permission.replied"))).toBe(true) + expect(isReplayableGlobalEvent(event("session.created"))).toBe(true) + expect(isReplayableGlobalEvent(event("session.updated"))).toBe(true) + expect(isReplayableGlobalEvent(event("session.deleted"))).toBe(true) + expect(isReplayableGlobalEvent(event("session.status"))).toBe(true) + expect(isReplayableGlobalEvent(event("server.instance.disposed"))).toBe(true) + }) + + test("rejects high-volume and unrelated events", () => { + expect(isReplayableGlobalEvent(event("message.part.delta"))).toBe(false) + expect(isReplayableGlobalEvent(event("message.part.updated"))).toBe(false) + expect(isReplayableGlobalEvent(event("message.updated"))).toBe(false) + expect(isReplayableGlobalEvent(event("todo.updated"))).toBe(false) + expect(isReplayableGlobalEvent(event("lsp.updated"))).toBe(false) + expect(isReplayableGlobalEvent(event("vcs.branch.updated"))).toBe(false) + }) +}) + +describe("EventReplayStore", () => { + test("rejects invalid retention configuration", () => { + expect(() => new EventReplayStore({ maxRecords: -1 })).toThrow(RangeError) + expect(() => new EventReplayStore({ maxRecords: 1.5 })).toThrow(RangeError) + expect(() => new EventReplayStore({ maxAgeMs: -1 })).toThrow(RangeError) + expect(() => new EventReplayStore({ maxAgeMs: Number.NaN })).toThrow(RangeError) + }) + + test("assigns monotonic ids under one boot id", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + const first = store.append(question("q1")) + const second = store.append(question("q2")) + + expect(first?.id).toBe("boot:1") + expect(second?.id).toBe("boot:2") + expect(store.latestID()).toBe("boot:2") + }) + + test("does not store non-replayable events", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + expect(store.append(event("message.part.delta"))).toBeUndefined() + expect(store.recordsForTest()).toEqual([]) + expect(store.latestID()).toBe("boot:0") + }) + + test("stores cloned envelopes", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + const input = question("q1") + const record = store.append(input) + ;(input.payload.properties as { id: string }).id = "mutated" + + expect((record?.envelope.payload.properties as { id: string }).id).toBe("q1") + expect((store.recordsForTest()[0].envelope.payload.properties as { id: string }).id).toBe("q1") + }) + + test("prunes by max record count", () => { + const store = new EventReplayStore({ bootID: "boot", maxRecords: 2, now: () => 1000 }) + store.append(question("q1")) + store.append(question("q2")) + store.append(question("q3")) + + expect(store.recordsForTest().map((record) => record.seq)).toEqual([2, 3]) + }) + + test("prunes by max age", () => { + let now = 1000 + const store = new EventReplayStore({ bootID: "boot", maxAgeMs: 100, now: () => now }) + store.append(question("q1")) + now = 1200 + store.append(question("q2")) + + expect(store.recordsForTest().map((record) => record.seq)).toEqual([2]) + }) + + test("replays records after cursor in ascending order", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + store.append(question("q1")) + store.append(question("q2")) + store.append(question("q3")) + + const opened = store.snapshot("boot:1") + + expect(opened.invalidCursor).toBe(false) + expect(opened.gap).toBe(false) + expect(opened.replay.map((record) => record.id)).toEqual(["boot:2", "boot:3"]) + }) + + test("returns no replay for invalid boot id", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + store.append(question("q1")) + + const opened = store.snapshot("old:1") + + expect(opened.invalidCursor).toBe(true) + expect(opened.gap).toBe(false) + expect(opened.replay).toEqual([]) + expect(opened.fenceID).toBe("boot:1") + }) + + test("treats a same-boot future cursor as invalid", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + store.append(question("q1")) + + const opened = store.snapshot("boot:99") + + expect(opened.invalidCursor).toBe(true) + expect(opened.replay).toEqual([]) + expect(opened.fenceID).toBe("boot:1") + }) + + test("detects a gap when cursor is older than retained records", () => { + const store = new EventReplayStore({ bootID: "boot", maxRecords: 2, now: () => 1000 }) + store.append(question("q1")) + store.append(question("q2")) + store.append(question("q3")) + + const opened = store.snapshot("boot:0") + + expect(opened.gap).toBe(true) + expect(opened.replay.map((record) => record.id)).toEqual(["boot:2", "boot:3"]) + }) + + test("detects a gap when cursor is behind but all retained records expired", () => { + let now = 1000 + const store = new EventReplayStore({ bootID: "boot", maxAgeMs: 100, now: () => now }) + store.append(question("q1")) + now = 1200 + store.append(question("q2")) + now = 1400 + store.append(question("q3")) + now = 1600 + + const opened = store.snapshot("boot:1") + + expect(store.recordsForTest()).toEqual([]) + expect(opened.gap).toBe(true) + expect(opened.replay).toEqual([]) + expect(opened.fenceID).toBe("boot:3") + }) + + test("clears records for one disposed directory", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + store.append(question("q1")) + store.append({ ...question("q2"), directory: "/other" }) + + store.clearDirectory("/repo") + + expect(store.recordsForTest().map((record) => record.envelope.directory)).toEqual(["/other"]) + expect(store.latestID()).toBe("boot:2") + }) + + test("clears all retained records without changing the cursor generation", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + store.append(question("q1")) + store.append(question("q2")) + + store.clear() + + expect(store.recordsForTest()).toEqual([]) + expect(store.latestID()).toBe("boot:2") + }) + + test("reset starts a new empty generation", () => { + const store = new EventReplayStore({ bootID: "boot", now: () => 1000 }) + store.append(question("q1")) + + store.reset() + + expect(store.recordsForTest()).toEqual([]) + expect(store.latestID()).toMatch(/^[a-z0-9]+-[a-f0-9-]+:0$/) + }) +}) diff --git a/packages/opencode/test/server/global-event-replay.test.ts b/packages/opencode/test/server/global-event-replay.test.ts new file mode 100644 index 000000000..4f29f2102 --- /dev/null +++ b/packages/opencode/test/server/global-event-replay.test.ts @@ -0,0 +1,192 @@ +import { describe, expect, test } from "bun:test" +import { Hono } from "hono" +import { EventReplayStore } from "../../src/server/event-replay" +import { createGlobalEventReplayBridge, openGlobalEventReplayConnection, streamGlobalEvents } from "../../src/server/instance/global" + +const envelope = (type: string, id = type) => ({ + directory: "/repo", + payload: { + type, + properties: type === "question.asked" ? { id, sessionID: "ses_1", questions: [] } : { id }, + }, +}) + +async function readSseUntil(response: Response, predicate: (text: string) => boolean, timeoutMs = 2_000) { + const reader = response.body?.getReader() + if (!reader) throw new Error("Expected SSE response body") + const decoder = new TextDecoder() + let text = "" + const deadline = Date.now() + timeoutMs + try { + while (!predicate(text)) { + const remaining = deadline - Date.now() + if (remaining <= 0) throw new Error(`Timed out waiting for SSE payload. Received: ${text}`) + const next = await Promise.race([ + reader.read(), + new Promise((_, reject) => + setTimeout(() => reject(new Error(`Timed out reading SSE stream. Received: ${text}`)), remaining), + ), + ]) + if (next.done) break + text += decoder.decode(next.value, { stream: true }) + } + if (!predicate(text)) throw new Error(`SSE stream ended before expected payload. Received: ${text}`) + return text + } finally { + await reader.cancel() + } +} + +describe("createGlobalEventReplayBridge", () => { + test("live broadcasts replayable events with ids and non-replayable events without ids", () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + const packets: Array<{ id?: string; replaySeq?: number; data: string }> = [] + const unsubscribe = bridge.subscribe((packet) => packets.push(packet)) + + bridge.append(envelope("question.asked", "q1")) + bridge.append(envelope("message.part.delta", "delta1")) + + expect(packets[0].id).toBe("boot:1") + expect(packets[0].replaySeq).toBe(1) + expect(JSON.parse(packets[0].data).payload.type).toBe("question.asked") + expect(packets[1].id).toBeUndefined() + expect(packets[1].replaySeq).toBeUndefined() + expect(JSON.parse(packets[1].data).payload.type).toBe("message.part.delta") + unsubscribe() + }) + + test("connection seeds a fresh cursor and valid reconnect replays without server.connected", () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + bridge.append(envelope("question.asked", "q1")) + const fresh: Array<{ id?: string; replaySeq?: number; data: string }> = [] + + openGlobalEventReplayConnection({ bridge, push: (packet) => fresh.push(packet) })() + + expect(fresh[0].id).toBe("boot:1") + + bridge.append(envelope("question.replied", "q1")) + const reconnect: Array<{ id?: string; replaySeq?: number; data: string }> = [] + + openGlobalEventReplayConnection({ + bridge, + lastEventID: "boot:1", + push: (packet) => reconnect.push(packet), + })() + + expect(reconnect.map((packet) => JSON.parse(packet.data).payload.type)).toEqual(["question.replied"]) + expect(reconnect[0].id).toBe("boot:2") + }) + + test("invalid reconnect emits one server.connected with the fence id", () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + bridge.append(envelope("question.asked", "q1")) + const packets: Array<{ id?: string; replaySeq?: number; data: string }> = [] + + openGlobalEventReplayConnection({ + bridge, + lastEventID: "old:1", + push: (packet) => packets.push(packet), + })() + + expect(packets.map((packet) => JSON.parse(packet.data).payload.type)).toEqual(["server.connected"]) + expect(packets[0].id).toBe("boot:1") + }) + + test("gapped reconnect sends only server.connected and skips partial replay", () => { + const bridge = createGlobalEventReplayBridge({ + replayStore: new EventReplayStore({ bootID: "boot", maxRecords: 1 }), + }) + bridge.append(envelope("question.asked", "q1")) + bridge.append(envelope("question.replied", "q1")) + const packets: Array<{ id?: string; replaySeq?: number; data: string }> = [] + + openGlobalEventReplayConnection({ + bridge, + lastEventID: "boot:0", + push: (packet) => packets.push(packet), + })() + + expect(packets.map((packet) => JSON.parse(packet.data).payload.type)).toEqual(["server.connected"]) + expect(packets[0].id).toBe("boot:2") + }) + + test("dispose events invalidate retained replay records", () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + bridge.append(envelope("question.asked", "q1")) + bridge.append({ + directory: "/repo", + payload: { + type: "server.instance.disposed", + properties: { directory: "/repo" }, + }, + }) + const packets: Array<{ id?: string; replaySeq?: number; data: string }> = [] + + openGlobalEventReplayConnection({ + bridge, + lastEventID: "boot:0", + push: (packet) => packets.push(packet), + })() + + expect(packets.map((packet) => JSON.parse(packet.data).payload.type)).toEqual(["server.connected"]) + expect(packets[0].id).toBe("boot:2") + }) + + test("global dispose resets the replay generation and clears retained records", () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + bridge.append(envelope("question.asked", "q1")) + + const packet = bridge.append({ + payload: { + type: "global.disposed", + properties: {}, + }, + }) + + expect(JSON.parse(packet.data).payload.type).toBe("global.disposed") + expect(packet.id).toBeUndefined() + expect(bridge.replayStore.recordsForTest()).toEqual([]) + expect(bridge.replayStore.latestID()).not.toBe("boot:1") + }) + + test("missed instance dispose advances reconnect recovery from the previous fence", () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + bridge.append(envelope("question.asked", "q1")) + bridge.append({ + directory: "/repo", + payload: { + type: "server.instance.disposed", + properties: { directory: "/repo" }, + }, + }) + const packets: Array<{ id?: string; replaySeq?: number; data: string }> = [] + + openGlobalEventReplayConnection({ + bridge, + lastEventID: "boot:1", + push: (packet) => packets.push(packet), + })() + + expect(packets.map((packet) => JSON.parse(packet.data).payload.type)).toEqual(["server.instance.disposed"]) + expect(packets[0].id).toBe("boot:2") + }) + + test("route reads Last-Event-ID and writes replay ids into SSE output", async () => { + const bridge = createGlobalEventReplayBridge({ replayStore: new EventReplayStore({ bootID: "boot" }) }) + bridge.append(envelope("question.asked", "q1")) + const app = new Hono().get("/event", (c) => streamGlobalEvents(c, bridge)) + const controller = new AbortController() + + const response = await app.request("/event", { + headers: { "Last-Event-ID": "boot:0" }, + signal: controller.signal, + }) + expect(response.status).toBe(200) + const text = await readSseUntil(response, (value) => value.includes("question.asked")) + controller.abort() + + expect(text).toContain("id: boot:1") + expect(text).toContain("data: {\"directory\":\"/repo\",\"payload\":{\"type\":\"question.asked\"") + expect(text).not.toContain("server.connected") + }) +})