type Subscriber = (handleId: THandleId) => void; type StageSubscriber = (stage: TStage) => void; type LifecycleStage = "uninitialized" | "ready" | "invalidated"; type StageMap = { uninitialized: TStage; ready: TStage; invalidated: TStage; }; type ActiveCommand = { epoch: number; invalidatedHandleIds: Set; }; const DEFAULT_LIFECYCLE_STAGES: StageMap = { uninitialized: "uninitialized", ready: "ready", invalidated: "invalidated", }; /** * Serializes leaf command dispatches for a lazily resolved handle. * Command callbacks must not enqueue another command on the same client. */ export type LazyHandleClient< THandleId extends string, TSettlement, TStage extends string = LifecycleStage, > = { runCommand( buildPayload: (handleId: THandleId | undefined) => TPayload, dispatchAndWait: (payload: TPayload) => Promise, ): Promise; runCommandWithRetry( buildPayload: (handleId: THandleId | undefined) => TPayload, dispatchAndWait: (payload: TPayload) => Promise, options: { shouldRetry: (error: unknown) => boolean; onBeforeRetry?: (error: unknown) => void; }, ): Promise; /** At its FIFO turn, dispatches with the current handle or no-ops if absent. */ runResolvedHandleCommand( buildPayload: (handleId: THandleId) => TPayload, dispatchAndWait: (payload: TPayload) => Promise, ): Promise; getHandleId(): THandleId | undefined; /** Clears the resolved handle before retrying a stale-handle command. */ clearHandleForRetry(): THandleId | undefined; /** Records an authoritative invalidation for the active command. */ recordHandleInvalidation(handleId: string): THandleId | undefined; /** * Clears a matching resolved handle without cancelling active command * recovery. */ clearHandleIfMatching(handleId: string): THandleId | undefined; hasPendingCommands(): boolean; subscribeHandle(subscriber: Subscriber): () => void; getStage(): TStage; setStage(stage: TStage): void; subscribeStage(subscriber: StageSubscriber): () => void; reset(): void; invalidate(): void; }; export type CreateLazyHandleClientOptions< THandleId extends string, TSettlement, TStage extends string = LifecycleStage, > = { getHandleIdFromSettlement: (settlement: TSettlement) => THandleId | undefined; stages?: StageMap; initialStage?: TStage; }; type CreateLazyHandleClientDefaultOptions< THandleId extends string, TSettlement, > = { getHandleIdFromSettlement: (settlement: TSettlement) => THandleId | undefined; initialStage?: LifecycleStage; stages?: undefined; }; type CreateLazyHandleClientCustomOptions< THandleId extends string, TSettlement, TStage extends string, > = { getHandleIdFromSettlement: (settlement: TSettlement) => THandleId | undefined; stages: StageMap; initialStage?: TStage; }; /** * Serializes commands for a lazily initialized handle. * Each command payload is built only when its turn runs, so it always sees * the latest settled handle id from prior commands. */ const createLazyHandleClientWithStages = < THandleId extends string, TSettlement, TStage extends string, >(input: { getHandleIdFromSettlement: (settlement: TSettlement) => THandleId | undefined; stages: StageMap; initialStage?: TStage; }): LazyHandleClient => { let handleId: THandleId | undefined; let externallyInvalidatedHandleId: THandleId | undefined; let tail = Promise.resolve(); let queuedCommandCount = 0; let epoch = 0; let activeCommand: ActiveCommand | undefined; const { stages } = input; let stage = input.initialStage ?? stages.uninitialized; const subscribers = new Set>(); const stageSubscribers = new Set>(); const enqueue = ( execute: () => Promise, ): Promise => { const run = queuedCommandCount === 0 ? execute() : tail.then(execute, execute); queuedCommandCount += 1; tail = run.then( () => { queuedCommandCount -= 1; }, () => { queuedCommandCount -= 1; }, ); return run; }; const setStage = (nextStage: TStage): void => { if (stage === nextStage) { return; } if (stage === stages.invalidated && nextStage !== stages.invalidated) { return; } stage = nextStage; stageSubscribers.forEach((subscriber) => { subscriber(nextStage); }); }; const publishHandle = ( nextHandleId: THandleId, command: ActiveCommand, ): void => { if (stage === stages.invalidated) { return; } if ( nextHandleId === externallyInvalidatedHandleId || command.invalidatedHandleIds.has(nextHandleId) ) { return; } externallyInvalidatedHandleId = undefined; handleId = nextHandleId; setStage(stages.ready); subscribers.forEach((subscriber) => { subscriber(nextHandleId); }); }; const publishHandleFromSettlement = ( settlement: TSettlement, command: ActiveCommand, ): void => { const settledHandleId = input.getHandleIdFromSettlement(settlement); if (settledHandleId && command.epoch === epoch) { publishHandle(settledHandleId, command); } }; const runCommand: LazyHandleClient< THandleId, TSettlement, TStage >["runCommand"] = (buildPayload, dispatchAndWait) => { const execute = async (): Promise => { if (stage === stages.invalidated) { throw new Error("Handle client is invalidated"); } const command = { epoch, invalidatedHandleIds: new Set(), }; activeCommand = command; try { const payload = buildPayload(handleId); const settlement = await dispatchAndWait(payload); publishHandleFromSettlement(settlement, command); return settlement; } finally { if (activeCommand === command) { activeCommand = undefined; } } }; return enqueue(execute); }; const runCommandWithRetry: LazyHandleClient< THandleId, TSettlement, TStage >["runCommandWithRetry"] = (buildPayload, dispatchAndWait, options) => { const execute = async (): Promise => { if (stage === stages.invalidated) { throw new Error("Handle client is invalidated"); } const command = { epoch, invalidatedHandleIds: new Set(), }; activeCommand = command; try { let didRetry = false; for (;;) { if (stage === stages.invalidated) { throw new Error("Handle client is invalidated"); } const payload = buildPayload(handleId); try { const settlement = await dispatchAndWait(payload); publishHandleFromSettlement(settlement, command); return settlement; } catch (error) { if ( command.epoch === epoch && !didRetry && options.shouldRetry(error) ) { didRetry = true; options.onBeforeRetry?.(error); if ( activeCommand !== command || command.epoch !== epoch || stage === stages.invalidated ) { throw error; } continue; } throw error; } } } finally { if (activeCommand === command) { activeCommand = undefined; } } }; return enqueue(execute); }; const runResolvedHandleCommand: LazyHandleClient< THandleId, TSettlement, TStage >["runResolvedHandleCommand"] = (buildPayload, dispatchAndWait) => enqueue(async () => { if (stage === stages.invalidated) { throw new Error("Handle client is invalidated"); } if (handleId === undefined) { return; } await dispatchAndWait(buildPayload(handleId)); }); const clearHandleIfMatching = ( expectedHandleId: string, ): THandleId | undefined => { if (handleId !== expectedHandleId) { return undefined; } const previousHandleId = handleId; handleId = undefined; externallyInvalidatedHandleId = previousHandleId; setStage(stages.uninitialized); return previousHandleId; }; const recordHandleInvalidation = ( invalidatedHandleId: string, ): THandleId | undefined => { if (activeCommand?.epoch === epoch) { activeCommand.invalidatedHandleIds.add(invalidatedHandleId); } return clearHandleIfMatching(invalidatedHandleId); }; const clearHandleForRetry = (): THandleId | undefined => { if (!handleId) { return undefined; } const previousHandleId = handleId; handleId = undefined; externallyInvalidatedHandleId = undefined; setStage(stages.uninitialized); return previousHandleId; }; return { runCommand, runCommandWithRetry, runResolvedHandleCommand, getHandleId: () => handleId, clearHandleForRetry, recordHandleInvalidation, clearHandleIfMatching, hasPendingCommands: () => queuedCommandCount > 0, subscribeHandle: (subscriber) => { subscribers.add(subscriber); if (handleId) { subscriber(handleId); } return () => { subscribers.delete(subscriber); }; }, getStage: () => stage, setStage, subscribeStage: (subscriber) => { stageSubscribers.add(subscriber); subscriber(stage); return () => { stageSubscribers.delete(subscriber); }; }, reset: () => { epoch += 1; activeCommand = undefined; handleId = undefined; externallyInvalidatedHandleId = undefined; stage = stages.uninitialized; stageSubscribers.forEach((subscriber) => { subscriber(stage); }); }, invalidate: () => { epoch += 1; activeCommand = undefined; handleId = undefined; externallyInvalidatedHandleId = undefined; setStage(stages.invalidated); }, }; }; export function createLazyHandleClient( options: CreateLazyHandleClientDefaultOptions, ): LazyHandleClient; export function createLazyHandleClient< THandleId extends string, TSettlement, TStage extends string, >( options: CreateLazyHandleClientCustomOptions, ): LazyHandleClient; export function createLazyHandleClient( options: | CreateLazyHandleClientDefaultOptions | CreateLazyHandleClientCustomOptions, ): LazyHandleClient { if ("stages" in options && options.stages) { return createLazyHandleClientWithStages({ getHandleIdFromSettlement: options.getHandleIdFromSettlement, stages: options.stages, initialStage: options.initialStage, }); } return createLazyHandleClientWithStages({ getHandleIdFromSettlement: options.getHandleIdFromSettlement, stages: DEFAULT_LIFECYCLE_STAGES, initialStage: options.initialStage, }); }