📡 fix: Preserve Redis Abort Terminal Delivery (#14749)

* 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
This commit is contained in:
Danny Avila 2026-08-12 07:33:07 -04:00 committed by GitHub
parent 88e08c91e8
commit b2128a7d18
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 1448 additions and 65 deletions

View file

@ -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<void>>;
}
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<RuntimeJobState, FencedRuntimeRetirementContext>();
/** 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<string>();
/** 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<typeof setTimeout> | 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<void> {
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<void> {
this.shuttingDown = true;
this.cancelFencedRuntimeRetirements();
if (this.cleanupInterval) {
clearInterval(this.cleanupInterval);

View file

@ -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<void>((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();

View file

@ -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 () => {

View file

@ -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<void>((resolve) => {
signalLegacyMarkStarted = resolve;
});
let releaseLegacyMark: (() => void) | undefined;
const legacyMarkGate = new Promise<void>((resolve) => {
releaseLegacyMark = resolve;
});
const markStartedSpy = jest
.spyOn(jobStore, 'markIdempotencyKeyStarted')
.mockImplementationOnce(async (...args) => {
signalLegacyMarkStarted?.();
await legacyMarkGate;
return actualMarkStarted(...args);
});
let creatingReplacement: ReturnType<typeof replacer.createJob> | 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<null>(() => 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<void>((resolve) => {
signalLookupStarted = resolve;
});
let releaseLookup: (() => void) | undefined;
const lookupGate = new Promise<void>((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<void>((resolve) => {
signalLookupStarted = resolve;
});
let releaseLookup: (() => void) | undefined;
const lookupGate = new Promise<void>((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<void>((resolve) => {
signalPublicationStarted = resolve;
});
let resolvePublication: ((published: false) => void) | undefined;
const publicationGate = new Promise<false>((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<void> | 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<ReturnType<typeof owner.subscribe>> = 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 });

View file

@ -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<string, ReceiptEntry>;
@ -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();
}

View file

@ -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<void> {
async recordAbortAcknowledgement(streamId: string, generationId: number): Promise<boolean> {
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<void> {
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 });
});

View file

@ -1273,6 +1273,10 @@ export interface IEventTransport {
* the replica owning that generation processes the abort. */
emitAbortConfirmed?(streamId: string, generationId: number): Promise<boolean>;
/** 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<boolean>;
/** Publish a predecessor DONE only while the current job's opaque creation
* attempt still carries that predecessor in its durable receipt chain. */
emitReplacedDoneConfirmed?(

View file

@ -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;