Skip to content

Commit cb77725

Browse files
[core] Derive correlation ids from per-kind sequences (opt-in) (#3301)
1 parent f8f6e17 commit cb77725

32 files changed

Lines changed: 822 additions & 30 deletions
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
'@workflow/core': patch
3+
'workflow': patch
4+
---
5+
6+
Add env option to split correlation ID derivation into per-entity-type sequential ULIDs, instead of sharing one derivation source
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
'@workflow/core': patch
3+
'workflow': patch
4+
---
5+
6+
Fix stream IDs minted while serializing a step's arguments latching the host wall clock into a run's ID sequence, which gave every entity created afterwards a different correlation ID on each replay

‎docs/content/docs/v5/configuration/runtime-tuning.mdx‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,17 @@ For example, a workflow can run a 10-minute inline step even with `WORKFLOW_REPL
7575
- Delay before a re-invocation caused by a rejected event creation.
7676
- Unlike an in-process restart, which re-reads immediately, a re-invocation only happens once the in-process budget failed to catch up — so the delay gives the other writers a moment to quiesce.
7777

78+
### `WORKFLOW_PER_KIND_CORRELATION_IDS`
79+
80+
- Default: disabled
81+
- Experimental. Gives each kind of entity a workflow creates — steps, waits, hooks, attribute writes, abort controllers, stream IDs — its own sequence of correlation IDs.
82+
- With one sequence shared by every kind, an ID is an ordinal over the whole run, so a single extra draw of any kind shifts every ID after it. Two concurrent replays of the same run that disagree about one `sleep()` then assign different IDs to every step that follows, and each writes events the other can neither match nor consume, which fails the run with `CORRUPTED_EVENT_LOG`. Per-kind sequences confine that to the kind that actually differs.
83+
- IDs remain ordered within a kind, so hooks created by your workflow are still listed in creation order. A hook the runtime creates for you, such as the one backing an abort controller, draws from its own kind and so is listed at an arbitrary position relative to your hooks rather than at its creation position.
84+
- A run must replay under the scheme that minted its IDs. A replay that switches schemes mid-run assigns IDs its own earlier events do not carry, so it can consume none of them and the run fails.
85+
- On Vercel, a run keeps replaying on the deployment it started on, so it only ever sees the value baked into that deployment. Changing the setting affects new runs only.
86+
- Elsewhere — `@workflow/world-postgres`, `@workflow/world-local`, any self-hosted process — nothing pins a run to the code that started it. Turn the setting on during a quiet window with no runs in flight, and roll the new value out to your whole fleet at once: a rolling deploy that leaves both values live replays one run under two schemes concurrently, which is the failure the setting exists to reduce.
87+
- Set `1` to enable.
88+
7889
## Inline execution
7990

8091
### `WORKFLOW_V2_TIMEOUT_MS`

‎packages/core/src/abort-consistency.test.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,10 @@ import type { Event, WorkflowRun } from '@workflow/world';
1111
import * as nanoid from 'nanoid';
1212
import { monotonicFactory } from 'ulid';
1313
import { describe, expect, it, vi } from 'vitest';
14+
import {
15+
createCorrelationIdGenerator,
16+
isPerKindCorrelationIdsEnabled,
17+
} from './correlation-id.js';
1418
import { EventsConsumer } from './events-consumer.js';
1519
import type { WorkflowSuspension } from './global.js';
1620
import type { WorkflowOrchestratorContext } from './private.js';
@@ -44,7 +48,12 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext {
4448
getPromiseQueue: () => Promise.resolve(),
4549
}),
4650
invocationsQueue: new Map(),
47-
generateUlid: () => ulid(workflowStartedAt),
51+
generateCorrelationId: createCorrelationIdGenerator({
52+
seed: 'test',
53+
fixedTimestamp: workflowStartedAt,
54+
positional: () => ulid(workflowStartedAt),
55+
perKind: isPerKindCorrelationIdsEnabled(),
56+
}),
4857
generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) =>
4958
new Uint8Array(size).map(() => 256 * context.globalThis.Math.random())
5059
),

‎packages/core/src/abort-controller.test.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,10 @@ import type { Event } from '@workflow/world';
1212
import * as nanoid from 'nanoid';
1313
import { monotonicFactory } from 'ulid';
1414
import { describe, expect, it, vi } from 'vitest';
15+
import {
16+
createCorrelationIdGenerator,
17+
isPerKindCorrelationIdsEnabled,
18+
} from './correlation-id.js';
1519
import { DEFERRED_CHECK_DELAY_MS, EventsConsumer } from './events-consumer.js';
1620
import type { WorkflowOrchestratorContext } from './private.js';
1721
import { ReplayPayloadCache } from './replay-payload-cache.js';
@@ -42,7 +46,12 @@ function setupWorkflowContext(
4246
getPromiseQueue: () => ctx.promiseQueue,
4347
}),
4448
invocationsQueue: new Map(),
45-
generateUlid: () => ulid(workflowStartedAt),
49+
generateCorrelationId: createCorrelationIdGenerator({
50+
seed: 'test',
51+
fixedTimestamp: workflowStartedAt,
52+
positional: () => ulid(workflowStartedAt),
53+
perKind: isPerKindCorrelationIdsEnabled(),
54+
}),
4655
generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) =>
4756
new Uint8Array(size).map(() => 256 * context.globalThis.Math.random())
4857
),

‎packages/core/src/abort-replay-ordering.test.ts‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,10 @@ import type { Event } from '@workflow/world';
2727
import * as nanoid from 'nanoid';
2828
import { monotonicFactory } from 'ulid';
2929
import { describe, expect, it, vi } from 'vitest';
30+
import {
31+
createCorrelationIdGenerator,
32+
isPerKindCorrelationIdsEnabled,
33+
} from './correlation-id.js';
3034
import { EventsConsumer } from './events-consumer.js';
3135
import {
3236
scheduleWhenIdle,
@@ -77,7 +81,12 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext {
7781
getPromiseQueue: () => ctx.promiseQueue,
7882
}),
7983
invocationsQueue: new Map(),
80-
generateUlid: () => ulid(workflowStartedAt),
84+
generateCorrelationId: createCorrelationIdGenerator({
85+
seed: 'test',
86+
fixedTimestamp: workflowStartedAt,
87+
positional: () => ulid(workflowStartedAt),
88+
perKind: isPerKindCorrelationIdsEnabled(),
89+
}),
8190
generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) =>
8291
new Uint8Array(size).map(() => 256 * context.globalThis.Math.random())
8392
),

‎packages/core/src/async-deserialization-ordering.test.ts‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import * as nanoid from 'nanoid';
44
import { monotonicFactory } from 'ulid';
55
import { afterEach, beforeAll, describe, expect, it, vi } from 'vitest';
66
import { registerSerializationClass } from './class-serialization.js';
7+
import { createCorrelationIdGenerator } from './correlation-id.js';
78
import { EventsConsumer } from './events-consumer.js';
89
import type { WorkflowOrchestratorContext } from './private.js';
910
import { ReplayPayloadCache } from './replay-payload-cache.js';
@@ -57,7 +58,14 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext {
5758
getPromiseQueue: () => Promise.resolve(),
5859
}),
5960
invocationsQueue: new Map(),
60-
generateUlid: () => ulid(workflowStartedAt),
61+
generateCorrelationId: createCorrelationIdGenerator({
62+
seed: 'test',
63+
fixedTimestamp: workflowStartedAt,
64+
positional: () => ulid(workflowStartedAt),
65+
// The event logs in this file hardcode correlation ids the run-wide
66+
// shared sequence minted, so replay only matches under that scheme.
67+
perKind: false,
68+
}),
6169
generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) =>
6270
new Uint8Array(size).map(() => 256 * context.globalThis.Math.random())
6371
),
Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
import type { Event } from '@workflow/world';
2+
import * as nanoid from 'nanoid';
3+
import { monotonicFactory } from 'ulid';
4+
import { describe, expect, it, vi } from 'vitest';
5+
import { createCorrelationIdGenerator } from './correlation-id.js';
6+
import { EventsConsumer } from './events-consumer.js';
7+
import type { WorkflowOrchestratorContext } from './private.js';
8+
import { ReplayPayloadCache } from './replay-payload-cache.js';
9+
import { dehydrateStepReturnValue } from './serialization.js';
10+
import { createUseStep } from './step.js';
11+
import { createContext } from './vm/index.js';
12+
import { createCreateHook } from './workflow/hook.js';
13+
import { createSleep } from './workflow/sleep.js';
14+
15+
/**
16+
* Correlation-id stability seen through the primitives that actually mint ids,
17+
* rather than through the generator alone: that a step's id survives another
18+
* kind of entity being created alongside it, and that a replay consumes an event
19+
* log carrying the ids a same-seeded replay derives.
20+
*
21+
* The rest of the replay suites author their event logs with literal correlation
22+
* ids from the shared sequence and pin themselves to it. These fixtures derive
23+
* their ids instead, so they hold under either scheme.
24+
*/
25+
26+
const SEED = 'test';
27+
const FIXED_TIMESTAMP = 1753481739458;
28+
29+
function setupWorkflowContext(
30+
events: Event[],
31+
perKind: boolean
32+
): WorkflowOrchestratorContext {
33+
const context = createContext({
34+
seed: SEED,
35+
fixedTimestamp: FIXED_TIMESTAMP,
36+
});
37+
const ulid = monotonicFactory(() => context.globalThis.Math.random());
38+
return {
39+
runId: 'wrun_test',
40+
encryptionKey: undefined,
41+
replayPayloadCache: new ReplayPayloadCache(undefined),
42+
globalThis: context.globalThis,
43+
eventsConsumer: new EventsConsumer(events, {
44+
onUnconsumedEvent: () => {},
45+
getPromiseQueue: () => Promise.resolve(),
46+
}),
47+
invocationsQueue: new Map(),
48+
generateCorrelationId: createCorrelationIdGenerator({
49+
seed: SEED,
50+
fixedTimestamp: FIXED_TIMESTAMP,
51+
positional: () => ulid(FIXED_TIMESTAMP),
52+
perKind,
53+
}),
54+
generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) =>
55+
new Uint8Array(size).map(() => 256 * context.globalThis.Math.random())
56+
),
57+
onWorkflowError: vi.fn(),
58+
promiseQueue: Promise.resolve(),
59+
pendingDeliveries: 0,
60+
pendingDeliveryBarriers: new Map(),
61+
};
62+
}
63+
64+
/**
65+
* The id the next step of a replay would claim. Nothing in the log resolves the
66+
* step, so the returned promise stays pending by design: the queue item is what
67+
* we are after.
68+
*/
69+
function probeStepId(
70+
perKind: boolean,
71+
before?: (ctx: WorkflowOrchestratorContext) => void
72+
): string {
73+
const ctx = setupWorkflowContext([], perKind);
74+
before?.(ctx);
75+
void createUseStep(ctx)('add')(1, 2).catch(() => {});
76+
const item = [...ctx.invocationsQueue.values()].find(
77+
(entry) => entry.type === 'step'
78+
);
79+
if (!item) {
80+
throw new Error('expected a step invocation');
81+
}
82+
return item.correlationId;
83+
}
84+
85+
function createHookAndSleep(ctx: WorkflowOrchestratorContext): void {
86+
createCreateHook(ctx)();
87+
void createSleep(ctx)('1h').catch(() => {});
88+
}
89+
90+
describe('correlation ids through the replay primitives', () => {
91+
it('keeps a step id when a hook and a sleep are created before it', () => {
92+
expect(probeStepId(true, createHookAndSleep)).toBe(probeStepId(true));
93+
});
94+
95+
it('renumbers that step under one sequence shared by every kind', () => {
96+
// The failure this PR removes, and the reason the assertion above is worth
97+
// making: with a shared sequence the hook and the sleep consume the two
98+
// ordinals the step would otherwise have drawn from.
99+
expect(probeStepId(false, createHookAndSleep)).not.toBe(probeStepId(false));
100+
});
101+
102+
it('consumes a step_completed authored with the derived id', async () => {
103+
const correlationId = probeStepId(true, createHookAndSleep);
104+
const ctx = setupWorkflowContext(
105+
[
106+
{
107+
eventId: 'evnt_0',
108+
runId: 'wrun_test',
109+
eventType: 'step_completed',
110+
correlationId,
111+
eventData: {
112+
stepName: 'add',
113+
result: await dehydrateStepReturnValue(3, 'wrun_test', undefined),
114+
},
115+
createdAt: new Date(),
116+
},
117+
],
118+
true
119+
);
120+
createHookAndSleep(ctx);
121+
await expect(createUseStep(ctx)('add')(1, 2)).resolves.toBe(3);
122+
expect(ctx.onWorkflowError).not.toHaveBeenCalled();
123+
});
124+
});

0 commit comments

Comments
 (0)