Skip to content

Commit ad71b58

Browse files
authored
Report corrupted event logs distinctly (#2046)
1 parent 753e39e commit ad71b58

21 files changed

Lines changed: 225 additions & 61 deletions
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
---
2+
"@workflow/core": patch
3+
"@workflow/errors": patch
4+
"@workflow/world": patch
5+
---
6+
7+
Report corrupted event logs with a distinct `CorruptedEventLogError` type and `CORRUPTED_EVENT_LOG` run error code.

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
* (for real-time step propagation).
77
*/
88

9-
import { WorkflowRuntimeError } from '@workflow/errors';
9+
import { CorruptedEventLogError } from '@workflow/errors';
1010
import { withResolvers } from '@workflow/utils';
1111
import type { Event } from '@workflow/world';
1212
import * as nanoid from 'nanoid';
@@ -118,7 +118,7 @@ describe('AbortController in workflow VM', () => {
118118
expect(controller.signal.aborted).toBe(true);
119119
});
120120

121-
it('reports a WorkflowRuntimeError when abort hook_received token mismatches the controller', async () => {
121+
it('reports a CorruptedEventLogError when abort hook_received token mismatches the controller', async () => {
122122
ctx = setupWorkflowContext([]);
123123
const ProbeAbortController = createCreateAbortController(ctx);
124124
new ProbeAbortController();
@@ -157,7 +157,7 @@ describe('AbortController in workflow VM', () => {
157157
new AbortController();
158158

159159
const workflowError = await errorReceived.promise;
160-
expect(workflowError).toBeInstanceOf(WorkflowRuntimeError);
160+
expect(workflowError).toBeInstanceOf(CorruptedEventLogError);
161161
expect(workflowError?.message).toContain('hook_received');
162162
expect(workflowError?.message).toContain('wrong-token');
163163
expect(workflowError?.message).toContain(probeHookItem.token);

‎packages/core/src/classify-error.test.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import {
2+
CorruptedEventLogError,
23
HookConflictError,
34
RUN_ERROR_CODES,
45
WorkflowNotRegisteredError,
@@ -9,6 +10,12 @@ import { describe, expect, it } from 'vitest';
910
import { classifyRunError } from './classify-error.js';
1011

1112
describe('classifyRunError', () => {
13+
it('classifies CorruptedEventLogError as CORRUPTED_EVENT_LOG', () => {
14+
expect(
15+
classifyRunError(new CorruptedEventLogError('corrupted event log'))
16+
).toBe(RUN_ERROR_CODES.CORRUPTED_EVENT_LOG);
17+
});
18+
1219
it('classifies WorkflowRuntimeError as RUNTIME_ERROR', () => {
1320
expect(
1421
classifyRunError(new WorkflowRuntimeError('corrupted event log'))

‎packages/core/src/classify-error.ts‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import {
2+
CorruptedEventLogError,
23
RUN_ERROR_CODES,
34
type RunErrorCode,
45
StepNotRegisteredError,
@@ -7,7 +8,7 @@ import {
78
} from '@workflow/errors';
89

910
/**
10-
* Set of error names that should classify as `RUNTIME_ERROR`. Each
11+
* Set of error names that should classify as generic `RUNTIME_ERROR`. Each
1112
* `*.is()` static does a name-based duck check, so subclassing alone is
1213
* not enough — we have to enumerate every concrete subclass we want to
1314
* recognize. Keep in sync with the `WorkflowRuntimeError` class hierarchy
@@ -25,8 +26,8 @@ const RUNTIME_ERROR_CHECKS = [
2526
* After the structural separation of infrastructure vs user code error
2627
* handling, the only errors that reach the `run_failed` try/catch are:
2728
* - User code errors (throws from workflow functions, propagated step failures)
28-
* - WorkflowRuntimeError and subclasses (corrupted event log, missing
29-
* timestamps, workflow/step not registered, etc.)
29+
* - WorkflowRuntimeError and subclasses (missing timestamps, workflow/step
30+
* not registered, corrupted event log, etc.)
3031
*
3132
* Uses each subclass's `.is()` static (a name-based duck check) instead of
3233
* a single `instanceof` check because workflows execute in a separate
@@ -36,6 +37,10 @@ const RUNTIME_ERROR_CHECKS = [
3637
* errors as user errors.
3738
*/
3839
export function classifyRunError(err: unknown): RunErrorCode {
40+
if (CorruptedEventLogError.is(err)) {
41+
return RUN_ERROR_CODES.CORRUPTED_EVENT_LOG;
42+
}
43+
3944
for (const isMatch of RUNTIME_ERROR_CHECKS) {
4045
if (isMatch(err)) {
4146
return RUN_ERROR_CODES.RUNTIME_ERROR;

‎packages/core/src/describe-error.test.ts‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import {
2+
CorruptedEventLogError,
23
RUN_ERROR_CODES,
34
SerializationError,
45
StepNotRegisteredError,
@@ -63,6 +64,15 @@ describe('describeError', () => {
6364
expect(result.hint).toContain('internal workflow SDK error');
6465
});
6566

67+
test('CorruptedEventLogError is attributed to the SDK with a distinct code', () => {
68+
const result = describeError(
69+
new CorruptedEventLogError('corrupted event log')
70+
);
71+
expect(result.attribution).toBe('sdk');
72+
expect(result.errorCode).toBe(RUN_ERROR_CODES.CORRUPTED_EVENT_LOG);
73+
expect(result.hint).toContain('event log contains');
74+
});
75+
6676
test('StepNotRegisteredError (subclass of WorkflowRuntimeError) is attributed to the SDK', () => {
6777
const result = describeError(new StepNotRegisteredError('missingStep'));
6878
expect(result.attribution).toBe('sdk');
@@ -144,6 +154,23 @@ describe('describeRunError', () => {
144154
expect(result.hint).toContain('replay took too long');
145155
});
146156

157+
test('CORRUPTED_EVENT_LOG errorCode is attributed to the SDK', () => {
158+
const result = describeRunError({
159+
errorCode: RUN_ERROR_CODES.CORRUPTED_EVENT_LOG,
160+
});
161+
expect(result.attribution).toBe('sdk');
162+
expect(result.hint).toContain('event log contains');
163+
});
164+
165+
test('CorruptedEventLogError name restores the distinct code', () => {
166+
const result = describeRunError({
167+
errorName: 'CorruptedEventLogError',
168+
});
169+
expect(result.attribution).toBe('sdk');
170+
expect(result.errorCode).toBe(RUN_ERROR_CODES.CORRUPTED_EVENT_LOG);
171+
expect(result.hint).toContain('event log contains');
172+
});
173+
147174
test('MAX_DELIVERIES_EXCEEDED errorCode is attributed to the SDK', () => {
148175
const result = describeRunError({
149176
errorCode: RUN_ERROR_CODES.MAX_DELIVERIES_EXCEEDED,
@@ -226,6 +253,18 @@ describe('describeError — payload shape snapshots', () => {
226253
`);
227254
});
228255

256+
test('CorruptedEventLogError payload', () => {
257+
expect(
258+
describeError(new CorruptedEventLogError('event mismatch'))
259+
).toMatchInlineSnapshot(`
260+
{
261+
"attribution": "sdk",
262+
"errorCode": "CORRUPTED_EVENT_LOG",
263+
"hint": "The workflow event log contains orphaned or mismatched events and cannot be replayed. This is an internal workflow SDK error; please report it with the runId.",
264+
}
265+
`);
266+
});
267+
229268
test('REPLAY_TIMEOUT via precomputed errorCode payload', () => {
230269
expect(
231270
describeError(undefined, RUN_ERROR_CODES.REPLAY_TIMEOUT)

‎packages/core/src/describe-error.ts‎

Lines changed: 39 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import {
2+
CorruptedEventLogError,
23
RUN_ERROR_CODES,
34
type RunErrorCode,
45
SerializationError,
@@ -76,6 +77,8 @@ const CONTEXT_ERROR_HINT =
7677
'A workflow-only or step-only API was called from the wrong context. The error message includes the exact API and how to move the call.';
7778
const RUNTIME_ERROR_HINT =
7879
'This is an internal workflow SDK error, not a bug in your code. If it keeps happening, please report it with the stack trace and the runId.';
80+
const CORRUPTED_EVENT_LOG_HINT =
81+
'The workflow event log contains orphaned or mismatched events and cannot be replayed. This is an internal workflow SDK error; please report it with the runId.';
7982
const REPLAY_TIMEOUT_HINT =
8083
'The workflow replay took too long. This usually means the event log is unusually large or the workflow function is doing heavy synchronous work between step boundaries.';
8184
const MAX_DELIVERIES_HINT =
@@ -100,24 +103,41 @@ function normalizeErrorCode(code: string | undefined): RunErrorCode {
100103
export function describeRunError(
101104
signal: PersistedErrorSignal
102105
): ErrorDescription {
103-
const errorCode = normalizeErrorCode(signal.errorCode);
104106
const name = signal.errorName;
107+
const errorCode =
108+
name === 'CorruptedEventLogError'
109+
? RUN_ERROR_CODES.CORRUPTED_EVENT_LOG
110+
: normalizeErrorCode(signal.errorCode);
105111

106112
if (name === 'SerializationError') {
107113
return { attribution: 'user', errorCode, hint: SERIALIZATION_ERROR_HINT };
108114
}
109115
if (name && CONTEXT_ERROR_NAMES.has(name)) {
110116
return { attribution: 'user', errorCode, hint: CONTEXT_ERROR_HINT };
111117
}
112-
if (name === 'WorkflowRuntimeError' || name === 'StepNotRegisteredError') {
113-
return { attribution: 'sdk', errorCode, hint: RUNTIME_ERROR_HINT };
118+
if (name === 'CorruptedEventLogError') {
119+
return {
120+
attribution: 'sdk',
121+
errorCode,
122+
hint: CORRUPTED_EVENT_LOG_HINT,
123+
};
114124
}
115125
if (errorCode === RUN_ERROR_CODES.REPLAY_TIMEOUT) {
116126
return { attribution: 'sdk', errorCode, hint: REPLAY_TIMEOUT_HINT };
117127
}
118128
if (errorCode === RUN_ERROR_CODES.MAX_DELIVERIES_EXCEEDED) {
119129
return { attribution: 'sdk', errorCode, hint: MAX_DELIVERIES_HINT };
120130
}
131+
if (errorCode === RUN_ERROR_CODES.CORRUPTED_EVENT_LOG) {
132+
return {
133+
attribution: 'sdk',
134+
errorCode,
135+
hint: CORRUPTED_EVENT_LOG_HINT,
136+
};
137+
}
138+
if (name === 'WorkflowRuntimeError' || name === 'StepNotRegisteredError') {
139+
return { attribution: 'sdk', errorCode, hint: RUNTIME_ERROR_HINT };
140+
}
121141
if (errorCode === RUN_ERROR_CODES.RUNTIME_ERROR) {
122142
return { attribution: 'sdk', errorCode, hint: RUNTIME_ERROR_HINT };
123143
}
@@ -167,6 +187,14 @@ export function describeError(
167187
};
168188
}
169189

190+
if (CorruptedEventLogError.is(err)) {
191+
return {
192+
attribution: 'sdk',
193+
errorCode: effectiveCode,
194+
hint: CORRUPTED_EVENT_LOG_HINT,
195+
};
196+
}
197+
170198
if (err instanceof WorkflowRuntimeError) {
171199
return {
172200
attribution: 'sdk',
@@ -191,5 +219,13 @@ export function describeError(
191219
};
192220
}
193221

222+
if (effectiveCode === RUN_ERROR_CODES.CORRUPTED_EVENT_LOG) {
223+
return {
224+
attribution: 'sdk',
225+
errorCode: effectiveCode,
226+
hint: CORRUPTED_EVENT_LOG_HINT,
227+
};
228+
}
229+
194230
return { attribution: 'user', errorCode: effectiveCode };
195231
}

‎packages/core/src/log-format.test.ts‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -74,21 +74,21 @@ describe('composeLogLine', () => {
7474
test('renders sdk-attributed errors with the sdk badge', () => {
7575
const out = composeLogLine(
7676
PREFIX,
77-
'Workflow myFlow failed due to an SDK runtime error\nWorkflowRuntimeError: corrupted event log',
77+
'Workflow myFlow failed due to an SDK runtime error\nCorruptedEventLogError: corrupted event log',
7878
{
79-
errorCode: 'RUNTIME_ERROR',
79+
errorCode: 'CORRUPTED_EVENT_LOG',
8080
errorAttribution: 'sdk',
81-
errorName: 'WorkflowRuntimeError',
81+
errorName: 'CorruptedEventLogError',
8282
errorMessage: 'corrupted event log',
8383
hint: 'This is an internal workflow SDK error.',
8484
}
8585
);
8686
expect(out).toMatchInlineSnapshot(`
8787
"[workflow-sdk] Workflow myFlow failed due to an SDK runtime error
88-
sdk error · WorkflowRuntimeError
89-
code RUNTIME_ERROR
88+
sdk error · CorruptedEventLogError
89+
code CORRUPTED_EVENT_LOG
9090
hint: This is an internal workflow SDK error.
91-
WorkflowRuntimeError: corrupted event log"
91+
CorruptedEventLogError: corrupted event log"
9292
`);
9393
});
9494

‎packages/core/src/runtime.test.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -152,7 +152,7 @@ describe('workflowEntrypoint replay guards', () => {
152152
expect.objectContaining({
153153
eventType: 'run_failed',
154154
eventData: expect.objectContaining({
155-
errorCode: RUN_ERROR_CODES.RUNTIME_ERROR,
155+
errorCode: RUN_ERROR_CODES.CORRUPTED_EVENT_LOG,
156156
}),
157157
})
158158
);
@@ -210,7 +210,7 @@ describe('workflowEntrypoint replay guards', () => {
210210
expect.objectContaining({
211211
eventType: 'run_failed',
212212
eventData: expect.objectContaining({
213-
errorCode: RUN_ERROR_CODES.RUNTIME_ERROR,
213+
errorCode: RUN_ERROR_CODES.CORRUPTED_EVENT_LOG,
214214
}),
215215
})
216216
);

‎packages/core/src/runtime.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1117,9 +1117,9 @@ export function workflowEntrypoint(
11171117
);
11181118
}
11191119

1120-
// Classify the error: WorkflowRuntimeError indicates an
1121-
// internal issue (corrupted event log, missing data);
1122-
// everything else is a user code error.
1120+
// Classify the error: WorkflowRuntimeError indicates
1121+
// an SDK/runtime issue, and selected subclasses use
1122+
// more specific codes for backend tracking.
11231123
const errorCode = classifyRunError(err);
11241124

11251125
runtimeLogger.error('Error while running workflow', {

‎packages/core/src/step.test.ts‎

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,8 @@
1-
import { FatalError, WorkflowRuntimeError } from '@workflow/errors';
1+
import {
2+
CorruptedEventLogError,
3+
FatalError,
4+
WorkflowRuntimeError,
5+
} from '@workflow/errors';
26
import { withResolvers } from '@workflow/utils';
37
import type { Event } from '@workflow/world';
48
import * as nanoid from 'nanoid';
@@ -459,7 +463,7 @@ describe('createUseStep', () => {
459463
void add(1, 2);
460464

461465
const workflowError = await errorReceived.promise;
462-
expect(workflowError).toBeInstanceOf(WorkflowRuntimeError);
466+
expect(workflowError).toBeInstanceOf(CorruptedEventLogError);
463467
expect(workflowError.message).toContain('Corrupted event log');
464468
expect(workflowError.message).toContain('step_created');
465469
expect(workflowError.message).toContain('subtract');
@@ -568,7 +572,7 @@ describe('createUseStep', () => {
568572
void add(1, 2);
569573

570574
const workflowError = await errorReceived.promise;
571-
expect(workflowError).toBeInstanceOf(WorkflowRuntimeError);
575+
expect(workflowError).toBeInstanceOf(CorruptedEventLogError);
572576
expect(workflowError.message).toContain('Corrupted event log');
573577
expect(workflowError.message).toContain('step_completed');
574578
expect(workflowError.message).toContain('subtract');
@@ -629,7 +633,7 @@ describe('createUseStep', () => {
629633
void add(1, 2);
630634

631635
const workflowError = await errorReceived.promise;
632-
expect(workflowError).toBeInstanceOf(WorkflowRuntimeError);
636+
expect(workflowError).toBeInstanceOf(CorruptedEventLogError);
633637
expect(workflowError.message).toContain('Corrupted event log');
634638
expect(workflowError.message).toContain('step_failed');
635639
expect(workflowError.message).toContain('subtract');
@@ -753,7 +757,7 @@ describe('createUseStep', () => {
753757
expect(error?.message).toBe('Plain error message');
754758
});
755759

756-
it('should invoke workflow error handler with WorkflowRuntimeError for unexpected event type', async () => {
760+
it('should invoke workflow error handler with CorruptedEventLogError for unexpected event type', async () => {
757761
// Simulate a corrupted event log where a step receives an unexpected event type
758762
// (e.g., a wait_completed event when expecting step_completed/step_failed)
759763
const ctx = setupWorkflowContext([
@@ -779,7 +783,7 @@ describe('createUseStep', () => {
779783
const stepPromise = add(1, 2);
780784

781785
const workflowError = await errorReceived.promise;
782-
expect(workflowError).toBeInstanceOf(WorkflowRuntimeError);
786+
expect(workflowError).toBeInstanceOf(CorruptedEventLogError);
783787
expect(workflowError?.message).toContain('Unexpected event type for step');
784788
expect(workflowError?.message).toContain('step_01K11TFZ62YS0YYFDQ3E8B9YCV');
785789
expect(workflowError?.message).toContain('add');

0 commit comments

Comments
 (0)