From b2128a7d189ac020ebb6e49a57ee986e98326b77 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Wed, 12 Aug 2026 07:33:07 -0400 Subject: [PATCH] =?UTF-8?q?=F0=9F=93=A1=20fix:=20Preserve=20Redis=20Abort?= =?UTF-8?q?=20Terminal=20Delivery=20(#14749)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(stream): preserve terminal delivery after Redis fences * fix(stream): cover Redis abort acknowledgment window * chore: sort stream timing imports * fix(stream): follow durable replacement handoffs * fix(stream): bound replacement handoff retirement * test(stream): cover handoff deadline drain * test(stream): cover fenced steer retirement grace * fix(stream): scope fenced retirement lifecycle * fix(stream): retire fenced subscribers safely * test(stream): settle subscription fixtures --- .../api/src/stream/GenerationJobManager.ts | 369 ++++++- .../__tests__/RedisEventTransport.spec.ts | 38 + .../stream/__tests__/idempotencyClaim.spec.ts | 93 +- .../api/src/stream/__tests__/startup.spec.ts | 954 ++++++++++++++++++ .../__tests__/steerReceiptIntegrity.spec.ts | 8 + .../implementations/RedisEventTransport.ts | 31 +- .../api/src/stream/interfaces/IJobStore.ts | 4 + packages/api/src/stream/internal/timing.ts | 16 + 8 files changed, 1448 insertions(+), 65 deletions(-) create mode 100644 packages/api/src/stream/internal/timing.ts diff --git a/packages/api/src/stream/GenerationJobManager.ts b/packages/api/src/stream/GenerationJobManager.ts index 166273f79e..9f99cf0a82 100644 --- a/packages/api/src/stream/GenerationJobManager.ts +++ b/packages/api/src/stream/GenerationJobManager.ts @@ -69,6 +69,11 @@ import { InMemoryJobStore } from './implementations/InMemoryJobStore'; import { attachAskUserQuestionAnswers, normalizeResumeRunStepIndices } from '~/agents/hitl/resume'; import { emitChunkWithReceipt } from './internal/chunkPublication'; import { resolveCoalesceWindowMs } from './internal/coalescing'; +import { + REDIS_ABORT_TERMINAL_GRACE_MS, + REDIS_EVENT_REORDER_TIMEOUT_MS, + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS, +} from './internal/timing'; import { filterPersistableAbortContent } from './abortContent'; import { toClientPendingAction } from '~/agents/hitl/policy'; import { ApprovalLifecycle, pausePersistenceActionId } from './ApprovalLifecycle'; @@ -586,6 +591,12 @@ interface RuntimeJobState { allSubscribersLeftHandlers?: Array<(...args: unknown[]) => void | Promise>; } +interface FencedRuntimeRetirementContext { + controller: AbortController; + timer?: NodeJS.Timeout; + cleanupStarted?: boolean; +} + interface PreparedSubscription { runtime: RuntimeJobState; jobData: SerializableJobData | null; @@ -656,6 +667,15 @@ class GenerationJobManagerClass { private cleanupInterval: NodeJS.Timeout | null = null; + /** Generation-scoped retirement callbacks must not outlive the configured + * store/transport pair that created them. */ + private fencedRuntimeRetirements = new Map(); + + /** Suppresses the stream-global disconnect callback while an exact fenced + * predecessor detaches. The callback may otherwise target a same-stream + * successor that never owned the departing SSE response. */ + private fencedSubscriberDetachments = new Set(); + /** Rejects new jobs once graceful shutdown has started. */ private shuttingDown = false; @@ -723,6 +743,7 @@ class GenerationJobManagerClass { cleanupOnComplete?: boolean; }): void { assertJobStoreV2(services.jobStore); + this.cancelFencedRuntimeRetirements(); const previousStore = this.storeLabel; if (this.cleanupInterval) { logger.warn( @@ -1281,6 +1302,9 @@ class GenerationJobManagerClass { private registerAllSubscribersLeft(streamId: string): void { this.eventTransport.onAllSubscribersLeft(streamId, () => { + if (this.fencedSubscriberDetachments.has(streamId)) { + return; + } const runtime = this.runtimeState.get(streamId); if (!runtime) { return; @@ -4079,8 +4103,11 @@ class GenerationJobManagerClass { streamId, { onChunk: (event, generationId) => { + const currentRuntime = this.runtimeState.get(streamId); + const isMatchingTerminalDrain = + currentRuntime == null && generationId === runtime.createdAt; if ( - this.runtimeState.get(streamId) !== runtime || + (currentRuntime !== runtime && !isMatchingTerminalDrain) || (generationId != null && generationId !== runtime.createdAt) ) { return; @@ -4134,7 +4161,29 @@ class GenerationJobManagerClass { runtime.earlyReplayHandlers.delete(queueChunk); runtime.localErrorHandlers.delete(queueError); detachSignal?.removeEventListener('abort', detachOnAbort); - transportSubscription.unsubscribe(); + const fencedLastSubscriber = + runtime.localErrorHandlers.size === 0 && + (this.fencedRuntimeRetirements.has(runtime) || + this.runtimeState.get(streamId) !== runtime); + const addFencedDetachmentSuppression = + fencedLastSubscriber && !this.fencedSubscriberDetachments.has(streamId); + if (fencedLastSubscriber) { + runtime.syncSent = false; + runtime.hasSubscriber = false; + runtime.attachmentGeneration++; + runtime.lastSubscriberCleanupGeneration = runtime.attachmentGeneration; + } + if (addFencedDetachmentSuppression) { + this.fencedSubscriberDetachments.add(streamId); + } + try { + this.cleanupUnobservedFencedRuntime(streamId, runtime); + transportSubscription.unsubscribe(); + } finally { + if (addFencedDetachmentSuppression) { + this.fencedSubscriberDetachments.delete(streamId); + } + } resolveDetached(); }, }; @@ -5444,14 +5493,290 @@ class GenerationJobManagerClass { } } + private async getReplacementHandoffState( + streamId: string, + predecessorCreatedAt: number, + ): Promise<'none' | 'pending' | 'settled'> { + try { + const current = (await this.jobStore.getJob(streamId)) as CreatedJobData | null; + if (current == null) { + /** A fenced append proves that a newer durable owner existed. Its job + * may disappear after publishing predecessor DONE but before this poll, + * so absence is a settled handoff that still needs delivery grace. */ + return 'settled'; + } + if (current.createdAt <= predecessorCreatedAt) { + return 'none'; + } + const receipts = + current.replacedJobs ?? (current.replacedJob != null ? [current.replacedJob] : []); + return receipts.some((receipt) => receipt.createdAt === predecessorCreatedAt) + ? 'pending' + : 'settled'; + } catch (error) { + logger.warn( + `[GenerationJobManager] Failed to inspect replacement handoff for fenced generation ${streamId}:`, + error, + ); + // A transient read failure cannot prove that DONE publication finished. + return 'pending'; + } + } + + private async getReplacementHandoffStateBeforeDeadline( + streamId: string, + predecessorCreatedAt: number, + deadlineAt: number, + lifecycleSignal: AbortSignal, + ): Promise<'none' | 'pending' | 'settled' | 'deadline' | 'cancelled'> { + if (lifecycleSignal.aborted) { + return 'cancelled'; + } + const remainingMs = deadlineAt - Date.now(); + if (remainingMs <= 0) { + return 'deadline'; + } + + let deadlineTimer: ReturnType | undefined; + let removeCancellationListener: (() => void) | undefined; + try { + const cancellation = new Promise<'cancelled'>((resolve) => { + const onCancelled = (): void => resolve('cancelled'); + removeCancellationListener = (): void => + lifecycleSignal.removeEventListener('abort', onCancelled); + lifecycleSignal.addEventListener('abort', onCancelled, { once: true }); + if (lifecycleSignal.aborted) { + onCancelled(); + } + }); + const deadline = new Promise<'deadline'>((resolve) => { + deadlineTimer = setTimeout(() => resolve('deadline'), remainingMs); + deadlineTimer.unref?.(); + }); + return await Promise.race([ + this.getReplacementHandoffState(streamId, predecessorCreatedAt), + deadline, + cancellation, + ]); + } finally { + if (deadlineTimer != null) { + clearTimeout(deadlineTimer); + } + removeCancellationListener?.(); + } + } + + private cancelFencedRuntimeRetirements(): void { + for (const retirement of this.fencedRuntimeRetirements.values()) { + if (retirement.timer != null) { + clearTimeout(retirement.timer); + retirement.timer = undefined; + } + retirement.controller.abort(); + } + this.fencedRuntimeRetirements.clear(); + } + + private scheduleFencedRuntimeRetirement( + streamId: string, + runtime: RuntimeJobState, + startedAt: number, + retirement: FencedRuntimeRetirementContext, + handoffPendingObserved = false, + postHandoffGraceApplied = false, + requestedDelayMs = REDIS_ABORT_TERMINAL_GRACE_MS, + ): void { + const lifecycleSignal = retirement.controller.signal; + if (lifecycleSignal.aborted || this.fencedRuntimeRetirements.get(runtime) !== retirement) { + return; + } + const remainingMs = REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS - (Date.now() - startedAt); + /** Handoff inspection is capped at 30 seconds. Once inspection ends, the + * final delivery grace is intentionally un-clamped and may extend the + * generation-scoped hold by one reorder window. */ + const delayMs = postHandoffGraceApplied + ? Math.max(1, requestedDelayMs) + : Math.max(1, Math.min(requestedDelayMs, remainingMs)); + const retirementTimer = setTimeout(() => { + if (retirement.timer === retirementTimer) { + retirement.timer = undefined; + } + if (lifecycleSignal.aborted) { + return; + } + void this.finishFencedRuntimeRetirement( + streamId, + runtime, + startedAt, + retirement, + handoffPendingObserved, + postHandoffGraceApplied, + lifecycleSignal, + ); + }, delayMs); + retirement.timer = retirementTimer; + retirementTimer.unref?.(); + } + + private async finishFencedRuntimeRetirement( + streamId: string, + runtime: RuntimeJobState, + startedAt: number, + retirement: FencedRuntimeRetirementContext, + handoffPendingObserved: boolean, + postHandoffGraceApplied: boolean, + lifecycleSignal: AbortSignal, + ): Promise { + if ( + this.shuttingDown || + lifecycleSignal.aborted || + this.fencedRuntimeRetirements.get(runtime) !== retirement + ) { + return; + } + if ( + runtime.localErrorHandlers.size === 0 && + this.cleanupUnobservedFencedRuntime(streamId, runtime) + ) { + return; + } + if (!postHandoffGraceApplied) { + const deadlineAt = startedAt + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS; + const handoffState = await this.getReplacementHandoffStateBeforeDeadline( + streamId, + runtime.createdAt, + deadlineAt, + lifecycleSignal, + ); + if (this.shuttingDown || lifecycleSignal.aborted || handoffState === 'cancelled') { + return; + } + if (handoffState === 'pending' && Date.now() < deadlineAt) { + this.scheduleFencedRuntimeRetirement(streamId, runtime, startedAt, retirement, true); + return; + } + if ( + handoffPendingObserved || + handoffState === 'pending' || + handoffState === 'settled' || + handoffState === 'deadline' + ) { + this.scheduleFencedRuntimeRetirement( + streamId, + runtime, + startedAt, + retirement, + handoffPendingObserved || handoffState === 'pending', + true, + REDIS_EVENT_REORDER_TIMEOUT_MS * 2, + ); + return; + } + } + + if (this.fencedRuntimeRetirements.get(runtime) !== retirement) { + return; + } + retirement.cleanupStarted = true; + this.cleanupFencedRuntime(streamId, runtime); + if (this.fencedRuntimeRetirements.get(runtime) === retirement) { + this.fencedRuntimeRetirements.delete(runtime); + } + } + + private cleanupUnobservedFencedRuntime(streamId: string, runtime: RuntimeJobState): boolean { + if (runtime.localErrorHandlers.size !== 0) { + return false; + } + const retirement = this.fencedRuntimeRetirements.get(runtime); + if (retirement == null || retirement.cleanupStarted === true) { + return false; + } + retirement.cleanupStarted = true; + if (retirement.timer != null) { + clearTimeout(retirement.timer); + retirement.timer = undefined; + } + retirement.controller.abort(); + this.fencedRuntimeRetirements.delete(runtime); + this.cleanupFencedRuntime(streamId, runtime); + return true; + } + + private recordFencedRuntimeAbortProof( + streamId: string, + runtime: RuntimeJobState, + ownsExactProvider: boolean, + ): void { + const recordAbortAcknowledgement = this.eventTransport.recordAbortAcknowledgement; + if (!this._isRedis || recordAbortAcknowledgement == null || !ownsExactProvider) { + return; + } + + try { + void recordAbortAcknowledgement + .call(this.eventTransport, streamId, runtime.createdAt) + .then((confirmed) => { + if (!confirmed) { + logger.warn( + `[GenerationJobManager] Abort proof was not persisted for fenced generation ${streamId}`, + ); + } + }) + .catch((error) => { + logger.error( + `[GenerationJobManager] Failed to persist abort proof for fenced generation ${streamId}:`, + error, + ); + }); + } catch (error) { + logger.error( + `[GenerationJobManager] Failed to start abort proof for fenced generation ${streamId}:`, + error, + ); + } + } + + private cleanupFencedRuntime(streamId: string, runtime: RuntimeJobState): void { + if (runtime.replacementTransportHold !== true) { + this.releaseAbortSubscription(runtime); + this.releaseJobOwnership(streamId, runtime.createdAt); + this.jobStore.clearContentState(streamId, runtime.createdAt); + + if (this.runtimeState.get(streamId) === runtime) { + this.runtimeState.delete(streamId); + this.runStepBuffers?.delete(streamId); + this.replayEventWriteQueues.delete(streamId); + this.tokenUsageWriteQueues.delete(streamId); + } + } + + // finalEvent/errorEvent are cached before transport dispatch, so they do + // not prove that an attached SSE response closed. Each captured handler + // ignores this reconnect signal when its terminal is already queued. + for (const notify of [...runtime.localErrorHandlers]) { + try { + notify(TERMINAL_PUBLICATION_RECONNECT_ERROR); + } catch (error) { + logger.error( + `[GenerationJobManager] Failed to recycle a subscriber for fenced generation ${streamId}:`, + error, + ); + } + } + } + /** A generation-fenced append returning false is same-slot proof that this - * epoch is no longer the durable owner. This is the provider's backstop when - * a replacement abort publication is lost. Retire only the exact captured - * runtime; a newer local successor may already occupy the same stream id. */ + * epoch is no longer the durable owner. Stop its provider immediately, then + * preserve captured subscribers while a durable successor still owns their + * handoff receipt. A newer runtime is never touched. */ private retireRuntimeAfterDurableFence(streamId: string, runtime: RuntimeJobState): void { if (this.runtimeState.get(streamId) !== runtime) { return; } + if (this.fencedRuntimeRetirements.has(runtime)) { + return; + } /** A runtime whose stop signal already landed (cross-replica abort, * replacement handshake) observes this fence as a consequence of its own * termination — most often a coalesced window draining after the abort @@ -5464,30 +5789,22 @@ class GenerationJobManagerClass { } runtime.startupTelemetry?.end('replaced'); runtime.startupTelemetry = undefined; + const ownsExactProvider = this.ownedJobs.get(streamId) === runtime.createdAt; runtime.abortController.abort(); + this.recordFencedRuntimeAbortProof(streamId, runtime, ownsExactProvider); + if (this.shuttingDown) { + return; + } if (runtime.replacementTransportHold === true) { return; } - this.releaseAbortSubscription(runtime); - if (this.runtimeState.get(streamId) === runtime) { - this.runtimeState.delete(streamId); - this.releaseJobOwnership(streamId, runtime.createdAt); - this.runStepBuffers?.delete(streamId); - this.replayEventWriteQueues.delete(streamId); - this.tokenUsageWriteQueues.delete(streamId); - this.jobStore.clearContentState(streamId, runtime.createdAt); - try { - // Delete runtime ownership before closing. The transport's - // all-subscribers-left callback may run synchronously; it must not - // persist a partial response for this already-replaced epoch. - this.eventTransport.closeLocalSubscribers?.(streamId, TERMINAL_PUBLICATION_RECONNECT_ERROR); - } catch (error) { - logger.error( - `[GenerationJobManager] Failed to recycle subscribers for fenced generation ${streamId}:`, - error, - ); - } + if (runtime.localErrorHandlers.size === 0) { + this.cleanupFencedRuntime(streamId, runtime); + return; } + const retirement: FencedRuntimeRetirementContext = { controller: new AbortController() }; + this.fencedRuntimeRetirements.set(runtime, retirement); + this.scheduleFencedRuntimeRetirement(streamId, runtime, Date.now(), retirement); } private isCurrentRuntime(streamId: string, runtime: RuntimeJobState): boolean { @@ -6736,6 +7053,7 @@ class GenerationJobManagerClass { /** Returns sizes of internal runtime maps for diagnostics */ getRuntimeStats(): { runtimeStateSize: number; + fencedRuntimeRetirements: number; runStepBufferSize: number; eventTransportStreams: number; earlyBufferedEvents: number; @@ -6749,6 +7067,7 @@ class GenerationJobManagerClass { } return { runtimeStateSize: this.runtimeState.size, + fencedRuntimeRetirements: this.fencedRuntimeRetirements.size, runStepBufferSize: this.runStepBuffers?.size ?? 0, eventTransportStreams: this.eventTransport.getTrackedStreamIds().length, earlyBufferedEvents, @@ -6859,6 +7178,7 @@ class GenerationJobManagerClass { return; } this.shuttingDown = true; + this.cancelFencedRuntimeRetirements(); for (const runtime of this.runtimeState.values()) { runtime.startupTelemetry?.end('aborted'); @@ -6880,6 +7200,7 @@ class GenerationJobManagerClass { */ async destroy(): Promise { this.shuttingDown = true; + this.cancelFencedRuntimeRetirements(); if (this.cleanupInterval) { clearInterval(this.cleanupInterval); diff --git a/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts b/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts index ef9a26b0f2..cd37e16f34 100644 --- a/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts +++ b/packages/api/src/stream/__tests__/RedisEventTransport.spec.ts @@ -1109,6 +1109,44 @@ describe('RedisEventTransport', () => { } }); + it('recovers a confirmed abort from durable owner proof after the initial read', async () => { + jest.useFakeTimers(); + const mockPublisher = createMockPublisher(); + const mockSubscriber = createMockSubscriber(); + const transport = new RedisEventTransport( + mockPublisher as unknown as Redis, + mockSubscriber as unknown as Redis, + ); + const streamId = 'delayed-owner-abort-proof'; + let resolveAbortPublished!: () => void; + const abortPublished = new Promise((resolve) => { + resolveAbortPublished = resolve; + }); + mockPublisher.publish.mockImplementationOnce(async () => { + resolveAbortPublished(); + return 0; + }); + + try { + const confirmation = transport.emitAbortConfirmed(streamId, 1234); + await abortPublished; + + await expect(transport.recordAbortAcknowledgement(streamId, 1234)).resolves.toBe(true); + expect(mockPublisher.set).toHaveBeenCalledWith( + `stream:{${streamId}}:abort-ack:1234`, + '1', + 'EX', + 86400, + ); + + await jest.advanceTimersByTimeAsync(3000); + await expect(confirmation).resolves.toBe(true); + } finally { + transport.destroy(); + jest.useRealTimers(); + } + }); + it('waits for a delayed owner acknowledgement when cluster publish reports zero local receivers', async () => { const mockPublisher = createMockPublisher(); const mockSubscriber = createMockSubscriber(); diff --git a/packages/api/src/stream/__tests__/idempotencyClaim.spec.ts b/packages/api/src/stream/__tests__/idempotencyClaim.spec.ts index d8b4c407b8..e9ed727c78 100644 --- a/packages/api/src/stream/__tests__/idempotencyClaim.spec.ts +++ b/packages/api/src/stream/__tests__/idempotencyClaim.spec.ts @@ -1,3 +1,8 @@ +import { + REDIS_ABORT_TERMINAL_GRACE_MS, + REDIS_EVENT_REORDER_TIMEOUT_MS, + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS, +} from '~/stream/internal/timing'; import { JobCreationSupersededError, type CreatedJobData, @@ -760,19 +765,33 @@ describe('GenerationJobManager start-generation claim', () => { 'remote-replacement-attempt', ); - await manager.emitChunk(streamId, { - event: 'on_message_delta', - data: { delta: 'stale provider output' }, - }); - await Promise.resolve(); + jest.useFakeTimers(); + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { delta: 'stale provider output' }, + }); + await Promise.resolve(); - expect(predecessor.abortController.signal.aborted).toBe(true); - expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); - expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); - expect(await store.getJob(streamId)).toMatchObject({ - createdAt: replacement.createdAt, - status: 'running', - }); + expect(predecessor.abortController.signal.aborted).toBe(true); + expect(onError).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync(REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS); + + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + + expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + expect(await store.getJob(streamId)).toMatchObject({ + createdAt: replacement.createdAt, + status: 'running', + }); + } finally { + jest.useRealTimers(); + } }); it('self-fences an ordinary provider when its Redis append rejects', async () => { @@ -794,16 +813,25 @@ describe('GenerationJobManager start-generation claim', () => { expect(await manager.subscribe(streamId, jest.fn(), jest.fn(), onError)).not.toBeNull(); jest.spyOn(store, 'appendChunk').mockRejectedValueOnce(new Error('simulated Redis outage')); - await manager.emitChunk(streamId, { - event: 'on_message_delta', - data: { delta: 'uncoordinated provider output' }, - }); - await Promise.resolve(); - await Promise.resolve(); + jest.useFakeTimers(); + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { delta: 'uncoordinated provider output' }, + }); + await Promise.resolve(); + await Promise.resolve(); - expect(job.abortController.signal.aborted).toBe(true); - expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); - expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + expect(job.abortController.signal.aborted).toBe(true); + expect(onError).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + + expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + } }); it('self-fences when a chunk append wins but its active-only publication is fenced', async () => { @@ -822,14 +850,23 @@ describe('GenerationJobManager start-generation claim', () => { expect(await manager.subscribe(streamId, jest.fn(), jest.fn(), onError)).not.toBeNull(); registerChunkPublicationCapability(transport, async () => false); - await manager.emitChunk(streamId, { - event: 'on_message_delta', - data: { delta: 'publication lost its active-generation CAS' }, - }); + jest.useFakeTimers(); + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { delta: 'publication lost its active-generation CAS' }, + }); - expect(job.abortController.signal.aborted).toBe(true); - expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); - expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + expect(job.abortController.signal.aborted).toBe(true); + expect(onError).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + + expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + } }); it('does not let a late stale-append result retire a newer local runtime', async () => { diff --git a/packages/api/src/stream/__tests__/startup.spec.ts b/packages/api/src/stream/__tests__/startup.spec.ts index a45e977331..1ebf2c7064 100644 --- a/packages/api/src/stream/__tests__/startup.spec.ts +++ b/packages/api/src/stream/__tests__/startup.spec.ts @@ -1,6 +1,12 @@ import type { AbortResult } from '~/stream/interfaces/IJobStore'; import type { AgentStartupTelemetry } from '~/agents/startup'; import type { ServerSentEvent } from '~/types'; +import { + REDIS_ABORT_ACK_TIMEOUT_MS, + REDIS_ABORT_TERMINAL_GRACE_MS, + REDIS_EVENT_REORDER_TIMEOUT_MS, + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS, +} from '~/stream/internal/timing'; import { GenerationJobManagerClass, TERMINAL_PUBLICATION_RECONNECT_ERROR, @@ -30,6 +36,33 @@ function createManager(): GenerationJobManagerClass { return manager; } +class DeferredEventTransport extends InMemoryEventTransport { + private pendingDeliveries: Array<() => void> = []; + + override emitChunk(streamId: string, event: unknown, generationId?: number): void { + this.pendingDeliveries.push(() => super.emitChunk(streamId, event, generationId)); + } + + override emitDone(streamId: string, event: unknown, generationId?: number): void { + this.pendingDeliveries.push(() => super.emitDone(streamId, event, generationId)); + } + + flushDeliveries(): void { + const deliveries = this.pendingDeliveries; + this.pendingDeliveries = []; + for (const deliver of deliveries) { + deliver(); + } + } +} + +function getEventLabel(event: ServerSentEvent): string | undefined { + if (!('data' in event) || typeof event.data !== 'object' || event.data == null) { + return undefined; + } + return typeof event.data.label === 'string' ? event.data.label : undefined; +} + function createPendingAction(streamId: string) { return buildPendingAction( buildToolApprovalPayload([ @@ -641,6 +674,927 @@ describe('GenerationJobManager startup telemetry', () => { await manager.destroy(); }); + it('delivers only matching generation-tagged chunks after terminal cleanup before done', async () => { + const streamId = 'stream-terminal-delayed-chunk'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new DeferredEventTransport(); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const order: string[] = []; + const subscription = await manager.subscribe( + streamId, + (event) => { + const label = getEventLabel(event); + if (label != null) { + order.push(label); + } + }, + () => order.push('done'), + ); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { label: 'matching' }, + } as ServerSentEvent); + eventTransport.emitChunk(streamId, { data: { label: 'untagged' } }); + eventTransport.emitChunk(streamId, { data: { label: 'mismatched' } }, job.createdAt + 1); + + await expect(manager.abortJob(streamId)).resolves.toMatchObject({ success: true }); + await expect(jobStore.getJob(streamId)).resolves.toBeNull(); + expect(order).toEqual([]); + + eventTransport.flushDeliveries(); + + expect(order).toEqual(['matching', 'done']); + } finally { + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('rejects a delayed predecessor chunk after a replacement runtime is installed', async () => { + const now = jest.spyOn(Date, 'now').mockReturnValue(1000); + const streamId = 'stream-delayed-predecessor-chunk'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new DeferredEventTransport(); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + manager.initialize(); + const predecessor = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const received: string[] = []; + const subscription = await manager.subscribe(streamId, (event) => { + const label = getEventLabel(event); + if (label != null) { + received.push(label); + } + }); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { label: 'predecessor' }, + } as ServerSentEvent); + now.mockReturnValue(2000); + const replacement = await manager.createJob(streamId, 'user-1', 'conversation-1'); + + expect(replacement.createdAt).toBeGreaterThan(predecessor.createdAt); + eventTransport.flushDeliveries(); + + expect(received).toEqual([]); + } finally { + now.mockRestore(); + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('lets a matching terminal drain win after a pre-signal durable fence', async () => { + const streamId = 'stream-fenced-terminal-drain'; + const eventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ + jobStore: new InMemoryJobStore({ ttlAfterComplete: 60_000 }), + eventTransport, + isRedis: false, + }); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onDone = jest.fn(); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, onDone, onError); + jest.useFakeTimers(); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + + expect(job.abortController.signal.aborted).toBe(true); + expect(onError).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync( + REDIS_ABORT_ACK_TIMEOUT_MS + REDIS_EVENT_REORDER_TIMEOUT_MS, + ); + expect(onError).not.toHaveBeenCalled(); + + eventTransport.emitDone(streamId, { final: true }, job.createdAt); + expect(onDone).toHaveBeenCalledWith({ final: true }); + + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS); + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('keeps a fenced predecessor attached while its durable handoff receipt is pending', async () => { + const streamId = 'stream-fenced-handoff-pending'; + const clientRequestId = 'req-fenced-handoff-pending'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const owner = new GenerationJobManagerClass(); + const replacer = new GenerationJobManagerClass(); + owner.configure({ jobStore, eventTransport, isRedis: false }); + replacer.configure({ jobStore, eventTransport, isRedis: false }); + owner.initialize(); + replacer.initialize(); + const predecessor = await owner.createJob(streamId, 'user-1', 'conversation-1'); + const predecessorDone = jest.fn(); + const predecessorError = jest.fn(); + const predecessorSubscription = await owner.subscribe( + streamId, + () => undefined, + predecessorDone, + predecessorError, + ); + const claim = await replacer.claimGeneration( + 'user-1', + clientRequestId, + streamId, + 'conversation-1', + 2, + ); + const actualMarkStarted = jobStore.markIdempotencyKeyStarted.bind(jobStore); + let signalLegacyMarkStarted: (() => void) | undefined; + const legacyMarkStarted = new Promise((resolve) => { + signalLegacyMarkStarted = resolve; + }); + let releaseLegacyMark: (() => void) | undefined; + const legacyMarkGate = new Promise((resolve) => { + releaseLegacyMark = resolve; + }); + const markStartedSpy = jest + .spyOn(jobStore, 'markIdempotencyKeyStarted') + .mockImplementationOnce(async (...args) => { + signalLegacyMarkStarted?.(); + await legacyMarkGate; + return actualMarkStarted(...args); + }); + let creatingReplacement: ReturnType | undefined; + jest.useFakeTimers(); + + try { + creatingReplacement = replacer.createJob(streamId, 'user-1', 'conversation-1', { + idempotencyClientRequestId: clientRequestId, + idempotencyClaimToken: claim.existing!.claimToken, + initialMetadata: { generationProtocolVersion: 2 }, + }); + await legacyMarkStarted; + + await owner.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + expect(predecessor.abortController.signal.aborted).toBe(true); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + + expect(predecessorDone).not.toHaveBeenCalled(); + expect(predecessorError).not.toHaveBeenCalled(); + expect(eventTransport.getSubscriberCount(streamId)).toBe(1); + + releaseLegacyMark?.(); + const replacement = await creatingReplacement; + + expect(replacement.createdAt).toBeGreaterThan(predecessor.createdAt); + expect(predecessorDone).toHaveBeenCalledWith( + expect.objectContaining({ + final: true, + reconcile: true, + reconcileReason: 'generation_replaced', + generationCreatedAt: predecessor.createdAt, + }), + ); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + expect(predecessorError).not.toHaveBeenCalled(); + expect(owner.getRuntimeStats()).toMatchObject({ + runtimeStateSize: 0, + fencedRuntimeRetirements: 0, + }); + + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + expect(predecessorError).not.toHaveBeenCalled(); + expect(owner.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + releaseLegacyMark?.(); + await creatingReplacement?.catch(() => undefined); + markStartedSpy.mockRestore(); + predecessorSubscription?.unsubscribe(); + await Promise.all([owner.destroy(), replacer.destroy()]); + } + }); + + it('bounds a stalled replacement lookup before recycling the fenced subscriber', async () => { + const streamId = 'stream-fenced-handoff-lookup-stalled'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, undefined, onError); + const getJobSpy = jest.spyOn(jobStore, 'getJob'); + jest.useFakeTimers(); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + getJobSpy.mockImplementation(() => new Promise(() => undefined)); + + expect(job.abortController.signal.aborted).toBe(true); + await jest.advanceTimersByTimeAsync(REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS); + + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + + expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + getJobSpy.mockRestore(); + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('lets queued predecessor done drain after its successor job disappears', async () => { + const streamId = 'stream-fenced-handoff-successor-gone'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new DeferredEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + manager.initialize(); + const predecessor = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onDone = jest.fn(); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, onDone, onError); + jest.useFakeTimers(); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + await jest.advanceTimersByTimeAsync(1); + const successor = await jobStore.createJob( + streamId, + 'user-1', + 'conversation-1', + undefined, + { generationProtocolVersion: 2 }, + undefined, + undefined, + undefined, + undefined, + undefined, + 'successor-gone-attempt', + ); + eventTransport.emitDone(streamId, { final: true }, predecessor.createdAt); + await jobStore.deleteJob(streamId, successor.createdAt); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS - 1); + + expect(onDone).not.toHaveBeenCalled(); + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + + eventTransport.flushDeliveries(); + + expect(onDone).toHaveBeenCalledWith({ final: true }); + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('preserves a full terminal drain when a pending handoff settles at its deadline', async () => { + const streamId = 'stream-fenced-handoff-deadline-drain'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new DeferredEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + manager.initialize(); + const predecessor = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onDone = jest.fn(); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, onDone, onError); + jest.useFakeTimers(); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + await jest.advanceTimersByTimeAsync(1); + const successor = await jobStore.createJob( + streamId, + 'user-1', + 'conversation-1', + undefined, + { generationProtocolVersion: 2 }, + undefined, + undefined, + undefined, + undefined, + undefined, + 'deadline-successor-attempt', + ); + + await jest.advanceTimersByTimeAsync(REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS - 2); + eventTransport.emitDone(streamId, { final: true }, predecessor.createdAt); + await jobStore.acknowledgeReplacedJobs?.(streamId, successor.creationAttemptId!, [ + predecessor.createdAt, + ]); + await jest.advanceTimersByTimeAsync(1); + + expect(onDone).not.toHaveBeenCalled(); + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + + eventTransport.flushDeliveries(); + + expect(onDone).toHaveBeenCalledWith({ final: true }); + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('retires a durably fenced runtime immediately when no subscriber is attached', async () => { + const streamId = 'stream-fenced-without-subscriber'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = Object.assign(new InMemoryEventTransport(), { + recordAbortAcknowledgement: jest.fn().mockResolvedValue(true), + }); + registerChunkPublicationCapability(eventTransport, async () => false); + const getJob = jest.spyOn(jobStore, 'getJob'); + const clearContentState = jest.spyOn(jobStore, 'clearContentState'); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: true }); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + getJob.mockClear(); + clearContentState.mockClear(); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + + expect(job.abortController.signal.aborted).toBe(true); + expect(manager.getRuntimeStats()).toMatchObject({ + runtimeStateSize: 0, + earlyBufferedEvents: 0, + earlyBufferedBytes: 0, + }); + expect(clearContentState).toHaveBeenCalledWith(streamId, job.createdAt); + expect(getJob).not.toHaveBeenCalled(); + expect(eventTransport.recordAbortAcknowledgement).toHaveBeenCalledWith( + streamId, + job.createdAt, + ); + } finally { + await manager.destroy(); + } + }); + + it('does not record abort proof for a runtime this manager does not own', async () => { + const streamId = 'stream-fenced-non-owner'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = Object.assign(new InMemoryEventTransport(), { + recordAbortAcknowledgement: jest.fn().mockResolvedValue(true), + }); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: true }); + manager.initialize(); + const durableJob = await jobStore.createJob(streamId, 'user-1', 'conversation-1'); + + try { + await expect(manager.getJob(streamId)).resolves.toBeDefined(); + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + + expect(eventTransport.recordAbortAcknowledgement).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + await expect(jobStore.getJob(streamId)).resolves.toMatchObject({ + createdAt: durableJob.createdAt, + status: 'running', + }); + } finally { + await manager.destroy(); + } + }); + + it('cancels a fenced retirement when its last subscriber detaches', async () => { + const streamId = 'stream-fenced-before-subscriber-detach'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const originalGetJob = jobStore.getJob.bind(jobStore); + const getJob = jest.spyOn(jobStore, 'getJob'); + const clearContentState = jest.spyOn(jobStore, 'clearContentState'); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: true }); + jest.useFakeTimers(); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onAllSubscribersLeft = jest.fn(); + job.emitter.on('allSubscribersLeft', onAllSubscribersLeft); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, undefined, onError); + await jest.advanceTimersByTimeAsync(0); + getJob.mockClear(); + clearContentState.mockClear(); + let signalLookupStarted: (() => void) | undefined; + const lookupStarted = new Promise((resolve) => { + signalLookupStarted = resolve; + }); + let releaseLookup: (() => void) | undefined; + const lookupGate = new Promise((resolve) => { + releaseLookup = resolve; + }); + getJob.mockImplementationOnce(async (...args) => { + signalLookupStarted?.(); + await lookupGate; + return originalGetJob(...args); + }); + const timerCountBeforeFence = jest.getTimerCount(); + let detached = false; + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + + expect(job.abortController.signal.aborted).toBe(true); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + expect(jest.getTimerCount()).toBe(timerCountBeforeFence + 1); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + await lookupStarted; + + subscription?.unsubscribe(); + detached = true; + await jest.advanceTimersByTimeAsync(0); + + expect(manager.getRuntimeStats()).toMatchObject({ + runtimeStateSize: 0, + earlyBufferedEvents: 0, + earlyBufferedBytes: 0, + }); + expect(clearContentState).toHaveBeenCalledWith(streamId, job.createdAt); + expect(getJob).toHaveBeenCalledTimes(1); + expect(onError).not.toHaveBeenCalled(); + expect(onAllSubscribersLeft).not.toHaveBeenCalled(); + expect(jest.getTimerCount()).toBe(timerCountBeforeFence); + + releaseLookup?.(); + await jest.advanceTimersByTimeAsync(0); + + await jest.advanceTimersByTimeAsync( + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS + REDIS_EVENT_REORDER_TIMEOUT_MS * 2, + ); + + expect(getJob).toHaveBeenCalledTimes(1); + expect(onError).not.toHaveBeenCalled(); + expect(onAllSubscribersLeft).not.toHaveBeenCalled(); + } finally { + releaseLookup?.(); + getJob.mockRestore(); + if (!detached) { + subscription?.unsubscribe(); + } + await manager.destroy(); + jest.useRealTimers(); + } + }); + + it.each(['manual', 'timer'] as const)( + 'does not apply a %s fenced predecessor detach to an unattached successor', + async (detachment) => { + const streamId = 'stream-fenced-detach-with-successor'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new DeferredEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + jest.useFakeTimers(); + manager.initialize(); + const predecessor = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const predecessorError = jest.fn(); + const predecessorSubscription = await manager.subscribe( + streamId, + () => undefined, + undefined, + predecessorError, + ); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + expect(predecessor.abortController.signal.aborted).toBe(true); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(1); + + const successor = await manager.createJob(streamId, 'user-1', 'conversation-2'); + const successorAllSubscribersLeft = jest.fn(); + successor.emitter.on('allSubscribersLeft', successorAllSubscribersLeft); + const updateJob = jest.spyOn(jobStore, 'updateJob'); + updateJob.mockClear(); + + if (detachment === 'manual') { + predecessorSubscription?.unsubscribe(); + } else { + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + } + await jest.advanceTimersByTimeAsync(0); + + expect(successor.abortController.signal.aborted).toBe(false); + if (detachment === 'timer') { + expect(predecessorError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + } else { + expect(predecessorError).not.toHaveBeenCalled(); + } + expect(successorAllSubscribersLeft).not.toHaveBeenCalled(); + expect(updateJob).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats()).toMatchObject({ + runtimeStateSize: 1, + fencedRuntimeRetirements: 0, + }); + await expect(manager.hasJob(streamId)).resolves.toBe(true); + } finally { + predecessorSubscription?.unsubscribe(); + await manager.destroy(); + jest.useRealTimers(); + } + }, + ); + + it('cancels a scheduled fenced retirement when services are reconfigured', async () => { + const streamId = 'stream-fenced-before-reconfigure'; + const oldJobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const oldEventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(oldEventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore: oldJobStore, eventTransport: oldEventTransport, isRedis: false }); + jest.useFakeTimers(); + manager.initialize(); + const oldJob = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const oldError = jest.fn(); + const oldSubscription = await manager.subscribe(streamId, () => undefined, undefined, oldError); + await jest.advanceTimersByTimeAsync(0); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + expect(oldJob.abortController.signal.aborted).toBe(true); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(1); + + const currentJobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const currentEventTransport = new InMemoryEventTransport(); + const currentGetJob = jest.spyOn(currentJobStore, 'getJob'); + const currentClearContentState = jest.spyOn(currentJobStore, 'clearContentState'); + manager.configure({ + jobStore: currentJobStore, + eventTransport: currentEventTransport, + isRedis: false, + }); + await manager.createJob(streamId, 'user-1', 'conversation-2'); + currentGetJob.mockClear(); + currentClearContentState.mockClear(); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(0); + + await jest.advanceTimersByTimeAsync( + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS + REDIS_EVENT_REORDER_TIMEOUT_MS * 2, + ); + + expect(oldError).not.toHaveBeenCalled(); + expect(currentGetJob).not.toHaveBeenCalled(); + expect(currentClearContentState).not.toHaveBeenCalled(); + await expect(manager.hasJob(streamId)).resolves.toBe(true); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + } finally { + jest.useRealTimers(); + oldSubscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('cancels an in-flight fenced retirement lookup when services are reconfigured', async () => { + const streamId = 'stream-fenced-lookup-before-reconfigure'; + const oldJobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const oldEventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(oldEventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore: oldJobStore, eventTransport: oldEventTransport, isRedis: false }); + jest.useFakeTimers(); + manager.initialize(); + const oldJob = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const oldError = jest.fn(); + const oldSubscription = await manager.subscribe(streamId, () => undefined, undefined, oldError); + await jest.advanceTimersByTimeAsync(0); + const originalGetJob = oldJobStore.getJob.bind(oldJobStore); + let signalLookupStarted: (() => void) | undefined; + const lookupStarted = new Promise((resolve) => { + signalLookupStarted = resolve; + }); + let releaseLookup: (() => void) | undefined; + const lookupGate = new Promise((resolve) => { + releaseLookup = resolve; + }); + const getJobSpy = jest.spyOn(oldJobStore, 'getJob').mockImplementationOnce(async (...args) => { + signalLookupStarted?.(); + await lookupGate; + return originalGetJob(...args); + }); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + expect(oldJob.abortController.signal.aborted).toBe(true); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(1); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + await lookupStarted; + + const currentJobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const currentEventTransport = new InMemoryEventTransport(); + const currentGetJob = jest.spyOn(currentJobStore, 'getJob'); + const currentClearContentState = jest.spyOn(currentJobStore, 'clearContentState'); + manager.configure({ + jobStore: currentJobStore, + eventTransport: currentEventTransport, + isRedis: false, + }); + await manager.createJob(streamId, 'user-1', 'conversation-2'); + currentGetJob.mockClear(); + currentClearContentState.mockClear(); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(0); + releaseLookup?.(); + await Promise.resolve(); + + await jest.advanceTimersByTimeAsync( + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS + REDIS_EVENT_REORDER_TIMEOUT_MS * 2, + ); + + expect(oldError).not.toHaveBeenCalled(); + expect(currentGetJob).not.toHaveBeenCalled(); + expect(currentClearContentState).not.toHaveBeenCalled(); + await expect(manager.hasJob(streamId)).resolves.toBe(true); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + } finally { + releaseLookup?.(); + getJobSpy.mockRestore(); + jest.useRealTimers(); + oldSubscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('cancels fenced retirement timers when the manager is destroyed', async () => { + const streamId = 'stream-fenced-before-destroy'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + jest.useFakeTimers(); + const timerCountBeforeInitialize = jest.getTimerCount(); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, undefined, onError); + await jest.advanceTimersByTimeAsync(0); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + expect(job.abortController.signal.aborted).toBe(true); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(1); + + await manager.destroy(); + const errorCountAfterDestroy = onError.mock.calls.length; + + expect(onError).not.toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(manager.getRuntimeStats().fencedRuntimeRetirements).toBe(0); + expect(jest.getTimerCount()).toBe(timerCountBeforeInitialize); + + await jest.advanceTimersByTimeAsync( + REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS + REDIS_EVENT_REORDER_TIMEOUT_MS * 2, + ); + + expect(onError).toHaveBeenCalledTimes(errorCountAfterDestroy); + expect(onError).not.toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(jest.getTimerCount()).toBe(timerCountBeforeInitialize); + } finally { + jest.useRealTimers(); + subscription?.unsubscribe(); + } + }); + + it('does not start a fenced retirement after shutdown preparation', async () => { + const streamId = 'stream-fenced-during-shutdown'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new InMemoryEventTransport(); + let signalPublicationStarted: (() => void) | undefined; + const publicationStarted = new Promise((resolve) => { + signalPublicationStarted = resolve; + }); + let resolvePublication: ((published: false) => void) | undefined; + const publicationGate = new Promise((resolve) => { + resolvePublication = resolve; + }); + registerChunkPublicationCapability(eventTransport, async () => { + signalPublicationStarted?.(); + return publicationGate; + }); + const clearContentState = jest.spyOn(jobStore, 'clearContentState'); + const manager = new GenerationJobManagerClass(); + manager.configure({ jobStore, eventTransport, isRedis: false }); + jest.useFakeTimers(); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const subscription = await manager.subscribe(streamId, () => undefined); + clearContentState.mockClear(); + let emitting: Promise | undefined; + + try { + emitting = manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + await publicationStarted; + + manager.prepareForShutdown(); + resolvePublication?.(false); + await emitting; + + expect(job.abortController.signal.aborted).toBe(true); + expect(clearContentState).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + } finally { + resolvePublication?.(false); + await emitting?.catch(() => undefined); + subscription?.unsubscribe(); + await manager.destroy(); + jest.useRealTimers(); + } + }); + + it('reconnect-closes a fenced generation when no terminal arrives within the drain window', async () => { + const streamId = 'stream-fenced-terminal-missing'; + const eventTransport = new InMemoryEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const manager = new GenerationJobManagerClass(); + manager.configure({ + jobStore: new InMemoryJobStore({ ttlAfterComplete: 60_000 }), + eventTransport, + isRedis: false, + }); + manager.initialize(); + const job = await manager.createJob(streamId, 'user-1', 'conversation-1'); + const onError = jest.fn(); + const subscription = await manager.subscribe(streamId, () => undefined, undefined, onError); + jest.useFakeTimers(); + + try { + await manager.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + + expect(job.abortController.signal.aborted).toBe(true); + expect(onError).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS - 1); + expect(onError).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync(1); + expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); + } finally { + jest.useRealTimers(); + subscription?.unsubscribe(); + await manager.destroy(); + } + }); + + it('reconnect-closes only the fenced predecessor after a replacement installs during grace', async () => { + const streamId = 'stream-fenced-predecessor-replaced'; + const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); + const eventTransport = new DeferredEventTransport(); + registerChunkPublicationCapability(eventTransport, async () => false); + const owner = new GenerationJobManagerClass(); + const replacer = new GenerationJobManagerClass(); + owner.configure({ jobStore, eventTransport, isRedis: false }); + replacer.configure({ jobStore, eventTransport, isRedis: false }); + owner.initialize(); + replacer.initialize(); + const predecessor = await owner.createJob(streamId, 'user-1', 'conversation-1'); + const predecessorError = jest.fn(); + const predecessorSubscription = await owner.subscribe( + streamId, + () => undefined, + undefined, + predecessorError, + ); + let replacementSubscription: Awaited> = null; + jest.useFakeTimers(); + + try { + await owner.emitChunk(streamId, { + event: 'on_message_delta', + data: { id: 'step-1', delta: { content: [{ type: 'text', text: 'tail' }] } }, + }); + expect(predecessor.abortController.signal.aborted).toBe(true); + + const replacement = await replacer.createJob(streamId, 'user-1', 'conversation-1'); + const current = await owner.getJob(streamId); + const replacementDone = jest.fn(); + const replacementError = jest.fn(); + replacementSubscription = await owner.subscribe( + streamId, + () => undefined, + replacementDone, + replacementError, + { expectedCreatedAt: replacement.createdAt }, + ); + + expect(replacement.createdAt).toBeGreaterThan(predecessor.createdAt); + expect(current?.createdAt).toBe(replacement.createdAt); + expect(replacementSubscription).not.toBeNull(); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + + expect(predecessorError).not.toHaveBeenCalled(); + expect(replacementError).not.toHaveBeenCalled(); + expect(replacementDone).not.toHaveBeenCalled(); + expect(owner.getRuntimeStats().runtimeStateSize).toBe(1); + expect(eventTransport.getSubscriberCount(streamId)).toBe(2); + + eventTransport.flushDeliveries(); + + expect(predecessorError).not.toHaveBeenCalled(); + expect(replacementError).not.toHaveBeenCalled(); + expect(replacementDone).not.toHaveBeenCalled(); + expect(owner.getRuntimeStats().runtimeStateSize).toBe(1); + expect(eventTransport.getSubscriberCount(streamId)).toBe(1); + + await jest.advanceTimersByTimeAsync(REDIS_EVENT_REORDER_TIMEOUT_MS * 2); + + expect(predecessorError).not.toHaveBeenCalled(); + expect(replacementError).not.toHaveBeenCalled(); + expect(replacementDone).not.toHaveBeenCalled(); + expect(owner.getRuntimeStats().runtimeStateSize).toBe(1); + expect(eventTransport.getSubscriberCount(streamId)).toBe(1); + } finally { + jest.useRealTimers(); + predecessorSubscription?.unsubscribe(); + replacementSubscription?.unsubscribe(); + await Promise.all([owner.destroy(), replacer.destroy()]); + } + }); + it('does not return a lazy runtime replaced while its abort listener activates', async () => { const now = jest.spyOn(Date, 'now').mockReturnValue(1000); const jobStore = new InMemoryJobStore({ ttlAfterComplete: 60_000 }); diff --git a/packages/api/src/stream/__tests__/steerReceiptIntegrity.spec.ts b/packages/api/src/stream/__tests__/steerReceiptIntegrity.spec.ts index 9602960389..248f2e1d03 100644 --- a/packages/api/src/stream/__tests__/steerReceiptIntegrity.spec.ts +++ b/packages/api/src/stream/__tests__/steerReceiptIntegrity.spec.ts @@ -7,6 +7,7 @@ import { InMemoryEventTransport } from '../implementations/InMemoryEventTranspor import { registerChunkPublicationCapability } from '../internal/chunkPublication'; import { InMemoryJobStore } from '../implementations/InMemoryJobStore'; import { STEER_ENQUEUE_RECEIPT_FULL } from '../interfaces/IJobStore'; +import { REDIS_ABORT_TERMINAL_GRACE_MS } from '../internal/timing'; type ReceiptEntry = { receipt: SteerReceipt; expiresAt: number }; type ReceiptMap = Map; @@ -327,6 +328,7 @@ describe('InMemoryJobStore steer receipt integrity', () => { }; const hostContent: Array<{ steerId: string; text: string }> = []; let subscription: { unsubscribe: () => void } | null = null; + jest.useFakeTimers(); try { const job = await manager.createJob(streamId, item.userId, streamId, { @@ -376,9 +378,15 @@ describe('InMemoryJobStore steer receipt integrity', () => { item, }); expect(job.abortController.signal.aborted).toBe(true); + expect(onError).not.toHaveBeenCalled(); + expect(manager.getRuntimeStats().runtimeStateSize).toBe(1); + + await jest.advanceTimersByTimeAsync(REDIS_ABORT_TERMINAL_GRACE_MS); + expect(onError).toHaveBeenCalledWith(TERMINAL_PUBLICATION_RECONNECT_ERROR); expect(manager.getRuntimeStats().runtimeStateSize).toBe(0); } finally { + jest.useRealTimers(); subscription?.unsubscribe(); await manager.destroy(); } diff --git a/packages/api/src/stream/implementations/RedisEventTransport.ts b/packages/api/src/stream/implementations/RedisEventTransport.ts index 1e97b8b52e..d172fb388c 100644 --- a/packages/api/src/stream/implementations/RedisEventTransport.ts +++ b/packages/api/src/stream/implementations/RedisEventTransport.ts @@ -8,6 +8,10 @@ import { MAX_COALESCED_EVENTS, resolveCoalesceWindowMs, } from '~/stream/internal/coalescing'; +import { + REDIS_ABORT_ACK_TIMEOUT_MS, + REDIS_EVENT_REORDER_TIMEOUT_MS, +} from '~/stream/internal/timing'; import { registerChunkPublicationCapability } from '~/stream/internal/chunkPublication'; import { instrumentIORedisClient, RedisUseCases } from '~/cache/redisTelemetry'; @@ -196,16 +200,10 @@ const PUBLISH_REPLACED_DONE_LUA = 'redis.call("EXPIRE", KEYS[1], ttl) end local seq = val - 1 ' + 'redis.call("PUBLISH", ARGV[1], ARGV[2] .. string.format("%d", seq) .. ARGV[3]) return seq'; -/** Max time (ms) to wait for out-of-order messages before force-flushing */ -const REORDER_TIMEOUT_MS = 500; /** Max messages to buffer before force-flushing (prevents memory issues) */ const MAX_BUFFER_SIZE = 100; /** Rolling-upgrade recovery window after a legacy job hash expires without an epoch marker. */ const GENERATION_EPOCH_GRACE_TTL_SECONDS = 300; -/** A replacement remains fail-closed if its exact generation owner cannot - * acknowledge promptly. Redis pub/sub is local-network traffic; this budget - * tolerates reconnect jitter without holding an HTTP request indefinitely. */ -const ABORT_ACK_TIMEOUT_MS = 3000; /** Durable owner proof outlives receipt retries and process-local subscriptions. */ const ABORT_ACK_TTL_SECONDS = 86400; @@ -870,7 +868,7 @@ export class RedisEventTransport implements IEventTransport { ); this.forceFlushBuffer(streamId, streamState); } - }, REORDER_TIMEOUT_MS); + }, REDIS_EVENT_REORDER_TIMEOUT_MS); } /** Deliver a message to all handlers */ @@ -1257,11 +1255,7 @@ export class RedisEventTransport implements IEventTransport { return (await this.publisher.get(KEYS.abortAck(streamId, generationId))) === '1'; } - private async publishAbortAcknowledgement( - streamId: string, - generationId: number, - abortRequestId: string, - ): Promise { + async recordAbortAcknowledgement(streamId: string, generationId: number): Promise { try { await this.publisher.set( KEYS.abortAck(streamId, generationId), @@ -1269,8 +1263,19 @@ export class RedisEventTransport implements IEventTransport { 'EX', ABORT_ACK_TTL_SECONDS, ); + return true; } catch (error) { logger.error(`[RedisEventTransport] Failed to persist generation abort proof:`, error); + return false; + } + } + + private async publishAbortAcknowledgement( + streamId: string, + generationId: number, + abortRequestId: string, + ): Promise { + if (!(await this.recordAbortAcknowledgement(streamId, generationId))) { // A live acknowledgement is only useful if a racing or inherited receipt // can prove the same owner stop after this subscription disappears. A // SET that committed despite a lost reply is recovered by the requester's @@ -1316,7 +1321,7 @@ export class RedisEventTransport implements IEventTransport { (acknowledged) => this.settleAbortAck(streamId, state, abortRequestId, acknowledged), () => this.settleAbortAck(streamId, state, abortRequestId, false), ); - }, ABORT_ACK_TIMEOUT_MS); + }, REDIS_ABORT_ACK_TIMEOUT_MS); state.abortAckWaiters.set(abortRequestId, { generationId, resolve, timeout }); }); diff --git a/packages/api/src/stream/interfaces/IJobStore.ts b/packages/api/src/stream/interfaces/IJobStore.ts index 2650a9e257..ef49e2b472 100644 --- a/packages/api/src/stream/interfaces/IJobStore.ts +++ b/packages/api/src/stream/interfaces/IJobStore.ts @@ -1273,6 +1273,10 @@ export interface IEventTransport { * the replica owning that generation processes the abort. */ emitAbortConfirmed?(streamId: string, generationId: number): Promise; + /** Persist proof that this process synchronously stopped the exact generation. + * A delayed replacement can use the proof after the owner's listeners retire. */ + recordAbortAcknowledgement?(streamId: string, generationId: number): Promise; + /** Publish a predecessor DONE only while the current job's opaque creation * attempt still carries that predecessor in its durable receipt chain. */ emitReplacedDoneConfirmed?( diff --git a/packages/api/src/stream/internal/timing.ts b/packages/api/src/stream/internal/timing.ts new file mode 100644 index 0000000000..523302f722 --- /dev/null +++ b/packages/api/src/stream/internal/timing.ts @@ -0,0 +1,16 @@ +/** Maximum time Redis delivery may buffer an out-of-order event. */ +export const REDIS_EVENT_REORDER_TIMEOUT_MS = 500; + +/** Maximum time a replacement waits for its exact generation owner to + * acknowledge an abort before it falls back to durable proof. */ +export const REDIS_ABORT_ACK_TIMEOUT_MS = 3_000; + +/** A replacement publishes predecessor DONE only after abort confirmation. + * Preserve the captured subscriber for that full acknowledgement budget plus + * the DONE reorder window and an equal scheduling/Redis settlement margin. */ +export const REDIS_ABORT_TERMINAL_GRACE_MS = + REDIS_ABORT_ACK_TIMEOUT_MS + REDIS_EVENT_REORDER_TIMEOUT_MS * 2; + +/** Maximum time spent inspecting a durable replacement handoff. One final + * event-reordering grace may follow before the captured subscriber is recycled. */ +export const REDIS_REPLACEMENT_HANDOFF_MAX_WAIT_MS = 30_000;