diff --git a/apps/cli/tests/github-webhook-real.e2e.ts b/apps/cli/tests/github-webhook-real.e2e.ts index a49c529aa5..93421ea93a 100644 --- a/apps/cli/tests/github-webhook-real.e2e.ts +++ b/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 { return port } -/** Invoke one public Web RPC method. */ -async function rpc(baseUrl: string, method: string, payload: unknown): Promise { - const response = await fetch(`${baseUrl}/api/${method}`, { +/** Invoke one public Remote method over its HTTP carrier. */ +async function remoteRpc(baseUrl: string, endpoint: string, args: object): Promise { + 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> { + const socket = new WebSocket(`${baseUrl.replace(/^http/u, 'ws')}/api/remote.mux`) + const streamId = `github-webhook-real-${endpoint}-${randomUUID()}` + try { + await new Promise((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>((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): 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): 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 { + 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 { + 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 { + const throughSeq = await sessionCursor(baseUrl, sessionId) + return remoteRpc(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( 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(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(baseUrl, 'session.list', {}) + const sessions = await remoteRpc(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(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(baseUrl, 'session.history', { sessionId, maxMessages: 100 }), + async () => await history(baseUrl, sessionId), page => assistantText(page).includes(MARKER), 150_000, ) diff --git a/packages/api/gateway/src/client/journal-stream.ts b/packages/api/gateway/src/client/journal-stream.ts index 9b64cc5605..f21e79f241 100644 --- a/packages/api/gateway/src/client/journal-stream.ts +++ b/packages/api/gateway/src/client/journal-stream.ts @@ -73,7 +73,6 @@ export interface RemoteJournalStreamOptions { export abstract class RemoteJournalStream { private readonly stream: RemoteStream> private initialRequest!: PageRequest - private hasInitialRequest = false private resumeCursor: Cursor | undefined private hasResumeCursor = false private generation = 0 @@ -152,7 +151,6 @@ export abstract class RemoteJournalStream, iterator: AsyncIterator>, ): Promise { - 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>>): 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 { - readonly values?: readonly Item[] + readonly values?: readonly (Item | Promise)[] readonly terminal?: Error readonly hold?: boolean readonly afterAbortError?: Error @@ -51,7 +51,7 @@ function scripted(generations: Generation[], 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((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() + 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() + 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 + }) }) diff --git a/packages/api/gateway/tests/journal-stream.client.spec.ts b/packages/api/gateway/tests/journal-stream.client.spec.ts index 17d02cf8c5..38362c5598 100644 --- a/packages/api/gateway/tests/journal-stream.client.spec.ts +++ b/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[] + readonly frames: readonly ( + RemoteJournalFrame | Promise> + )[] readonly terminal?: Error readonly hold?: boolean + readonly waitAfterFrames?: Promise + readonly afterFrame?: (index: number) => void } type PageSource = Page | Promise | ((signal: AbortSignal) => Promise) @@ -64,8 +70,9 @@ class FixtureJournal extends RemoteJournalStream[], 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((resolve) => { @@ -119,6 +130,7 @@ class FixtureJournal extends RemoteJournalStream readonly changes: RemoteJournalChange[] @@ -143,10 +155,39 @@ function journalFixture( followCursors, changes, failed, + factory, ) return { journal, changes, failed, calls, pageRequests, pageCursors, followCursors } } +function remoteItem( + generation: number, + value: RemoteJournalFrame, + signal: AbortSignal, +): RemoteStreamItem> { + return { generation, value, signal, accept: vi.fn() } +} + +function controlledFactory( + next: () => Promise>>>, +): RemoteStreamFactory { + const lifetime = new AbortController() + return { + $stream(): RemoteStream { + 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 + }, + } +} + 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() + 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>() + 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((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>() + 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('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>() + const liveConsumed = Promise.withResolvers() + const openingPage = Promise.withResolvers() + const finish = Promise.withResolvers() + 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('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>() + const liveConsumed = Promise.withResolvers() + const openingPage = Promise.withResolvers() + 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>() + const secondLive = Promise.withResolvers>() + const firstConsumed = Promise.withResolvers() + const secondConsumed = Promise.withResolvers() + const openingPage = Promise.withResolvers() + const repairPage = Promise.withResolvers() + 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() + 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>() + 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(() => {}) + 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>() + const duplicatePage = new Promise(() => {}) + 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>>>({ + done: false, + value: remoteItem(1, { type: 'opened', cursor: 0 }, generation.signal), + }), + Promise.resolve>>>({ + done: true, + value: undefined, + }), + ] + const fixture = journalFixture( + [], + [new Promise(() => {})], + 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>>>() + const results = [ + Promise.resolve>>>({ + 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>>>() + const nextPending = Promise.withResolvers>>>() + const results = [ + Promise.resolve>>>({ + done: false, + value: remoteItem(1, { type: 'opened', cursor: 0 }, generation.signal), + }), + pending.promise, + nextPending.promise, + ] + const fixture = journalFixture( + [], + [signal => new Promise((_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 } }] }], diff --git a/packages/api/session-controller/tests/control-queue.host.spec.ts b/packages/api/session-controller/tests/control-queue.host.spec.ts index 38902726dd..9ce0c453a0 100644 --- a/packages/api/session-controller/tests/control-queue.host.spec.ts +++ b/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 }) + }) }) diff --git a/packages/api/session-controller/tests/transport.client.spec.ts b/packages/api/session-controller/tests/transport.client.spec.ts index 493b41d40d..2dd1c2602d 100644 --- a/packages/api/session-controller/tests/transport.client.spec.ts +++ b/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 } class ScriptedSessionRemote implements SessionTransportRemote { @@ -68,6 +69,7 @@ class ScriptedSessionRemote implements SessionTransportRemote { private readonly generations: FollowGeneration[], private readonly pages: RemoteResult[], private readonly controlFrames: readonly SessionControlFrame[] = [], + private readonly holdControl = true, ) {} async *follow(request: SessionFollowRequest, signal = new AbortController().signal): AsyncIterable { @@ -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((resolve) => { @@ -93,7 +96,7 @@ class ScriptedSessionRemote implements SessionTransportRemote { async *control(signal = new AbortController().signal): AsyncIterable { for (const frame of this.controlFrames) yield frame - if (!signal.aborted) { + if (this.holdControl && !signal.aborted) { await new Promise((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() + 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() + }) }) diff --git a/packages/api/workspace-controller/tests/transport.client.spec.ts b/packages/api/workspace-controller/tests/transport.client.spec.ts index a38859daca..dd4517600c 100644 --- a/packages/api/workspace-controller/tests/transport.client.spec.ts +++ b/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() + 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 | 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',