/** * Private replay/live-drain algorithm shared by every `createReplayLiveFeed()` * instance (`workflow-event-feed.ts`'s per-workflow feed and * `fleet-event-feed.ts`'s fleet-wide feed). Split out of * `workflow-event-feed.ts` purely to keep that file under the repository's * implementation-file line limit — nothing here is public API. * * @module server/replay-live-feed-internals */ import type { ReplayLiveFeedBackend, ReplayLiveSubscribeOptions, SequencedEventEnvelope } from './workflow-event-feed.ts'; export declare function createSerialOperationQueue(): (operation: () => Promise) => Promise; type DurableSubscriptionOptions = { readonly pollIntervalMs: number; readonly lifecycleSignal?: AbortSignal; }; export declare function createDurableSubscription(backend: ReplayLiveFeedBackend, options: DurableSubscriptionOptions, args?: ReplayLiveSubscribeOptions): AsyncIterable; /** Thrown when a replay window exceeds the caller's configured limit. */ export declare class ReplayWindowExceededError extends Error { readonly count: number; readonly limit: number; constructor(count: number, limit: number); } export declare function replayUpTo(backend: ReplayLiveFeedBackend, afterSequence: number, snapshot: number, signal: AbortSignal | undefined, replayOptions: ReplayLiveSubscribeOptions | undefined): AsyncIterable; export declare function shouldDeliverEnvelope(envelope: TEnvelope, replayOptions: ReplayLiveSubscribeOptions | undefined): boolean; export declare function drainLive(buffer: TEnvelope[], snapshot: number, signal: AbortSignal | undefined, overflowed: () => boolean, installWaker: (fn: (() => void) | null) => void): AsyncIterable; export {};