Skip to content

Commit 7a1ea5a

Browse files
authored
Fix namespaced active run recovery (#2888)
Signed-off-by: Casey Gowrie <ctgowrie@gmail.com>
1 parent 0b956f6 commit 7a1ea5a

5 files changed

Lines changed: 123 additions & 4 deletions

File tree

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
---
2+
'@workflow/world': patch
3+
'@workflow/world-local': patch
4+
'@workflow/world-postgres': patch
5+
---
6+
7+
Use the active queue namespace when re-enqueuing workflow runs during world startup recovery.

‎packages/world-local/src/reenqueue.test.ts‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { rm } from 'node:fs/promises';
22
import os from 'node:os';
33
import path from 'node:path';
4+
import { WorkflowInvokePayloadSchema } from '@workflow/world';
45
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
56
import { createWorld } from './index.js';
67
import { createRun, updateRun } from './test-helpers.js';
@@ -14,10 +15,12 @@ describe('re-enqueue active runs on start', () => {
1415
let dataDir: string;
1516

1617
beforeEach(() => {
18+
vi.stubEnv('WORKFLOW_QUEUE_NAMESPACE', undefined);
1719
dataDir = path.join(os.tmpdir(), `wf-reenqueue-${Date.now()}`);
1820
});
1921

2022
afterEach(async () => {
23+
vi.unstubAllEnvs();
2124
await rm(dataDir, { recursive: true, force: true });
2225
});
2326

@@ -86,6 +89,49 @@ describe('re-enqueue active runs on start', () => {
8689
await world2.close();
8790
});
8891

92+
it('re-enqueues runs to the active queue namespace', async () => {
93+
vi.stubEnv('WORKFLOW_QUEUE_NAMESPACE', 'custom');
94+
95+
const world1 = createWorld({ dataDir });
96+
await world1.start();
97+
98+
const pendingRun = await createRun(world1, {
99+
deploymentId: 'dpl_1',
100+
workflowName: 'myWorkflow',
101+
input: new Uint8Array([1]),
102+
});
103+
104+
await world1.close();
105+
106+
const world2 = createWorld({ dataDir });
107+
const namespacedRunIds: string[] = [];
108+
const unnamespacedRunIds: string[] = [];
109+
const namespacedHandler = world2.createQueueHandler(
110+
'__custom_wkf_workflow_',
111+
async (message) => {
112+
const body = WorkflowInvokePayloadSchema.parse(message);
113+
namespacedRunIds.push(body.runId);
114+
}
115+
);
116+
world2.registerHandler('__custom_wkf_workflow_', namespacedHandler);
117+
// Capture an incorrectly reconstructed queue without allowing it to enter
118+
// the local queue's retry loop.
119+
world2.registerHandler('__wkf_workflow_', async (req) => {
120+
const body = await req.json();
121+
unnamespacedRunIds.push(body.runId);
122+
return Response.json({ ok: true });
123+
});
124+
125+
await world2.start();
126+
127+
await vi.waitFor(() => {
128+
expect(namespacedRunIds).toEqual([pendingRun.runId]);
129+
});
130+
expect(unnamespacedRunIds).toHaveLength(0);
131+
132+
await world2.close();
133+
});
134+
89135
it('does nothing when there are no active runs', async () => {
90136
// Create a world with only completed runs
91137
const world1 = createWorld({ dataDir });

‎packages/world-postgres/src/index.ts‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,12 @@ export function createWorld(
7070
}),
7171
async start() {
7272
await queue.start();
73-
await reenqueueActiveRuns(storage.runs, queue.queue, 'world-postgres');
73+
await reenqueueActiveRuns(
74+
storage.runs,
75+
queue.queue,
76+
'world-postgres',
77+
config.namespace
78+
);
7479
},
7580
async close() {
7681
await streamer.close();
Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
import { afterEach, describe, expect, it, vi } from 'vitest';
2+
import type { Storage } from './interfaces.js';
3+
import type { Queue } from './queue.js';
4+
import { reenqueueActiveRuns } from './recovery.js';
5+
6+
function createRuns(): Storage['runs'] {
7+
return {
8+
list: vi.fn(async ({ status }) => ({
9+
data:
10+
status === 'pending'
11+
? [
12+
{
13+
runId: 'wrun_AAA',
14+
workflowName: 'myWorkflow',
15+
status,
16+
},
17+
]
18+
: [],
19+
hasMore: false,
20+
cursor: null,
21+
})),
22+
} as unknown as Storage['runs'];
23+
}
24+
25+
describe('reenqueueActiveRuns', () => {
26+
afterEach(() => {
27+
vi.unstubAllEnvs();
28+
});
29+
30+
it('uses WORKFLOW_QUEUE_NAMESPACE for recovered runs', async () => {
31+
vi.stubEnv('WORKFLOW_QUEUE_NAMESPACE', 'custom');
32+
const enqueue = vi.fn<Queue['queue']>();
33+
34+
await reenqueueActiveRuns(createRuns(), enqueue, 'test');
35+
36+
expect(enqueue).toHaveBeenCalledWith('__custom_wkf_workflow_myWorkflow', {
37+
runId: 'wrun_AAA',
38+
});
39+
});
40+
41+
it('prefers an explicit namespace over WORKFLOW_QUEUE_NAMESPACE', async () => {
42+
vi.stubEnv('WORKFLOW_QUEUE_NAMESPACE', 'environment');
43+
const enqueue = vi.fn<Queue['queue']>();
44+
45+
await reenqueueActiveRuns(createRuns(), enqueue, 'test', 'explicit');
46+
47+
expect(enqueue).toHaveBeenCalledWith('__explicit_wkf_workflow_myWorkflow', {
48+
runId: 'wrun_AAA',
49+
});
50+
});
51+
});

‎packages/world/src/recovery.ts‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,10 @@
1-
import type { Queue } from './queue.js';
21
import type { Storage } from './interfaces.js';
32
import type { ValidQueueName } from './queue.js';
3+
import {
4+
getQueueTopicPrefix,
5+
type Queue,
6+
resolveQueueNamespace,
7+
} from './queue.js';
48

59
/**
610
* Re-enqueue all active (pending/running) workflow runs so they resume
@@ -10,12 +14,18 @@ import type { ValidQueueName } from './queue.js';
1014
* @param runs - Storage runs interface for listing active runs
1115
* @param enqueue - Queue's enqueue method
1216
* @param label - Log prefix for identifying the world implementation (e.g. "world-local")
17+
* @param namespace - Optional queue namespace. Defaults to WORKFLOW_QUEUE_NAMESPACE.
1318
*/
1419
export async function reenqueueActiveRuns(
1520
runs: Storage['runs'],
1621
enqueue: Queue['queue'],
17-
label: string
22+
label: string,
23+
namespace?: string
1824
): Promise<void> {
25+
const workflowQueuePrefix = getQueueTopicPrefix(
26+
'workflow',
27+
resolveQueueNamespace(namespace)
28+
);
1929
let reenqueued = 0;
2030
for (const status of ['pending', 'running'] as const) {
2131
let cursor: string | undefined;
@@ -28,7 +38,7 @@ export async function reenqueueActiveRuns(
2838
});
2939
for (const run of page.data) {
3040
try {
31-
const queueName: ValidQueueName = `__wkf_workflow_${run.workflowName}`;
41+
const queueName: ValidQueueName = `${workflowQueuePrefix}${run.workflowName}`;
3242
await enqueue(queueName, { runId: run.runId });
3343
reenqueued++;
3444
} catch (err) {

0 commit comments

Comments
 (0)