Răsfoiți Sursa

test(api): close stream transport coverage gaps

imccyu 1 lună în urmă
părinte
comite
ddcab34c0e

+ 125 - 15
apps/cli/tests/github-webhook-real.e2e.ts

@@ -2,7 +2,7 @@
 
 import type { ChildProcess } from 'node:child_process'
 import { spawn } from 'node:child_process'
-import { createHmac } from 'node:crypto'
+import { createHmac, randomUUID } from 'node:crypto'
 import { existsSync } from 'node:fs'
 import { mkdir, mkdtemp, realpath, rm } from 'node:fs/promises'
 import { createServer } from 'node:net'
@@ -33,7 +33,7 @@ interface SessionList {
   }>
 }
 
-interface WorkspaceList {
+interface WorkspaceBaseline {
   items: Array<{
     path: string
     sessionIds: string[]
@@ -109,28 +109,138 @@ async function freePort(): Promise<number> {
   return port
 }
 
-/** Invoke one public Web RPC method. */
-async function rpc<T>(baseUrl: string, method: string, payload: unknown): Promise<T> {
-  const response = await fetch(`${baseUrl}/api/${method}`, {
+/** Invoke one public Remote method over its HTTP carrier. */
+async function remoteRpc<T>(baseUrl: string, endpoint: string, args: object): Promise<T> {
+  const response = await fetch(`${baseUrl}/api/${endpoint}`, {
     method: 'POST',
     headers: { 'content-type': 'application/json' },
     body: JSON.stringify({
       type: 'client-request',
-      rpcId: `github-webhook-real-${method}`,
-      method,
-      payload,
+      rpcId: `github-webhook-real-${endpoint}-${randomUUID()}`,
+      method: endpoint,
+      payload: { args },
     }),
   })
-  if (!response.ok) throw new Error(`${method} returned HTTP ${String(response.status)}: ${await response.text()}`)
+  if (!response.ok) {
+    throw new Error(`${endpoint} returned HTTP ${String(response.status)}: ${await response.text()}`)
+  }
   const envelope = await response.json() as {
     result: { ok: true; value: T } | { ok: false; error: { code: string; message: string } }
   }
   if (!envelope.result.ok) {
-    throw new Error(`${method} failed: ${envelope.result.error.code}: ${envelope.result.error.message}`)
+    throw new Error(`${endpoint} failed: ${envelope.result.error.code}: ${envelope.result.error.message}`)
   }
   return envelope.result.value
 }
 
+/** Read one opening item from a public Remote stream. */
+async function openingStreamItem(
+  baseUrl: string,
+  endpoint: string,
+  args: object,
+  accepts: (value: unknown) => boolean,
+): Promise<Record<string, unknown>> {
+  const socket = new WebSocket(`${baseUrl.replace(/^http/u, 'ws')}/api/remote.mux`)
+  const streamId = `github-webhook-real-${endpoint}-${randomUUID()}`
+  try {
+    await new Promise<void>((resolve, reject) => {
+      const cleanup = (): void => {
+        socket.removeEventListener('open', opened)
+        socket.removeEventListener('error', failed)
+        socket.removeEventListener('close', closed)
+      }
+      const opened = (): void => {
+        cleanup()
+        resolve()
+      }
+      const failed = (): void => {
+        cleanup()
+        reject(new Error(`${endpoint} carrier failed before opening`))
+      }
+      const closed = (): void => {
+        cleanup()
+        reject(new Error(`${endpoint} carrier closed before opening`))
+      }
+      socket.addEventListener('open', opened)
+      socket.addEventListener('error', failed)
+      socket.addEventListener('close', closed)
+    })
+    return await new Promise<Record<string, unknown>>((resolve, reject) => {
+      const timer = setTimeout(() => { finish(new Error(`${endpoint} did not publish its opening item`)) }, 10_000)
+      const cleanup = (): void => {
+        clearTimeout(timer)
+        socket.removeEventListener('message', message)
+        socket.removeEventListener('error', failed)
+        socket.removeEventListener('close', closed)
+      }
+      const finish = (error: Error | undefined, value?: Record<string, unknown>): void => {
+        cleanup()
+        if (error !== undefined) reject(error)
+        else if (value === undefined) reject(new Error(`${endpoint} opening item was absent`))
+        else resolve(value)
+      }
+      const message = (event: MessageEvent<unknown>): void => {
+        try {
+          if (typeof event.data !== 'string') throw new Error(`${endpoint} published a non-text frame`)
+          const frame: unknown = JSON.parse(event.data)
+          if (!isRecord(frame) || frame.streamId !== streamId) return
+          if (frame.type === 'error') {
+            finish(new Error(`${endpoint} failed: ${JSON.stringify(frame.error)}`))
+            return
+          }
+          if (frame.type === 'end') {
+            finish(new Error(`${endpoint} ended before its opening item`))
+            return
+          }
+          if (frame.type === 'item' && isRecord(frame.value) && accepts(frame.value)) {
+            finish(undefined, frame.value)
+          }
+        } catch (error) {
+          finish(error instanceof Error ? error : new Error(String(error)))
+        }
+      }
+      const failed = (): void => { finish(new Error(`${endpoint} carrier failed before its opening item`)) }
+      const closed = (): void => { finish(new Error(`${endpoint} carrier closed before its opening item`)) }
+      socket.addEventListener('message', message)
+      socket.addEventListener('error', failed)
+      socket.addEventListener('close', closed)
+      socket.send(JSON.stringify({ type: 'open', streamId, endpoint, payload: { args } }))
+    })
+  } finally {
+    socket.close()
+  }
+}
+
+/** Read the current Workspace baseline from a fresh follow generation. */
+async function workspaceBaseline(baseUrl: string): Promise<WorkspaceBaseline> {
+  const frame = await openingStreamItem(
+    baseUrl,
+    'workspace/follow',
+    {},
+    value => isRecord(value) && value.type === 'baseline' && isRecord(value.value),
+  )
+  return frame.value as WorkspaceBaseline
+}
+
+/** Read the explicit page cut from a fresh Session follow generation. */
+async function sessionCursor(baseUrl: string, sessionId: string): Promise<number> {
+  const frame = await openingStreamItem(
+    baseUrl,
+    'session/follow',
+    { request: { address: { kind: 'session', sessionId } } },
+    value => isRecord(value) && value.type === 'opened' && Number.isSafeInteger(value.cursor),
+  )
+  return frame.cursor as number
+}
+
+/** Read Session history at the cursor explicitly opened for this page. */
+async function history(baseUrl: string, sessionId: string): Promise<HistoryPage> {
+  const throughSeq = await sessionCursor(baseUrl, sessionId)
+  return remoteRpc<HistoryPage>(baseUrl, 'session/page', {
+    request: { address: { kind: 'session', sessionId }, throughSeq, maxMessages: 100 },
+  })
+}
+
 /** Poll a public observation until it satisfies the test's behavior predicate. */
 async function eventually<T>(
   child: ChildProcess,
@@ -258,16 +368,16 @@ describe.skipIf(!process.env.DEEPSEEK_API_KEY)('GitHub webhook through the real
         child,
         observation.text,
         'one Workspace-attached Session',
-        async () => await rpc<WorkspaceList>(baseUrl, 'workspace.list', {}),
+        async () => await workspaceBaseline(baseUrl),
         value => value.items.some(workspace =>
           workspace.path === canonicalWorkspacePath && workspace.sessionIds.length === 1),
         30_000,
       )
       const workspace = workspaces.items.find(item => item.path === canonicalWorkspacePath)
       const sessionId = workspace?.sessionIds[0]
-      if (sessionId === undefined) throw new Error('workspace.list did not expose the webhook Session')
+      if (sessionId === undefined) throw new Error('workspace/follow did not expose the webhook Session')
 
-      const sessions = await rpc<SessionList>(baseUrl, 'session.list', {})
+      const sessions = await remoteRpc<SessionList>(baseUrl, 'session/list', { _request: {} })
       expect(sessions.items.find(session => session.sessionId === sessionId)).toMatchObject({
         agentPreset: 'minimal',
         blank: false,
@@ -278,7 +388,7 @@ describe.skipIf(!process.env.DEEPSEEK_API_KEY)('GitHub webhook through the real
         child,
         observation.text,
         'webhook provenance, title, and permission events',
-        async () => await rpc<HistoryPage>(baseUrl, 'session.history', { sessionId, maxMessages: 100 }),
+        async () => await history(baseUrl, sessionId),
         (page) => {
           const events = page.events.map(item => item.event)
           const title = events.find(event => event.type === 'session/title')
@@ -319,7 +429,7 @@ describe.skipIf(!process.env.DEEPSEEK_API_KEY)('GitHub webhook through the real
         child,
         observation.text,
         'a real DeepSeek assistant response',
-        async () => await rpc<HistoryPage>(baseUrl, 'session.history', { sessionId, maxMessages: 100 }),
+        async () => await history(baseUrl, sessionId),
         page => assistantText(page).includes(MARKER),
         150_000,
       )

+ 28 - 34
packages/api/gateway/src/client/journal-stream.ts

@@ -73,7 +73,6 @@ export interface RemoteJournalStreamOptions<Page, Entry, Cursor> {
 export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = void> {
   private readonly stream: RemoteStream<RemoteJournalFrame<Entry, Cursor>>
   private initialRequest!: PageRequest
-  private hasInitialRequest = false
   private resumeCursor: Cursor | undefined
   private hasResumeCursor = false
   private generation = 0
@@ -152,7 +151,6 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
     if (this.started) throw new Error(`${this.options.name} already opened`)
     this.started = true
     this.initialRequest = request
-    this.hasInitialRequest = true
     const iterator = this.stream[Symbol.asyncIterator]()
     try {
       const first = await this.takeNext(iterator)
@@ -233,7 +231,7 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
         if (item.value.type === 'opened') {
           throw new Error(`${this.options.name} emitted more than one opening cursor`)
         }
-        await this.acceptEntry(item, iterator)
+        await this.acceptEntry(item.value.entry, item, iterator)
       }
     } catch (error) {
       if (!this.disposed) this.options.failed(error)
@@ -285,32 +283,27 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
   }
 
   private async acceptEntry(
+    entry: Entry,
     item: JournalStreamItem<Entry, Cursor>,
     iterator: AsyncIterator<JournalStreamItem<Entry, Cursor>>,
   ): Promise<void> {
-    if (item.value.type !== 'entry') {
-      throw new Error(`${this.options.name} emitted more than one opening cursor`)
-    }
-    const entry = item.value.entry
     const cursor = this.options.cursor(entry)
-    const last = this.lastCursor
-    if (last !== undefined) {
-      if (this.options.compare(cursor, last) <= 0) return
-      if (!this.options.follows(last, cursor)) {
-        const request = this.repairPageRequest()
-        const superseded = await this.replaceThrough(
-          request,
-          cursor,
-          item.generation,
-          item.signal,
-          iterator,
-          [entry],
-        )
-        if (superseded !== undefined) {
-          await this.replaceGeneration(request, superseded, iterator, true)
-        }
-        return
+    const last = this.lastCursor as Cursor
+    if (this.options.compare(cursor, last) <= 0) return
+    if (!this.options.follows(last, cursor)) {
+      const request = this.repairPageRequest()
+      const superseded = await this.replaceThrough(
+        request,
+        cursor,
+        item.generation,
+        item.signal,
+        iterator,
+        [entry],
+      )
+      if (superseded !== undefined) {
+        await this.replaceGeneration(request, superseded, iterator, true)
       }
+      return
     }
     if (this.firstCursor === undefined) this.firstCursor = cursor
     this.lastCursor = cursor
@@ -400,7 +393,7 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
         if (!signal.aborted || this.stream.signal.aborted) throw result.error
         return this.awaitReplacementGeneration(generation, iterator, pending)
       }
-      this.releaseNext(pending)
+      this.releaseNext()
       if (result.type === 'next-error') throw result.error
       if (result.value.done) {
         signal.throwIfAborted()
@@ -426,7 +419,7 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
       try {
         next = await pending
       } finally {
-        this.releaseNext(pending)
+        this.releaseNext()
       }
       if (next.done) {
         this.stream.signal.throwIfAborted()
@@ -481,16 +474,15 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
     try {
       return await pending
     } finally {
-      this.releaseNext(pending)
+      this.releaseNext()
     }
   }
 
-  private releaseNext(pending: Promise<IteratorResult<JournalStreamItem<Entry, Cursor>>>): void {
-    if (this.pendingNext === pending) this.pendingNext = undefined
+  private releaseNext(): void {
+    this.pendingNext = undefined
   }
 
   private repairPageRequest(): PageRequest {
-    if (!this.hasInitialRequest) throw new Error(`${this.options.name} has no initial page request`)
     return this.repairRequest(this.initialRequest)
   }
 
@@ -509,13 +501,15 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
   }
 
   private assertPage(entries: readonly Entry[]): void {
-    for (let index = 1; index < entries.length; index++) {
-      const previous = entries[index - 1]
-      const entry = entries[index]
-      if (previous === undefined || entry === undefined) continue
+    const iterator = entries[Symbol.iterator]()
+    const first = iterator.next()
+    if (first.done) return
+    let previous = first.value
+    for (const entry of iterator) {
       if (!this.options.follows(this.options.cursor(previous), this.options.cursor(entry))) {
         throw new Error(`${this.options.name} page contains discontinuous entries`)
       }
+      previous = entry
     }
   }
 

+ 50 - 6
packages/api/gateway/tests/control-retry.client.spec.ts

@@ -37,7 +37,7 @@ function hostSource(initiallyAvailable: boolean): {
 }
 
 interface Generation<Item> {
-  readonly values?: readonly Item[]
+  readonly values?: readonly (Item | Promise<Item>)[]
   readonly terminal?: Error
   readonly hold?: boolean
   readonly afterAbortError?: Error
@@ -51,7 +51,7 @@ function scripted<Item>(generations: Generation<Item>[], opened?: () => void) {
       if (generation === undefined) throw new Error('fixture has no stream generation')
       opened?.()
       try {
-        for (const value of generation.values ?? []) yield value
+        for (const value of generation.values ?? []) yield await value
         if (generation.terminal !== undefined) throw generation.terminal
         if (generation.hold === true && !signal.aborted) {
           await new Promise<void>((resolve) => {
@@ -118,9 +118,21 @@ describe('RemoteStream', () => {
   })
 
   it('waits for a replacement Host generation after observing unavailability', async () => {
-    const source = hostSource(false)
+    let available = false
+    let listener: (() => void) | undefined
+    const subscribed = Promise.withResolvers<undefined>()
+    const connection = {
+      hostDescription: {
+        getSnapshot: () => available ? DESCRIPTION : undefined,
+        subscribe: (value: () => void) => {
+          listener = value
+          subscribed.resolve(undefined)
+          return () => { listener = undefined }
+        },
+      },
+    }
     let opened = 0
-    const stream = new RemoteStream(source.connection, {
+    const stream = new RemoteStream(connection, {
       name: 'fixture stream',
       open: scripted([
         { terminal: new RemoteStreamCarrierError('offline') },
@@ -130,10 +142,12 @@ describe('RemoteStream', () => {
     })
     const pending = stream[Symbol.asyncIterator]().next()
     await vi.waitFor(() => { expect(opened).toBe(1) })
+    await subscribed.promise
 
-    source.publish(false)
+    listener?.()
     expect(opened).toBe(1)
-    source.publish(true)
+    available = true
+    listener?.()
     await expect(pending).resolves.toMatchObject({
       done: false,
       value: { generation: 2, value: 'ready' },
@@ -141,6 +155,24 @@ describe('RemoteStream', () => {
     await stream.dispose()
   })
 
+  it('stops a pending retry when the logical stream is disposed', async () => {
+    const source = hostSource(false)
+    let opened = 0
+    const stream = new RemoteStream(source.connection, {
+      name: 'fixture stream',
+      open: scripted([
+        { terminal: new RemoteStreamCarrierError('offline') },
+      ], () => { opened++ }),
+      ended: () => new Error('ended'),
+    })
+    const pending = stream[Symbol.asyncIterator]().next()
+    await vi.waitFor(() => { expect(opened).toBe(1) })
+    source.publish(false)
+
+    await stream.dispose()
+    await expect(pending).resolves.toEqual({ done: true, value: undefined })
+  })
+
   it('contains a Host publication during subscription setup', async () => {
     let reads = 0
     let disposed = 0
@@ -300,4 +332,16 @@ describe('RemoteStream', () => {
       value: undefined,
     })
   })
+
+  it('drops a value that arrives after disposal begins', async () => {
+    const source = hostSource(true)
+    const late = Promise.withResolvers<string>()
+    const stream = supervisor(source.connection, [{ values: [late.promise] }])
+    const pending = stream[Symbol.asyncIterator]().next()
+    const disposing = stream.dispose()
+    late.resolve('late')
+
+    await expect(pending).resolves.toEqual({ done: true, value: undefined })
+    await disposing
+  })
 })

+ 396 - 3
packages/api/gateway/tests/journal-stream.client.spec.ts

@@ -5,6 +5,8 @@ import {
   RemoteStreamCarrierError,
   type RemoteJournalChange,
   type RemoteJournalFrame,
+  type RemoteStreamFactory,
+  type RemoteStreamItem,
   type RemoteStreamOptions,
 } from '../src/client/index.ts'
 
@@ -24,9 +26,13 @@ interface PageRequest {
 }
 
 interface Generation {
-  readonly frames: readonly RemoteJournalFrame<Entry, number>[]
+  readonly frames: readonly (
+    RemoteJournalFrame<Entry, number> | Promise<RemoteJournalFrame<Entry, number>>
+  )[]
   readonly terminal?: Error
   readonly hold?: boolean
+  readonly waitAfterFrames?: Promise<void>
+  readonly afterFrame?: (index: number) => void
 }
 
 type PageSource = Page | Promise<Page> | ((signal: AbortSignal) => Promise<Page>)
@@ -64,8 +70,9 @@ class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageReques
     private readonly followCursors: (number | undefined)[],
     changes: RemoteJournalChange<Page, Entry>[],
     failed: (error: unknown) => void,
+    factory: RemoteStreamFactory = STREAM_FACTORY,
   ) {
-    super(STREAM_FACTORY, {
+    super(factory, {
       name: 'fixture journal',
       emptyCursor: -1,
       entries: value => value.entries,
@@ -87,7 +94,11 @@ class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageReques
     this.followCursors.push(after)
     const generation = this.generations.shift()
     if (generation === undefined) throw new Error('no scripted journal generation')
-    for (const frame of generation.frames) yield frame
+    for (const [index, frame] of generation.frames.entries()) {
+      yield await frame
+      generation.afterFrame?.(index)
+    }
+    await generation.waitAfterFrames
     if (generation.terminal !== undefined) throw generation.terminal
     if (generation.hold === true && !signal.aborted) {
       await new Promise<void>((resolve) => {
@@ -119,6 +130,7 @@ class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageReques
 function journalFixture(
   generations: Generation[],
   pages: PageSource[],
+  factory: RemoteStreamFactory = STREAM_FACTORY,
 ): {
   readonly journal: RemoteJournalStream<Page, Entry, number, PageRequest>
   readonly changes: RemoteJournalChange<Page, Entry>[]
@@ -143,10 +155,39 @@ function journalFixture(
     followCursors,
     changes,
     failed,
+    factory,
   )
   return { journal, changes, failed, calls, pageRequests, pageCursors, followCursors }
 }
 
+function remoteItem(
+  generation: number,
+  value: RemoteJournalFrame<Entry, number>,
+  signal: AbortSignal,
+): RemoteStreamItem<RemoteJournalFrame<Entry, number>> {
+  return { generation, value, signal, accept: vi.fn() }
+}
+
+function controlledFactory(
+  next: () => Promise<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>,
+): RemoteStreamFactory {
+  const lifetime = new AbortController()
+  return {
+    $stream<Item>(): RemoteStream<Item> {
+      const iterator = {
+        next,
+        return: async () => ({ done: true as const, value: undefined }),
+      }
+      return {
+        signal: lifetime.signal,
+        restart: () => {},
+        dispose: async () => { lifetime.abort() },
+        [Symbol.asyncIterator]: () => iterator,
+      } as unknown as RemoteStream<Item>
+    },
+  }
+}
+
 describe('RemoteJournalStream', () => {
   it('opens follow before page, removes overlap, appends live entries, and prepends history', async () => {
     const fixture = journalFixture(
@@ -177,6 +218,69 @@ describe('RemoteJournalStream', () => {
     await fixture.journal.dispose()
   })
 
+  it('exposes its shared cancellation signal', async () => {
+    const fixture = journalFixture(
+      [{ frames: [{ type: 'opened', cursor: -1 }], hold: true }],
+      [page('empty', [])],
+    )
+
+    expect(fixture.journal.signal.aborted).toBe(false)
+    await fixture.journal.open({})
+    await fixture.journal.dispose()
+    expect(fixture.journal.signal.aborted).toBe(true)
+  })
+
+  it('classifies normal endings before initial and resumed opening cursors', async () => {
+    const initial = journalFixture([{ frames: [] }], [])
+    await expect(initial.journal.open({})).rejects.toThrow(
+      'fixture journal ended before its opening cursor',
+    )
+
+    const finish = Promise.withResolvers<undefined>()
+    const resumed = journalFixture(
+      [
+        { frames: [{ type: 'opened', cursor: 0 }], waitAfterFrames: finish.promise },
+        { frames: [] },
+      ],
+      [page('initial', [0])],
+    )
+    await resumed.journal.open({})
+    finish.resolve(undefined)
+    await vi.waitFor(() => { expect(resumed.failed).toHaveBeenCalledOnce() })
+    expect(resumed.failed.mock.calls[0]?.[0]).toMatchObject({
+      message: 'resumed fixture journal ended before its opening cursor',
+    })
+    await resumed.journal.dispose()
+  })
+
+  it('prepends into an empty window and accepts its first live entry', async () => {
+    const empty = journalFixture(
+      [{ frames: [{ type: 'opened', cursor: -1 }], hold: true }],
+      [page('empty', []), page('older', [0]), page('oldest', [])],
+    )
+    await empty.journal.open({})
+    await empty.journal.prepend({})
+    expect(empty.changes.at(-1)).toEqual({
+      type: 'prepend', page: page('older', [0]), entries: entries(0), hasMore: false,
+    })
+    await empty.journal.prepend({})
+    expect(empty.changes.at(-1)).toEqual({
+      type: 'prepend', page: page('oldest', []), entries: [], hasMore: false,
+    })
+    await empty.journal.dispose()
+
+    const live = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const followed = journalFixture(
+      [{ frames: [{ type: 'opened', cursor: -1 }, live.promise], hold: true }],
+      [page('empty', [])],
+    )
+    await followed.journal.open({})
+    live.resolve({ type: 'entry', entry: { seq: 0 } })
+    await vi.waitFor(() => { expect(followed.changes).toHaveLength(2) })
+    expect(followed.changes.at(-1)).toEqual({ type: 'append', entry: { seq: 0 } })
+    await followed.journal.dispose()
+  })
+
   it('publishes one sorted replacement from an exact page and live entries queued while it loads', async () => {
     let resolvePage!: (value: Page) => void
     const openingPage = new Promise<Page>((resolve) => { resolvePage = resolve })
@@ -304,6 +408,295 @@ describe('RemoteJournalStream', () => {
     await fixture.journal.dispose()
   })
 
+  it('replaces a superseded live-gap repair with the next generation', async () => {
+    const gap = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const fixture = journalFixture(
+      [
+        {
+          frames: [{ type: 'opened', cursor: 1 }, gap.promise],
+          terminal: new RemoteStreamCarrierError('generation lost'),
+        },
+        { frames: [{ type: 'opened', cursor: 4 }], hold: true },
+      ],
+      [
+        page('initial', [0, 1]),
+        () => new Promise<Page>(() => {}),
+        page('replacement', [0, 1, 2, 3, 4]),
+      ],
+    )
+
+    await fixture.journal.open({ limit: 5 })
+    gap.resolve({ type: 'entry', entry: { seq: 4 } })
+    await vi.waitFor(() => { expect(fixture.changes).toHaveLength(2) })
+    expect(fixture.changes.at(-1)).toMatchObject({
+      type: 'replace', page: { marker: 'replacement' }, entries: entries(0, 1, 2, 3, 4),
+    })
+    await fixture.journal.dispose()
+  })
+
+  it('replaces a superseded second repair page with the next generation', async () => {
+    const live = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const liveConsumed = Promise.withResolvers<undefined>()
+    const openingPage = Promise.withResolvers<Page>()
+    const finish = Promise.withResolvers<undefined>()
+    const fixture = journalFixture(
+      [
+        {
+          frames: [{ type: 'opened', cursor: 1 }, live.promise],
+          waitAfterFrames: finish.promise,
+          terminal: new RemoteStreamCarrierError('generation lost'),
+          afterFrame: (index) => { if (index === 1) liveConsumed.resolve(undefined) },
+        },
+        { frames: [{ type: 'opened', cursor: 4 }], hold: true },
+      ],
+      [
+        openingPage.promise,
+        () => new Promise<Page>(() => {}),
+        page('replacement', [0, 1, 2, 3, 4]),
+      ],
+    )
+
+    const opening = fixture.journal.open({})
+    await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([1]) })
+    live.resolve({ type: 'entry', entry: { seq: 3 } })
+    await liveConsumed.promise
+    openingPage.resolve(page('opening', [0, 1]))
+    await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([1, 3]) })
+    finish.resolve(undefined)
+    await opening
+
+    expect(fixture.pageCursors).toEqual([1, 3, 4])
+    expect(fixture.changes).toEqual([{
+      type: 'replace',
+      page: page('replacement', [0, 1, 2, 3, 4]),
+      entries: entries(0, 1, 2, 3, 4),
+      hasMore: false,
+    }])
+    await fixture.journal.dispose()
+  })
+
+  it('rereads the tail when queued entries advance beyond the opening page', async () => {
+    const live = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const liveConsumed = Promise.withResolvers<undefined>()
+    const openingPage = Promise.withResolvers<Page>()
+    const fixture = journalFixture(
+      [{
+        frames: [{ type: 'opened', cursor: 1 }, live.promise],
+        hold: true,
+        afterFrame: (index) => { if (index === 1) liveConsumed.resolve(undefined) },
+      }],
+      [openingPage.promise, page('repair', [0, 1, 2, 3])],
+    )
+
+    const opening = fixture.journal.open({ limit: 4 })
+    await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([1]) })
+    live.resolve({ type: 'entry', entry: { seq: 3 } })
+    await liveConsumed.promise
+    openingPage.resolve(page('opening', [0, 1]))
+    await opening
+
+    expect(fixture.pageCursors).toEqual([1, 3])
+    expect(fixture.changes).toEqual([{
+      type: 'replace', page: page('repair', [0, 1, 2, 3]), entries: entries(0, 1, 2, 3), hasMore: false,
+    }])
+    await fixture.journal.dispose()
+  })
+
+  it('rejects when queued entries advance beyond the second repair page', async () => {
+    const firstLive = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const secondLive = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const firstConsumed = Promise.withResolvers<undefined>()
+    const secondConsumed = Promise.withResolvers<undefined>()
+    const openingPage = Promise.withResolvers<Page>()
+    const repairPage = Promise.withResolvers<Page>()
+    const fixture = journalFixture(
+      [{
+        frames: [{ type: 'opened', cursor: 1 }, firstLive.promise, secondLive.promise],
+        hold: true,
+        afterFrame: (index) => {
+          if (index === 1) firstConsumed.resolve(undefined)
+          if (index === 2) secondConsumed.resolve(undefined)
+        },
+      }],
+      [openingPage.promise, repairPage.promise],
+    )
+
+    const opening = fixture.journal.open({})
+    await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([1]) })
+    firstLive.resolve({ type: 'entry', entry: { seq: 3 } })
+    await firstConsumed.promise
+    openingPage.resolve(page('opening', [0, 1]))
+    await vi.waitFor(() => { expect(fixture.pageCursors).toEqual([1, 3]) })
+    secondLive.resolve({ type: 'entry', entry: { seq: 5 } })
+    await secondConsumed.promise
+    repairPage.resolve(page('repair', [0, 1, 2, 3]))
+
+    await expect(opening).rejects.toThrow('page did not reach its opening cursor')
+  })
+
+  it('reports a resumed generation that emits an entry before its cursor', async () => {
+    const finish = Promise.withResolvers<undefined>()
+    const fixture = journalFixture(
+      [
+        {
+          frames: [{ type: 'opened', cursor: 0 }],
+          waitAfterFrames: finish.promise,
+          terminal: new RemoteStreamCarrierError('lost'),
+        },
+        { frames: [{ type: 'entry', entry: { seq: 1 } }] },
+      ],
+      [page('initial', [0])],
+    )
+
+    await fixture.journal.open({})
+    finish.resolve(undefined)
+    await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
+    expect(fixture.failed.mock.calls[0]?.[0]).toMatchObject({
+      message: 'resumed fixture journal emitted an entry before its opening cursor',
+    })
+    await fixture.journal.dispose()
+  })
+
+  it('reports a duplicate opening cursor after the initial page is published', async () => {
+    const duplicate = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const fixture = journalFixture(
+      [{ frames: [{ type: 'opened', cursor: 0 }, duplicate.promise], hold: true }],
+      [page('initial', [0])],
+    )
+
+    await fixture.journal.open({})
+    duplicate.resolve({ type: 'opened', cursor: 0 })
+    await vi.waitFor(() => { expect(fixture.failed).toHaveBeenCalledOnce() })
+    expect(fixture.failed.mock.calls[0]?.[0]).toMatchObject({
+      message: 'fixture journal emitted more than one opening cursor',
+    })
+    await fixture.journal.dispose()
+  })
+
+  it('propagates follow failures and duplicate cursors while an opening page is pending', async () => {
+    const pendingPage = new Promise<Page>(() => {})
+    const failedFollow = journalFixture(
+      [{ frames: [{ type: 'opened', cursor: 0 }], terminal: new Error('follow failed') }],
+      [pendingPage],
+    )
+    await expect(failedFollow.journal.open({})).rejects.toThrow('follow failed')
+
+    const duplicate = Promise.withResolvers<RemoteJournalFrame<Entry, number>>()
+    const duplicatePage = new Promise<Page>(() => {})
+    const duplicateOpening = journalFixture(
+      [{ frames: [{ type: 'opened', cursor: 0 }, duplicate.promise] }],
+      [duplicatePage],
+    )
+    const opening = duplicateOpening.journal.open({})
+    await vi.waitFor(() => { expect(duplicateOpening.pageCursors).toEqual([0]) })
+    duplicate.resolve({ type: 'opened', cursor: 0 })
+    await expect(opening).rejects.toThrow('more than one opening cursor')
+  })
+
+  it('rejects an iterator that ends while its opening page is pending', async () => {
+    const generation = new AbortController()
+    const results = [
+      Promise.resolve<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>({
+        done: false,
+        value: remoteItem(1, { type: 'opened', cursor: 0 }, generation.signal),
+      }),
+      Promise.resolve<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>({
+        done: true,
+        value: undefined,
+      }),
+    ]
+    const fixture = journalFixture(
+      [],
+      [new Promise<Page>(() => {})],
+      controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
+    )
+
+    await expect(fixture.journal.open({})).rejects.toThrow(
+      'ended while reading its replacement page',
+    )
+  })
+
+  it('rejects an iterator that ends before its opening cursor', async () => {
+    const factory = controlledFactory(() => Promise.resolve({ done: true, value: undefined }))
+    const fixture = journalFixture([], [], factory)
+
+    await expect(fixture.journal.open({})).rejects.toThrow(
+      'ended before its opening cursor',
+    )
+  })
+
+  it('suppresses a consumer failure after disposal begins', async () => {
+    const generation = new AbortController()
+    const next = Promise.withResolvers<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>()
+    const results = [
+      Promise.resolve<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>({
+        done: false,
+        value: remoteItem(1, { type: 'opened', cursor: 0 }, generation.signal),
+      }),
+      next.promise,
+    ]
+    const fixture = journalFixture(
+      [],
+      [page('initial', [0])],
+      controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
+    )
+
+    await fixture.journal.open({})
+    const closing = fixture.journal.dispose()
+    next.resolve({
+      done: false,
+      value: remoteItem(1, { type: 'opened', cursor: 0 }, generation.signal),
+    })
+    await closing
+    expect(fixture.failed).not.toHaveBeenCalled()
+  })
+
+  it.each([
+    { name: 'ends', final: { done: true as const, value: undefined }, message: 'ended while replacing' },
+    {
+      name: 'emits another opening cursor',
+      final: undefined,
+      message: 'more than one opening cursor',
+    },
+  ])('rejects when an aborted page generation $name', async ({ final, message }) => {
+    const generation = new AbortController()
+    const pending = Promise.withResolvers<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>()
+    const nextPending = Promise.withResolvers<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>()
+    const results = [
+      Promise.resolve<IteratorResult<RemoteStreamItem<RemoteJournalFrame<Entry, number>>>>({
+        done: false,
+        value: remoteItem(1, { type: 'opened', cursor: 0 }, generation.signal),
+      }),
+      pending.promise,
+      nextPending.promise,
+    ]
+    const fixture = journalFixture(
+      [],
+      [signal => new Promise<Page>((_resolve, reject) => {
+        signal.addEventListener('abort', () => { reject(new Error('page aborted')) }, { once: true })
+      })],
+      controlledFactory(() => results.shift() ?? Promise.resolve({ done: true, value: undefined })),
+    )
+
+    const opening = fixture.journal.open({})
+    await vi.waitFor(() => { expect(results).toHaveLength(1) })
+    generation.abort()
+    if (final === undefined) {
+      pending.resolve({
+        done: false,
+        value: remoteItem(1, { type: 'entry', entry: { seq: 1 } }, generation.signal),
+      })
+      await vi.waitFor(() => { expect(results).toHaveLength(0) })
+      nextPending.resolve({
+        done: false,
+        value: remoteItem(1, { type: 'opened', cursor: 1 }, generation.signal),
+      })
+    } else {
+      pending.resolve(final)
+    }
+    await expect(opening).rejects.toThrow(message)
+  })
+
   it('rejects malformed opening and page sequences', async () => {
     const beforeOpening = journalFixture(
       [{ frames: [{ type: 'entry', entry: { seq: 0 } }] }],

+ 15 - 0
packages/api/session-controller/tests/control-queue.host.spec.ts

@@ -102,4 +102,19 @@ describe('Session control queue projection', () => {
 
     await expect(waiting).resolves.toMatchObject({ done: true })
   })
+
+  it('ends active streams on context disposal after flushing buffered frames', async () => {
+    const { ctx, control, inbox } = await harness()
+    const iterator = control.control(new AbortController().signal)[Symbol.asyncIterator]()
+    await iterator.next()
+    inbox.append('next-turn', message('first'))
+    inbox.append('next-turn', message('second'))
+
+    const first = await iterator.next()
+    expect(first).toMatchObject({ done: false, value: { type: 'queue' } })
+    await ctx.fiber.dispose()
+    const second = await iterator.next()
+    expect(second).toMatchObject({ done: false, value: { type: 'queue' } })
+    await expect(iterator.next()).resolves.toMatchObject({ done: true })
+  })
 })

+ 75 - 4
packages/api/session-controller/tests/transport.client.spec.ts

@@ -57,6 +57,7 @@ interface FollowGeneration {
   readonly frames: readonly SessionFollowFrame[]
   readonly terminal?: Error
   readonly hold?: boolean
+  readonly waitAfterFrames?: Promise<void>
 }
 
 class ScriptedSessionRemote implements SessionTransportRemote {
@@ -68,6 +69,7 @@ class ScriptedSessionRemote implements SessionTransportRemote {
     private readonly generations: FollowGeneration[],
     private readonly pages: RemoteResult<SessionPage>[],
     private readonly controlFrames: readonly SessionControlFrame[] = [],
+    private readonly holdControl = true,
   ) {}
 
   async *follow(request: SessionFollowRequest, signal = new AbortController().signal): AsyncIterable<SessionFollowFrame> {
@@ -76,6 +78,7 @@ class ScriptedSessionRemote implements SessionTransportRemote {
     this.followRequests.push(request)
     this.signals.push(signal)
     for (const frame of generation.frames) yield frame
+    await generation.waitAfterFrames
     if (generation.terminal !== undefined) throw generation.terminal
     if (generation.hold === true && !signal.aborted) {
       await new Promise<void>((resolve) => {
@@ -93,7 +96,7 @@ class ScriptedSessionRemote implements SessionTransportRemote {
 
   async *control(signal = new AbortController().signal): AsyncIterable<SessionControlFrame> {
     for (const frame of this.controlFrames) yield frame
-    if (!signal.aborted) {
+    if (this.holdControl && !signal.aborted) {
       await new Promise<void>((resolve) => {
         signal.addEventListener('abort', () => { resolve() }, { once: true })
       })
@@ -168,7 +171,7 @@ describe('Session Client stream adapters', () => {
       failed: vi.fn(),
     })
 
-    await stream.open({})
+    await stream.open({ maxMessages: 50 })
     await vi.waitFor(() => { expect(remote.followRequests).toHaveLength(2) })
 
     expect(remote.followRequests).toEqual([
@@ -176,14 +179,45 @@ describe('Session Client stream adapters', () => {
       { address: ADDRESS, afterSeq: 2 },
     ])
     expect(remote.pageRequests).toEqual([
-      { address: ADDRESS, throughSeq: 1 },
-      { address: ADDRESS, throughSeq: 4 },
+      { address: ADDRESS, throughSeq: 1, maxMessages: 50 },
+      { address: ADDRESS, throughSeq: 4, maxMessages: 50 },
     ])
     expect(changes.map(change => change.type)).toEqual(['replace', 'append', 'replace'])
     expect(carrierFailed).toHaveBeenCalledWith(lost)
     await stream.dispose()
   })
 
+  it('repairs a resumed event stream without an optional message limit', async () => {
+    const finish = Promise.withResolvers<undefined>()
+    const remote = new ScriptedSessionRemote(
+      [
+        {
+          frames: [{ type: 'opened', cursor: 0 }],
+          waitAfterFrames: finish.promise,
+          terminal: new RemoteStreamCarrierError('lost'),
+        },
+        { frames: [{ type: 'opened', cursor: 1 }], hold: true },
+      ],
+      [
+        { ok: true, value: page([entry(0)]) },
+        { ok: true, value: page([entry(0), entry(1)]) },
+      ],
+    )
+    const stream = new SessionEventStream(sessionClient(remote), ADDRESS, {
+      publish: vi.fn(),
+      failed: vi.fn(),
+    })
+
+    await stream.open({})
+    finish.resolve(undefined)
+    await vi.waitFor(() => { expect(remote.pageRequests).toHaveLength(2) })
+    expect(remote.pageRequests).toEqual([
+      { address: ADDRESS, throughSeq: 0 },
+      { address: ADDRESS, throughSeq: 1 },
+    ])
+    await stream.dispose()
+  })
+
   it('turns a page failure into a typed stream failure and closes follow', async () => {
     const failure = { code: 'session-not-found', message: 'missing', details: { sessionId: 'session-1' } } as const
     const remote = new ScriptedSessionRemote(
@@ -226,4 +260,41 @@ describe('Session Client stream adapters', () => {
     await stream.dispose()
     await stream.dispose()
   })
+
+  it('classifies control streams that end before and after their opening baseline', async () => {
+    const beforeFailed = vi.fn()
+    const before = createSessionControlStream(
+      sessionClient(new ScriptedSessionRemote([], [], [], false)),
+      { accept: vi.fn(), failed: beforeFailed },
+    )
+    before.start()
+    await vi.waitFor(() => { expect(beforeFailed).toHaveBeenCalledOnce() })
+    expect(beforeFailed.mock.calls[0]?.[0]).toMatchObject({
+      message: 'session control stream ended before its opening snapshot',
+    })
+    await before.dispose()
+
+    const baseline: SessionControlFrame = {
+      type: 'baseline',
+      value: { queues: {}, jobs: {}, projections: {} },
+    }
+    const carrierFailed = vi.fn()
+    const failed = vi.fn()
+    const afterRemote = new ScriptedSessionRemote([], [], [baseline], false)
+    const after = createSessionControlStream(sessionClient(afterRemote), {
+      accept: vi.fn(),
+      carrierFailed: (error) => {
+        carrierFailed(error)
+        void after.dispose()
+      },
+      failed,
+    })
+    after.start()
+    await vi.waitFor(() => { expect(carrierFailed).toHaveBeenCalledOnce() })
+    expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
+      message: 'session control stream ended without a terminal result',
+    })
+    expect(failed).not.toHaveBeenCalled()
+    await after.dispose()
+  })
 })

+ 43 - 0
packages/api/workspace-controller/tests/transport.client.spec.ts

@@ -192,6 +192,49 @@ describe('Workspace Client snapshot adapter', () => {
     await stream.dispose()
   })
 
+  it('classifies a normal end after the opening baseline as carrier loss', async () => {
+    const remote = new ScriptedWorkspaceRemote([
+      { frames: [baseline('old')] },
+      { frames: [baseline('fresh')], hold: true },
+    ])
+    const replaceBaseline = vi.fn<WorkspaceFollowSink['replaceBaseline']>()
+    const carrierFailed = vi.fn()
+    const stream = createWorkspaceStateStream(workspaceClient(remote), {
+      accept: accepts({ replaceBaseline }),
+      carrierFailed,
+      failed: vi.fn(),
+    })
+
+    stream.start()
+    await vi.waitFor(() => { expect(replaceBaseline).toHaveBeenCalledTimes(2) })
+    expect(carrierFailed.mock.calls[0]?.[0]).toMatchObject({
+      message: 'Workspace state stream ended without a terminal result',
+    })
+    await stream.dispose()
+  })
+
+  it('suppresses callback failure after disposal begins', async () => {
+    const failed = vi.fn()
+    let closing: Promise<void> | undefined
+    const stream = createWorkspaceStateStream(
+      workspaceClient(new ScriptedWorkspaceRemote([{ frames: [baseline()] }])),
+      {
+        accept: accepts({
+          replaceBaseline: () => {
+            closing = stream.dispose()
+            throw new Error('disposed callback')
+          },
+        }),
+        failed,
+      },
+    )
+
+    stream.start()
+    await vi.waitFor(() => { expect(closing).toBeDefined() })
+    await closing
+    expect(failed).not.toHaveBeenCalled()
+  })
+
   it.each([
     {
       name: 'an increment before the baseline',