mirror of
https://github.com/deepseek-ai/deepseek-harness.git
synced 2026-08-29 04:26:38 +00:00
fix(api-gateway): retry journal pages after reconnect
This commit is contained in:
@@ -396,7 +396,10 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
|
||||
signal.throwIfAborted()
|
||||
return { type: 'page', page: result.value }
|
||||
}
|
||||
if (result.type === 'page-error') throw result.error
|
||||
if (result.type === 'page-error') {
|
||||
if (!signal.aborted || this.stream.signal.aborted) throw result.error
|
||||
return this.awaitReplacementGeneration(generation, iterator, pending)
|
||||
}
|
||||
this.releaseNext(pending)
|
||||
if (result.type === 'next-error') throw result.error
|
||||
if (result.value.done) {
|
||||
@@ -412,6 +415,32 @@ export abstract class RemoteJournalStream<Page, Entry, Cursor, PageRequest = voi
|
||||
}
|
||||
}
|
||||
|
||||
private async awaitReplacementGeneration(
|
||||
generation: number,
|
||||
iterator: AsyncIterator<JournalStreamItem<Entry, Cursor>>,
|
||||
initial: Promise<IteratorResult<JournalStreamItem<Entry, Cursor>>>,
|
||||
): Promise<{ readonly type: 'superseded'; readonly item: JournalStreamItem<Entry, Cursor> }> {
|
||||
let pending = initial
|
||||
while (true) {
|
||||
let next: IteratorResult<JournalStreamItem<Entry, Cursor>>
|
||||
try {
|
||||
next = await pending
|
||||
} finally {
|
||||
this.releaseNext(pending)
|
||||
}
|
||||
if (next.done) {
|
||||
this.stream.signal.throwIfAborted()
|
||||
throw new Error(`${this.options.name} ended while replacing an aborted page generation`)
|
||||
}
|
||||
const item = next.value
|
||||
if (item.generation !== generation) return { type: 'superseded', item }
|
||||
if (item.value.type === 'opened') {
|
||||
throw new Error(`${this.options.name} emitted more than one opening cursor`)
|
||||
}
|
||||
pending = this.nextResult(iterator)
|
||||
}
|
||||
}
|
||||
|
||||
private mergeReplacement(page: Page, queued: readonly Entry[]): Entry[] | undefined {
|
||||
const entries = [...this.options.entries(page)]
|
||||
this.assertPage(entries)
|
||||
|
||||
@@ -29,6 +29,8 @@ interface Generation {
|
||||
readonly hold?: boolean
|
||||
}
|
||||
|
||||
type PageSource = Page | Promise<Page> | ((signal: AbortSignal) => Promise<Page>)
|
||||
|
||||
const AVAILABLE_CONNECTION = {
|
||||
hostDescription: {
|
||||
getSnapshot: () => ({
|
||||
@@ -55,7 +57,7 @@ const STREAM_FACTORY = {
|
||||
class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageRequest> {
|
||||
constructor(
|
||||
private readonly generations: Generation[],
|
||||
private readonly pages: (Page | Promise<Page>)[],
|
||||
private readonly pages: PageSource[],
|
||||
private readonly calls: string[],
|
||||
private readonly pageRequests: PageRequest[],
|
||||
private readonly pageCursors: number[],
|
||||
@@ -95,13 +97,17 @@ class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageReques
|
||||
}
|
||||
|
||||
/** @inheritdoc */
|
||||
protected override readPage(request: PageRequest, through: number): Promise<Page> {
|
||||
protected override readPage(
|
||||
request: PageRequest,
|
||||
through: number,
|
||||
signal: AbortSignal,
|
||||
): Promise<Page> {
|
||||
this.calls.push('page')
|
||||
this.pageRequests.push(request)
|
||||
this.pageCursors.push(through)
|
||||
const value = this.pages.shift()
|
||||
if (value === undefined) throw new Error('no scripted journal page')
|
||||
return Promise.resolve(value)
|
||||
return typeof value === 'function' ? value(signal) : Promise.resolve(value)
|
||||
}
|
||||
|
||||
/** @inheritdoc */
|
||||
@@ -112,7 +118,7 @@ class FixtureJournal extends RemoteJournalStream<Page, Entry, number, PageReques
|
||||
|
||||
function journalFixture(
|
||||
generations: Generation[],
|
||||
pages: (Page | Promise<Page>)[],
|
||||
pages: PageSource[],
|
||||
): {
|
||||
readonly journal: RemoteJournalStream<Page, Entry, number, PageRequest>
|
||||
readonly changes: RemoteJournalChange<Page, Entry>[]
|
||||
@@ -241,6 +247,42 @@ describe('RemoteJournalStream', () => {
|
||||
await fixture.journal.dispose()
|
||||
})
|
||||
|
||||
it('restarts a page aborted with its carrier generation', async () => {
|
||||
const fixture = journalFixture(
|
||||
[
|
||||
{
|
||||
frames: [{ type: 'opened', cursor: 1 }],
|
||||
terminal: new RemoteStreamCarrierError('carrier lost during page'),
|
||||
},
|
||||
{
|
||||
frames: [{ type: 'opened', cursor: 2 }],
|
||||
hold: true,
|
||||
},
|
||||
],
|
||||
[
|
||||
signal => new Promise<Page>((_resolve, reject) => {
|
||||
const aborted = (): void => { reject(signal.reason) }
|
||||
signal.addEventListener('abort', aborted, { once: true })
|
||||
if (signal.aborted) aborted()
|
||||
}),
|
||||
page('replacement', [0, 1, 2]),
|
||||
],
|
||||
)
|
||||
|
||||
await fixture.journal.open({ limit: 3 })
|
||||
|
||||
expect(fixture.changes).toEqual([{
|
||||
type: 'replace',
|
||||
page: page('replacement', [0, 1, 2]),
|
||||
entries: entries(0, 1, 2),
|
||||
hasMore: false,
|
||||
}])
|
||||
expect(fixture.pageCursors).toEqual([1, 2])
|
||||
expect(fixture.followCursors).toEqual([undefined, 1])
|
||||
expect(fixture.failed).not.toHaveBeenCalled()
|
||||
await fixture.journal.dispose()
|
||||
})
|
||||
|
||||
it('repairs a live gap before publishing another change', async () => {
|
||||
const fixture = journalFixture(
|
||||
[{
|
||||
|
||||
Reference in New Issue
Block a user