import { describe, expect, it } from "bun:test"; import { createFlowCommandHandleController } from "./flow-command-handle-controller"; describe("flow command handle controller", () => { it("serializes commands and publishes settled handles to later payloads", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); await controller.runCommand( (flowHandleId) => ({ kind: "prefetch", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_123" } }; }, ); await controller.runCommand( (flowHandleId) => ({ kind: "open", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_123" } }; }, ); expect(dispatched).toEqual([ { kind: "prefetch", flowHandleId: undefined }, { kind: "open", flowHandleId: "handle_123" }, ]); expect(controller.getHandleId()).toBe("handle_123"); expect(controller.hasHandleId("handle_123")).toBe(true); }); it("clears stale handles and retries once without the handle", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; const clearedHandles: string[] = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", onClearResolvedHandle: (flowHandleId) => { clearedHandles.push(flowHandleId); }, }); await controller.runCommand( (flowHandleId) => ({ kind: "prefetch", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_stale" } }; }, ); let didFailOnce = false; await controller.runCommand( (flowHandleId) => ({ kind: "open", flowHandleId }), async (payload) => { dispatched.push(payload); if (!didFailOnce) { didFailOnce = true; throw new Error("stale"); } return { result: { flowHandleId: "handle_recovered" } }; }, ); expect(dispatched).toEqual([ { kind: "prefetch", flowHandleId: undefined }, { kind: "open", flowHandleId: "handle_stale" }, { kind: "open", flowHandleId: undefined }, ]); expect(clearedHandles).toEqual(["handle_stale"]); expect(controller.getHandleId()).toBe("handle_recovered"); }); it("keeps stale-handle retry in the same serialized command slot", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", }); await controller.runCommand( (flowHandleId) => ({ kind: "prefetch", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_stale" } }; }, ); let didFailOnce = false; const staleOpen = controller.runCommand( (flowHandleId) => ({ kind: "open", flowHandleId }), async (payload) => { dispatched.push(payload); if (!didFailOnce) { didFailOnce = true; throw new Error("stale"); } return { result: { flowHandleId: "handle_recovered" } }; }, ); const queuedPrerender = controller.runCommand( (flowHandleId) => ({ kind: "prerender", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_second" } }; }, ); await staleOpen; await queuedPrerender; expect(dispatched).toEqual([ { kind: "prefetch", flowHandleId: undefined }, { kind: "open", flowHandleId: "handle_stale" }, { kind: "open", flowHandleId: undefined }, { kind: "prerender", flowHandleId: "handle_recovered" }, ]); expect(controller.getHandleId()).toBe("handle_second"); }); it("runs handle-resolved commands in the existing serialized queue", async () => { const started: string[] = []; let resolveAllocation: | ((settlement: { result?: { flowHandleId?: string } }) => void) | undefined; let resolveResolvedCommand: (() => void) | undefined; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); const allocation = controller.runCommand( () => ({ kind: "open" }), () => new Promise((resolve) => { started.push("open"); resolveAllocation = resolve; }), ); const resolvedCommand = controller.runResolvedHandleCommand( (flowHandleId) => ({ flowHandleId, kind: "close" }), (payload) => new Promise((resolve) => { started.push(`${payload.kind}:${payload.flowHandleId}`); resolveResolvedCommand = resolve; }), ); const nextCommand = controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), async (payload) => { started.push(`${payload.kind}:${payload.flowHandleId}`); return { result: { flowHandleId: "handle_next" } }; }, ); const nextResolvedCommand = controller.runResolvedHandleCommand( (flowHandleId) => ({ flowHandleId, kind: "close" }), async (payload) => { started.push(`${payload.kind}:${payload.flowHandleId}`); }, ); expect(started).toEqual(["open"]); resolveAllocation?.({ result: { flowHandleId: "handle_allocated" } }); await allocation; await Promise.resolve(); expect(started).toEqual(["open", "close:handle_allocated"]); resolveResolvedCommand?.(); await resolvedCommand; await nextCommand; await nextResolvedCommand; expect(started).toEqual([ "open", "close:handle_allocated", "open:handle_allocated", "close:handle_next", ]); expect(controller.getHandleId()).toBe("handle_next"); }); it("skips handle-resolved commands when their queue slot has no handle", async () => { const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); let didDispatch = false; await controller.runResolvedHandleCommand( (flowHandleId) => flowHandleId, async () => { didDispatch = true; }, ); expect(didDispatch).toBe(false); }); it("skips queued handle-resolved commands after reset", async () => { let resolveAllocation: | ((settlement: { result?: { flowHandleId?: string } }) => void) | undefined; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); let didDispatch = false; const allocation = controller.runCommand( () => ({ kind: "open" }), () => new Promise((resolve) => { resolveAllocation = resolve; }), ); const resolvedCommand = controller.runResolvedHandleCommand( (flowHandleId) => flowHandleId, async () => { didDispatch = true; }, ); controller.reset(); resolveAllocation?.({ result: { flowHandleId: "handle_stale" } }); await allocation; await resolvedCommand; expect(didDispatch).toBe(false); expect(controller.getHandleId()).toBeUndefined(); }); it("clears stale handles after a resolved-handle command rejects", async () => { const clearedHandles: string[] = []; const seenHandles: Array = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", onClearResolvedHandle: (flowHandleId) => { clearedHandles.push(flowHandleId); }, }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_stale" } }), ); await expect( controller.runResolvedHandleCommand( (flowHandleId) => ({ flowHandleId, kind: "close" }), async () => { throw new Error("stale"); }, ), ).rejects.toThrow("stale"); await controller.runCommand( (flowHandleId) => { seenHandles.push(flowHandleId); return { kind: "open" }; }, async () => ({ result: { flowHandleId: "handle_recovered" } }), ); expect(clearedHandles).toEqual(["handle_stale"]); expect(seenHandles).toEqual([undefined]); expect(controller.getHandleId()).toBe("handle_recovered"); }); it("preserves stale dispatch errors when clear observers throw", async () => { const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale dispatch", onClearResolvedHandle: () => { throw new Error("observer failed"); }, }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_stale" } }), ); await expect( controller.runResolvedHandleCommand( (flowHandleId) => ({ flowHandleId, kind: "close" }), async () => { throw new Error("stale dispatch"); }, ), ).rejects.toThrow("stale dispatch"); expect(controller.getHandleId()).toBeUndefined(); }); it("retries stale commands when clear observers throw", async () => { const seenHandles: Array = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", onClearResolvedHandle: () => { throw new Error("observer failed"); }, }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_stale" } }), ); let didFailOnce = false; await controller.runCommand( (flowHandleId) => { seenHandles.push(flowHandleId); return { flowHandleId, kind: "open" }; }, async () => { if (!didFailOnce) { didFailOnce = true; throw new Error("stale"); } return { result: { flowHandleId: "handle_recovered" } }; }, ); expect(seenHandles).toEqual(["handle_stale", undefined]); expect(controller.getHandleId()).toBe("handle_recovered"); }); it("does not retry non-stale failures", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); await expect( controller.runCommand( (flowHandleId) => ({ kind: "open", flowHandleId }), async (payload) => { dispatched.push(payload); throw new Error("boom"); }, ), ).rejects.toThrow("boom"); expect(dispatched).toEqual([{ kind: "open", flowHandleId: undefined }]); }); it("ignores handles from commands that settle after reset", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; let resolveSettlement: | ((settlement: { result?: { flowHandleId?: string } }) => void) | undefined; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); const pendingCommand = controller.runCommand( (flowHandleId) => ({ kind: "prefetch", flowHandleId }), (payload) => { dispatched.push(payload); return new Promise((resolve) => { resolveSettlement = resolve; }); }, ); controller.reset(); resolveSettlement?.({ result: { flowHandleId: "handle_from_old_epoch" } }); await pendingCommand; await controller.runCommand( (flowHandleId) => ({ kind: "open", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_current" } }; }, ); expect(dispatched).toEqual([ { kind: "prefetch", flowHandleId: undefined }, { kind: "open", flowHandleId: undefined }, ]); expect(controller.getHandleId()).toBe("handle_current"); }); it("does not retry stale failures from commands that fail after reset", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; let rejectSettlement: ((error: Error) => void) | undefined; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", }); await controller.runCommand( (flowHandleId) => ({ kind: "prefetch", flowHandleId }), async (payload) => { dispatched.push(payload); return { result: { flowHandleId: "handle_before_reset" } }; }, ); const pendingCommand = controller.runCommand( (flowHandleId) => ({ kind: "open", flowHandleId }), (payload) => { dispatched.push(payload); return new Promise((_resolve, reject) => { rejectSettlement = reject; }); }, ); controller.reset(); rejectSettlement?.(new Error("stale")); await expect(pendingCommand).rejects.toThrow("stale"); expect(dispatched).toEqual([ { kind: "prefetch", flowHandleId: undefined }, { kind: "open", flowHandleId: "handle_before_reset" }, ]); expect(controller.getHandleId()).toBeUndefined(); }); it("retries an active stale command after authoritative handle invalidation", async () => { const dispatched: Array<{ flowHandleId?: string; kind: string }> = []; const clearedHandles: string[] = []; let rejectSettlement: ((error: Error) => void) | undefined; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", onClearResolvedHandle: (flowHandleId) => { clearedHandles.push(flowHandleId); throw new Error("observer failure is contained"); }, }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_stale" } }), ); const pendingCommand = controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), (payload) => { dispatched.push(payload); if (payload.flowHandleId === undefined) { return Promise.resolve({ result: { flowHandleId: "handle_recovered" }, }); } return new Promise((_resolve, reject) => { rejectSettlement = reject; }); }, ); expect(controller.invalidateResolvedHandle("handle_stale")).toBe(true); rejectSettlement?.(new Error("stale")); await pendingCommand; expect(dispatched).toEqual([ { flowHandleId: "handle_stale", kind: "open" }, { flowHandleId: undefined, kind: "open" }, ]); expect(clearedHandles).toEqual(["handle_stale"]); expect(controller.getHandleId()).toBe("handle_recovered"); }); it("keeps invalidation authoritative before the first settlement", async () => { let resolveSettlement: | ((settlement: { result?: { flowHandleId?: string } }) => void) | undefined; const seenHandles: Array = []; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); const pendingCommand = controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), () => new Promise((resolve) => { resolveSettlement = resolve; }), ); expect(controller.invalidateResolvedHandle("handle_pending")).toBe(false); resolveSettlement?.({ result: { flowHandleId: "handle_pending" } }); await pendingCommand; await controller.runCommand( (flowHandleId) => { seenHandles.push(flowHandleId); return { flowHandleId, kind: "open" }; }, async () => ({ result: { flowHandleId: "handle_current" } }), ); expect(seenHandles).toEqual([undefined]); expect(controller.getHandleId()).toBe("handle_current"); }); it("preserves authoritative invalidation across stale-handle retry", async () => { const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", }); let didFailOnce = false; await controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), async () => { if (!didFailOnce) { didFailOnce = true; controller.invalidateResolvedHandle("handle_recovered"); throw new Error("stale"); } return { result: { flowHandleId: "handle_recovered" } }; }, ); expect(controller.getHandleId()).toBeUndefined(); }); it("preserves authoritative invalidation received during stale-handle retry", async () => { const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", }); let attempts = 0; let resolveRetry: | ((settlement: { result?: { flowHandleId?: string } }) => void) | undefined; const pendingCommand = controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), async () => { attempts += 1; if (attempts === 1) { throw new Error("stale"); } return new Promise((resolve) => { resolveRetry = resolve; }); }, ); await Promise.resolve(); await Promise.resolve(); expect(attempts).toBe(2); expect(controller.invalidateResolvedHandle("handle_recovered")).toBe(false); resolveRetry?.({ result: { flowHandleId: "handle_recovered" } }); await pendingCommand; expect(controller.getHandleId()).toBeUndefined(); }); it("does not treat stale-handle retry cleanup as authoritative invalidation", async () => { const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: (error) => error instanceof Error && error.message === "stale", }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_reused" } }), ); let didFailOnce = false; await controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), async () => { if (!didFailOnce) { didFailOnce = true; throw new Error("stale"); } return { result: { flowHandleId: "handle_reused" } }; }, ); expect(controller.getHandleId()).toBe("handle_reused"); }); it("does not restore an invalidated handle from a late successful settlement", async () => { let resolveSettlement: | ((settlement: { result?: { flowHandleId?: string } }) => void) | undefined; const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_invalidated" } }), ); const pendingCommand = controller.runCommand( (flowHandleId) => ({ flowHandleId, kind: "open" }), () => new Promise((resolve) => { resolveSettlement = resolve; }), ); expect(controller.invalidateResolvedHandle("handle_invalidated")).toBe( true, ); resolveSettlement?.({ result: { flowHandleId: "handle_invalidated" } }); await pendingCommand; expect(controller.getHandleId()).toBeUndefined(); }); it("ignores delayed invalidation for a previously replaced handle", async () => { const controller = createFlowCommandHandleController< string, { result?: { flowHandleId?: string } } >({ getHandleIdFromSettlement: (settlement) => settlement.result?.flowHandleId, isStaleHandleError: () => false, }); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_old" } }), ); await controller.runCommand( () => ({ kind: "open" }), async () => ({ result: { flowHandleId: "handle_current" } }), ); expect(controller.invalidateResolvedHandle("handle_old")).toBe(false); expect(controller.getHandleId()).toBe("handle_current"); }); });