import { createLazyHandleClient, type LazyHandleClient, } from "./lazy-handle-client.js"; /** * Serializes leaf flow-command dispatches around a lazily resolved handle. * Command callbacks must not enqueue another command on the same controller. */ export type FlowCommandHandleController< THandleId extends string, TSettlement, > = { runCommand( buildPayload: (handleId: THandleId | undefined) => TPayload, dispatchAndWait: (payload: TPayload) => Promise, ): 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; hasHandleId(handleId: string): boolean; /** * Records an authoritative invalidation without cancelling active * stale-handle recovery. * * The return value reports only whether a matching resolved handle was * cleared. A pending command may still record the invalidation when this * returns `false`. */ invalidateResolvedHandle(handleId: string): boolean; hasPendingCommands(): boolean; subscribeHandle(subscriber: (handleId: THandleId) => void): () => void; clearResolvedHandle(): THandleId | undefined; reset(): void; }; export type CreateFlowCommandHandleControllerOptions< THandleId extends string, TSettlement, > = { getHandleIdFromSettlement: (settlement: TSettlement) => THandleId | undefined; isStaleHandleError: (error: unknown) => boolean; /** Best-effort notification after the resolved handle has been cleared. */ onClearResolvedHandle?: (handleId: THandleId) => void; }; export function createFlowCommandHandleController< THandleId extends string, TSettlement, >({ getHandleIdFromSettlement, isStaleHandleError, onClearResolvedHandle, }: CreateFlowCommandHandleControllerOptions< THandleId, TSettlement >): FlowCommandHandleController { const handleClient: LazyHandleClient = createLazyHandleClient({ getHandleIdFromSettlement, }); const notifyResolvedHandleCleared = (handleId: THandleId): void => { try { onClearResolvedHandle?.(handleId); } catch { // Observers cannot take ownership of the completed handle transition. } }; const clearResolvedHandle = (): THandleId | undefined => { const previousHandleId = handleClient.getHandleId(); if (!previousHandleId) { return undefined; } handleClient.reset(); notifyResolvedHandleCleared(previousHandleId); return previousHandleId; }; const clearResolvedHandleForRetry = (): THandleId | undefined => { const previousHandleId = handleClient.clearHandleForRetry(); if (!previousHandleId) { return undefined; } notifyResolvedHandleCleared(previousHandleId); return previousHandleId; }; const reset = (): void => { const previousHandleId = handleClient.getHandleId(); handleClient.reset(); if (previousHandleId) { notifyResolvedHandleCleared(previousHandleId); } }; const runCommand: FlowCommandHandleController< THandleId, TSettlement >["runCommand"] = async (buildPayload, dispatchAndWait) => { return handleClient.runCommandWithRetry(buildPayload, dispatchAndWait, { shouldRetry: isStaleHandleError, onBeforeRetry: () => { clearResolvedHandleForRetry(); }, }); }; const runResolvedHandleCommand: FlowCommandHandleController< THandleId, TSettlement >["runResolvedHandleCommand"] = async (buildPayload, dispatchAndWait) => { try { await handleClient.runResolvedHandleCommand( buildPayload, dispatchAndWait, ); } catch (error) { if (isStaleHandleError(error)) { clearResolvedHandle(); } throw error; } }; const invalidateResolvedHandle = (handleId: string): boolean => { const invalidatedHandleId = handleClient.recordHandleInvalidation(handleId); if (!invalidatedHandleId) { return false; } notifyResolvedHandleCleared(invalidatedHandleId); return true; }; return { runCommand, runResolvedHandleCommand, getHandleId: () => handleClient.getHandleId(), hasHandleId: (handleId) => handleClient.getHandleId() === handleId, invalidateResolvedHandle, hasPendingCommands: () => handleClient.hasPendingCommands(), subscribeHandle: (subscriber) => handleClient.subscribeHandle(subscriber), clearResolvedHandle, reset, }; }