Skip to content
Closed
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
5 changes: 5 additions & 0 deletions .changeset/tidy-cursors-fail.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@croco/batch-core": patch
---

Record checkpoint restoration failures through the step classifier and execution failure lifecycle while preserving the original error.
35 changes: 17 additions & 18 deletions packages/batch-core/src/libs/ChunkExecutor.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
import {
ExecutionProblems,
type Execution,
Expand Down Expand Up @@ -69,26 +69,25 @@
assertValidChunkSize(step?.chunkSize);
const execution = await this.resolveExecution(executionId, options);

// 2. Restore checkpoint if available
const checkpointKey = `${step.name}.cursor`;
const processedCountKey = `${step.name}.processedCount`;
let restoredCheckpoint = false;
if (this.hasCheckpoint(execution, checkpointKey) && isCheckpointable(step.reader)) {
step.reader.restoreCheckpoint(execution.checkpoints[checkpointKey]);
restoredCheckpoint = true;
}
try {
// 2. Restore checkpoint if available
const checkpointKey = `${step.name}.cursor`;
const processedCountKey = `${step.name}.processedCount`;
let restoredCheckpoint = false;
if (this.hasCheckpoint(execution, checkpointKey) && isCheckpointable(step.reader)) {
step.reader.restoreCheckpoint(execution.checkpoints[checkpointKey]);
restoredCheckpoint = true;
}

let items: O[] = [];
let chunkInputCount = 0;
let processedCount = this.resolveProcessedCount(
execution,
restoredCheckpoint,
processedCountKey,
);
const totalCount = execution.progress?.total;
let items: O[] = [];
let chunkInputCount = 0;
let processedCount = this.resolveProcessedCount(
execution,
restoredCheckpoint,
processedCountKey,
);
const totalCount = execution.progress?.total;

// 3. Read - Process - Write loop
try {
while (true) {
// Read
const item = await step.reader.read();
Expand Down
74 changes: 74 additions & 0 deletions packages/batch-core/src/tests/ChunkExecutor.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,80 @@ describe("ChunkExecutor", () => {
});
});

it.each([false, true])(
"records checkpoint restoration failures with retryable=%s before reading or writing",
async (retryable) => {
const manager = new ExecutionManagerImpl(new TestExecutionStore());
const execution = await manager.create({ type: "batch-job", maxAttempts: 2 });
await manager.checkpoint(execution.id, "restore-step.cursor", { offset: 2 });
const fail = vi.spyOn(manager, "fail");
const original = new Error("unsupported saved cursor");
const reader = createCheckpointReader([1, 2, 3]);
reader.restoreCheckpoint.mockImplementation(() => {
throw original;
});
const writer = { write: vi.fn().mockResolvedValue(undefined) };
const classifyFailure = vi.fn(() => ({ retryable, code: "invalid-cursor" }));

await expect(
new ChunkExecutor(manager).execute(
execution.id,
new Step<number, number>({
name: "restore-step",
reader,
writer,
chunkSize: 1,
classifyFailure,
}),
),
).rejects.toBe(original);

expect(reader.restoreCheckpoint).toHaveBeenCalledExactlyOnceWith({ offset: 2 });
expect(classifyFailure).toHaveBeenCalledExactlyOnceWith(original, {
executionId: execution.id,
stepName: "restore-step",
});
expect(fail).toHaveBeenCalledExactlyOnceWith(execution.id, {
message: original.message,
stack: original.stack,
retryable,
code: "invalid-cursor",
});
expect(await manager.get(execution.id)).toMatchObject({
status: retryable ? "retrying" : "failed",
checkpoints: { "restore-step.cursor": { offset: 2 } },
});
expect(reader.read).not.toHaveBeenCalled();
expect(writer.write).not.toHaveBeenCalled();
},
);

it("preserves start failures without classifying or recording a step failure", async () => {
const original = new Error("cannot start execution");
vi.mocked(executionManager.start).mockRejectedValue(original);
const reader = createCheckpointReader([1]);
const writer = { write: vi.fn() };
const classifyFailure = vi.fn(() => false);

await expect(
executor.execute(
"exec-1",
new Step<number, number>({
name: "restore-step",
reader,
writer,
classifyFailure,
}),
),
).rejects.toBe(original);

expect(classifyFailure).not.toHaveBeenCalled();
expect(executionManager.fail).not.toHaveBeenCalled();
expect(reader.restoreCheckpoint).not.toHaveBeenCalled();
expect(reader.read).not.toHaveBeenCalled();
expect(writer.write).not.toHaveBeenCalled();
});

it.each([0, 1, 4, 5])(
"should checkpoint all-filtered input at chunk boundaries and flush the remainder (%i items)",
async (count) => {
Expand Down
Loading