From cae4772fb9748c53f141ce0c4201e9318d1fa8dc Mon Sep 17 00:00:00 2001 From: Agent Date: Sun, 13 Sep 2026 06:50:42 +0000 Subject: [PATCH] fix(edge-worker): await worker startup in integration tests startWorker() fired worker.startOnlyOnce() fire-and-forget, so the startup transaction (ensure_flow_compiled: create_flow/add_step, or destructive local recompile via delete_flow_and_data) raced the test body and could commit across test boundaries, interleaving with the next test's reset_db. A stray commit then left a listed PGMQ queue with no owning flow row, and the next create_flow correctly rejected it: cannot create flow "X": queue "X" is already in use by another owner. Fix: make startWorker() async and await startup; await all 31 call sites and performanceMapFlow's startWorkers. Startup errors now fail at the call site instead of surfacing via a later worker.stop(). Product code, SQL, and migrations are untouched; the ownership guard stays as approved in #650. mapFlow.test.ts focused runs pass twice consecutively; full integration suite 63 passed, 0 failed. Refs #650. CI: edge-worker-integration on PR #679. --- pkgs/edge-worker/tests/integration/_helpers.ts | 10 ++++++++-- .../tests/integration/flow/conditionalFlow.test.ts | 14 +++++++------- .../integration/flow/contextResourcesFlow.test.ts | 2 +- .../tests/integration/flow/mapFlow.test.ts | 8 ++++---- .../tests/integration/flow/minimalFlow.test.ts | 2 +- .../integration/flow/performanceMapFlow.test.ts | 10 +++++----- .../tests/integration/flow/startDelayFlow.test.ts | 8 ++++---- .../tests/integration/flow/whenExhausted.test.ts | 12 ++++++------ .../tests/integration/flow_deprecation.test.ts | 6 +++--- 9 files changed, 39 insertions(+), 33 deletions(-) diff --git a/pkgs/edge-worker/tests/integration/_helpers.ts b/pkgs/edge-worker/tests/integration/_helpers.ts index 5dd4fbe42..aa94fe62f 100644 --- a/pkgs/edge-worker/tests/integration/_helpers.ts +++ b/pkgs/edge-worker/tests/integration/_helpers.ts @@ -49,7 +49,7 @@ export function createTestPlatformAdapter(sql: postgres.Sql): PlatformAdapter( +export async function startWorker( sql: postgres.Sql, flow: TFlow, options: FlowWorkerConfig @@ -84,7 +84,13 @@ export function startWorker( const worker = createFlowWorker(flow, mergedOptions, () => consoleLogger, createTestPlatformAdapter(sql)); - worker.startOnlyOnce({ + // Await startup: ensure_flow_compiled may create the flow, its queue, or + // destructively recompile (dropping queues). Fire-and-forget startup lets + // that transaction race the test body's own create_flow/add_step calls and, + // under load, commit after the test finished — leaving a listed queue with + // no owning flow row, which the create_flow ownership guard then rejects in + // a later test (#650). + await worker.startOnlyOnce({ edgeFunctionName: 'test_flow', workerId: crypto.randomUUID(), }); diff --git a/pkgs/edge-worker/tests/integration/flow/conditionalFlow.test.ts b/pkgs/edge-worker/tests/integration/flow/conditionalFlow.test.ts index 0a1ac5b41..3c16fc370 100644 --- a/pkgs/edge-worker/tests/integration/flow/conditionalFlow.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/conditionalFlow.test.ts @@ -49,7 +49,7 @@ Deno.test( } ); - const worker = startWorker(sql, ConditionalIfMetFlow, workerConfig); + const worker = await startWorker(sql, ConditionalIfMetFlow, workerConfig); try { await createFlowInDb(sql, ConditionalIfMetFlow); @@ -118,7 +118,7 @@ Deno.test( } ); - const worker = startWorker(sql, ConditionalIfUnmetFlow, workerConfig); + const worker = await startWorker(sql, ConditionalIfUnmetFlow, workerConfig); try { await createFlowInDb(sql, ConditionalIfUnmetFlow); @@ -192,7 +192,7 @@ Deno.test( } ); - const worker = startWorker(sql, ConditionalIfNotMetFlow, workerConfig); + const worker = await startWorker(sql, ConditionalIfNotMetFlow, workerConfig); try { await createFlowInDb(sql, ConditionalIfNotMetFlow); @@ -260,7 +260,7 @@ Deno.test( } ); - const worker = startWorker(sql, ConditionalIfNotUnmetFlow, workerConfig); + const worker = await startWorker(sql, ConditionalIfNotUnmetFlow, workerConfig); try { await createFlowInDb(sql, ConditionalIfNotUnmetFlow); @@ -333,7 +333,7 @@ Deno.test( } ); - const worker = startWorker(sql, SkipCascadeFlow, workerConfig); + const worker = await startWorker(sql, SkipCascadeFlow, workerConfig); try { await createFlowInDb(sql, SkipCascadeFlow); @@ -434,7 +434,7 @@ Deno.test( } ); - const worker = startWorker(sql, NonCascadeSkipFlow, workerConfig); + const worker = await startWorker(sql, NonCascadeSkipFlow, workerConfig); try { await createFlowInDb(sql, NonCascadeSkipFlow); @@ -529,7 +529,7 @@ Deno.test( } ); - const worker = startWorker(sql, FailOnUnmetFlow, workerConfig); + const worker = await startWorker(sql, FailOnUnmetFlow, workerConfig); try { await createFlowInDb(sql, FailOnUnmetFlow); diff --git a/pkgs/edge-worker/tests/integration/flow/contextResourcesFlow.test.ts b/pkgs/edge-worker/tests/integration/flow/contextResourcesFlow.test.ts index 334517f8c..1c5fa6a84 100644 --- a/pkgs/edge-worker/tests/integration/flow/contextResourcesFlow.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/contextResourcesFlow.test.ts @@ -55,7 +55,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, ContextResourcesFlow, { + const worker = await startWorker(sql, ContextResourcesFlow, { maxConcurrent: 1, batchSize: 10, maxPollSeconds: 1, diff --git a/pkgs/edge-worker/tests/integration/flow/mapFlow.test.ts b/pkgs/edge-worker/tests/integration/flow/mapFlow.test.ts index 85a0c961e..f5ff88e84 100644 --- a/pkgs/edge-worker/tests/integration/flow/mapFlow.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/mapFlow.test.ts @@ -66,7 +66,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, RootMapFlow, { + const worker = await startWorker(sql, RootMapFlow, { maxConcurrent: 3, batchSize: 10, maxPollSeconds: 1, @@ -114,7 +114,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, DependentMapFlow, { + const worker = await startWorker(sql, DependentMapFlow, { maxConcurrent: 3, batchSize: 10, maxPollSeconds: 1, @@ -176,7 +176,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, EmptyArrayMapFlow, { + const worker = await startWorker(sql, EmptyArrayMapFlow, { maxConcurrent: 3, batchSize: 10, maxPollSeconds: 1, @@ -226,7 +226,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, SlotRefillFlow, { + const worker = await startWorker(sql, SlotRefillFlow, { maxConcurrent: 2, batchSize: 2, maxPollSeconds: 1, diff --git a/pkgs/edge-worker/tests/integration/flow/minimalFlow.test.ts b/pkgs/edge-worker/tests/integration/flow/minimalFlow.test.ts index 5d6a40e0a..28e8ca956 100644 --- a/pkgs/edge-worker/tests/integration/flow/minimalFlow.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/minimalFlow.test.ts @@ -34,7 +34,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, MinimalFlow, { + const worker = await startWorker(sql, MinimalFlow, { maxConcurrent: 1, batchSize: 10, maxPollSeconds: 1, diff --git a/pkgs/edge-worker/tests/integration/flow/performanceMapFlow.test.ts b/pkgs/edge-worker/tests/integration/flow/performanceMapFlow.test.ts index d46f3c9e1..6bf1a10cd 100644 --- a/pkgs/edge-worker/tests/integration/flow/performanceMapFlow.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/performanceMapFlow.test.ts @@ -253,7 +253,7 @@ interface WorkerConfig { [key: string]: unknown; } -const startWorkers = ( +const startWorkers = async ( sql: postgres.Sql, flow: typeof LargeArrayMapFlow, numWorkers: number, @@ -262,7 +262,7 @@ const startWorkers = ( const workers = []; for (let i = 0; i < numWorkers; i++) { - const worker = startWorker(sql, flow, { + const worker = await startWorker(sql, flow, { ...config, // Add slight variation to worker configs to simulate real-world maxConcurrent: config.maxConcurrent + (i % 2) * 5, @@ -342,7 +342,7 @@ Deno.test( const NUM_WORKERS = 4; const ELEMENTS_PER_FLOW = 100; - const workers = startWorkers(sql, LargeArrayMapFlow, NUM_WORKERS, { + const workers = await startWorkers(sql, LargeArrayMapFlow, NUM_WORKERS, { maxConcurrent: 25, batchSize: 10, maxPollSeconds: 1, @@ -498,7 +498,7 @@ Deno.test( console.log = () => {}; console.debug = () => {}; - const worker = startWorker(sql, LargeArrayMapFlow, { + const worker = await startWorker(sql, LargeArrayMapFlow, { maxConcurrent: 50, // Higher concurrency for performance test batchSize: 25, // Larger batch size for better throughput maxPollSeconds: 1, @@ -725,7 +725,7 @@ Deno.test( console.log = () => {}; console.debug = () => {}; - const worker = startWorker(sql, LargeArrayMapFlow, { + const worker = await startWorker(sql, LargeArrayMapFlow, { maxConcurrent: 100, // Even higher concurrency for stress test batchSize: 50, maxPollSeconds: 2, diff --git a/pkgs/edge-worker/tests/integration/flow/startDelayFlow.test.ts b/pkgs/edge-worker/tests/integration/flow/startDelayFlow.test.ts index 03eae8b82..3180da35f 100644 --- a/pkgs/edge-worker/tests/integration/flow/startDelayFlow.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/startDelayFlow.test.ts @@ -72,7 +72,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, RootStepDelayFlow, { + const worker = await startWorker(sql, RootStepDelayFlow, { maxConcurrent: 1, batchSize: 10, maxPollSeconds: 1, @@ -127,7 +127,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, NormalStepDelayFlow, { + const worker = await startWorker(sql, NormalStepDelayFlow, { maxConcurrent: 2, batchSize: 10, maxPollSeconds: 1, @@ -187,7 +187,7 @@ Deno.test( withPgNoTransaction(async (sql) => { await sql`select pgflow_tests.reset_db();`; - const worker = startWorker(sql, CascadedDelayFlow, { + const worker = await startWorker(sql, CascadedDelayFlow, { maxConcurrent: 3, batchSize: 10, maxPollSeconds: 1, @@ -293,7 +293,7 @@ Deno.test( return `Success on attempt ${attemptCount}`; }); - const worker = startWorker(sql, FlowWithRetry, { + const worker = await startWorker(sql, FlowWithRetry, { maxConcurrent: 1, batchSize: 10, maxPollSeconds: 1, diff --git a/pkgs/edge-worker/tests/integration/flow/whenExhausted.test.ts b/pkgs/edge-worker/tests/integration/flow/whenExhausted.test.ts index 59b22f45a..e6f8645a3 100644 --- a/pkgs/edge-worker/tests/integration/flow/whenExhausted.test.ts +++ b/pkgs/edge-worker/tests/integration/flow/whenExhausted.test.ts @@ -53,7 +53,7 @@ Deno.test( } ); - const worker = startWorker(sql, FailOnErrorFlow, workerConfig); + const worker = await startWorker(sql, FailOnErrorFlow, workerConfig); try { await createFlowInDb(sql, FailOnErrorFlow); @@ -125,7 +125,7 @@ Deno.test( } ); - const worker = startWorker(sql, SkipOnErrorFlow, workerConfig); + const worker = await startWorker(sql, SkipOnErrorFlow, workerConfig); try { await createFlowInDb(sql, SkipOnErrorFlow); @@ -194,7 +194,7 @@ Deno.test( } ); - const worker = startWorker(sql, SkipCascadeOnErrorFlow, workerConfig); + const worker = await startWorker(sql, SkipCascadeOnErrorFlow, workerConfig); try { await createFlowInDb(sql, SkipCascadeOnErrorFlow); @@ -267,7 +267,7 @@ Deno.test( } ); - const worker = startWorker(sql, RetryThenSkipFlow, workerConfig); + const worker = await startWorker(sql, RetryThenSkipFlow, workerConfig); try { await createFlowInDb(sql, RetryThenSkipFlow); @@ -325,7 +325,7 @@ Deno.test( return { result: deps.maybe_skip ?? null }; }); - const worker = startWorker(sql, SuccessfulHandlerFlow, workerConfig); + const worker = await startWorker(sql, SuccessfulHandlerFlow, workerConfig); try { await createFlowInDb(sql, SuccessfulHandlerFlow); @@ -396,7 +396,7 @@ Deno.test( } ); - const worker = startWorker(sql, CombinedConditionsFlow, workerConfig); + const worker = await startWorker(sql, CombinedConditionsFlow, workerConfig); try { await createFlowInDb(sql, CombinedConditionsFlow); diff --git a/pkgs/edge-worker/tests/integration/flow_deprecation.test.ts b/pkgs/edge-worker/tests/integration/flow_deprecation.test.ts index 1a6cd44a8..0131ff91e 100644 --- a/pkgs/edge-worker/tests/integration/flow_deprecation.test.ts +++ b/pkgs/edge-worker/tests/integration/flow_deprecation.test.ts @@ -40,7 +40,7 @@ Deno.test( await sql`select pgflow.add_step(${flowSlug}::text, 'step2'::text, ARRAY['step1']::text[], null, null, null, null)`; // Start the flow worker - const worker = startWorker(sql, flow, { + const worker = await startWorker(sql, flow, { maxPollSeconds: 1, pollIntervalMs: 50, }); @@ -160,12 +160,12 @@ Deno.test( await sql`select pgflow.add_step(${flowSlug}::text, 'step2'::text, ARRAY['step1']::text[], null, null, null, null)`; // Start two workers for the same flow - const worker1 = startWorker(sql, flow1, { + const worker1 = await startWorker(sql, flow1, { maxPollSeconds: 1, pollIntervalMs: 50, }); - const worker2 = startWorker(sql, flow2, { + const worker2 = await startWorker(sql, flow2, { maxPollSeconds: 1, pollIntervalMs: 50, });