Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 8 additions & 2 deletions pkgs/edge-worker/tests/integration/_helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ export function createTestPlatformAdapter(sql: postgres.Sql): PlatformAdapter<Su
};
}

export function startWorker<TFlow extends AnyFlow>(
export async function startWorker<TFlow extends AnyFlow>(
sql: postgres.Sql,
flow: TFlow,
options: FlowWorkerConfig
Expand Down Expand Up @@ -84,7 +84,13 @@ export function startWorker<TFlow extends AnyFlow>(

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(),
});
Expand Down
14 changes: 7 additions & 7 deletions pkgs/edge-worker/tests/integration/flow/conditionalFlow.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ Deno.test(
}
);

const worker = startWorker(sql, ConditionalIfMetFlow, workerConfig);
const worker = await startWorker(sql, ConditionalIfMetFlow, workerConfig);

try {
await createFlowInDb(sql, ConditionalIfMetFlow);
Expand Down Expand Up @@ -118,7 +118,7 @@ Deno.test(
}
);

const worker = startWorker(sql, ConditionalIfUnmetFlow, workerConfig);
const worker = await startWorker(sql, ConditionalIfUnmetFlow, workerConfig);

try {
await createFlowInDb(sql, ConditionalIfUnmetFlow);
Expand Down Expand Up @@ -192,7 +192,7 @@ Deno.test(
}
);

const worker = startWorker(sql, ConditionalIfNotMetFlow, workerConfig);
const worker = await startWorker(sql, ConditionalIfNotMetFlow, workerConfig);

try {
await createFlowInDb(sql, ConditionalIfNotMetFlow);
Expand Down Expand Up @@ -260,7 +260,7 @@ Deno.test(
}
);

const worker = startWorker(sql, ConditionalIfNotUnmetFlow, workerConfig);
const worker = await startWorker(sql, ConditionalIfNotUnmetFlow, workerConfig);

try {
await createFlowInDb(sql, ConditionalIfNotUnmetFlow);
Expand Down Expand Up @@ -333,7 +333,7 @@ Deno.test(
}
);

const worker = startWorker(sql, SkipCascadeFlow, workerConfig);
const worker = await startWorker(sql, SkipCascadeFlow, workerConfig);

try {
await createFlowInDb(sql, SkipCascadeFlow);
Expand Down Expand Up @@ -434,7 +434,7 @@ Deno.test(
}
);

const worker = startWorker(sql, NonCascadeSkipFlow, workerConfig);
const worker = await startWorker(sql, NonCascadeSkipFlow, workerConfig);

try {
await createFlowInDb(sql, NonCascadeSkipFlow);
Expand Down Expand Up @@ -529,7 +529,7 @@ Deno.test(
}
);

const worker = startWorker(sql, FailOnUnmetFlow, workerConfig);
const worker = await startWorker(sql, FailOnUnmetFlow, workerConfig);

try {
await createFlowInDb(sql, FailOnUnmetFlow);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
8 changes: 4 additions & 4 deletions pkgs/edge-worker/tests/integration/flow/mapFlow.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -253,7 +253,7 @@ interface WorkerConfig {
[key: string]: unknown;
}

const startWorkers = (
const startWorkers = async (
sql: postgres.Sql,
flow: typeof LargeArrayMapFlow,
numWorkers: number,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
12 changes: 6 additions & 6 deletions pkgs/edge-worker/tests/integration/flow/whenExhausted.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ Deno.test(
}
);

const worker = startWorker(sql, FailOnErrorFlow, workerConfig);
const worker = await startWorker(sql, FailOnErrorFlow, workerConfig);

try {
await createFlowInDb(sql, FailOnErrorFlow);
Expand Down Expand Up @@ -125,7 +125,7 @@ Deno.test(
}
);

const worker = startWorker(sql, SkipOnErrorFlow, workerConfig);
const worker = await startWorker(sql, SkipOnErrorFlow, workerConfig);

try {
await createFlowInDb(sql, SkipOnErrorFlow);
Expand Down Expand Up @@ -194,7 +194,7 @@ Deno.test(
}
);

const worker = startWorker(sql, SkipCascadeOnErrorFlow, workerConfig);
const worker = await startWorker(sql, SkipCascadeOnErrorFlow, workerConfig);

try {
await createFlowInDb(sql, SkipCascadeOnErrorFlow);
Expand Down Expand Up @@ -267,7 +267,7 @@ Deno.test(
}
);

const worker = startWorker(sql, RetryThenSkipFlow, workerConfig);
const worker = await startWorker(sql, RetryThenSkipFlow, workerConfig);

try {
await createFlowInDb(sql, RetryThenSkipFlow);
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -396,7 +396,7 @@ Deno.test(
}
);

const worker = startWorker(sql, CombinedConditionsFlow, workerConfig);
const worker = await startWorker(sql, CombinedConditionsFlow, workerConfig);

try {
await createFlowInDb(sql, CombinedConditionsFlow);
Expand Down
6 changes: 3 additions & 3 deletions pkgs/edge-worker/tests/integration/flow_deprecation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
});
Expand Down Expand Up @@ -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,
});
Expand Down
Loading