fix: retry atomic workflow run reads (#1290)

This commit is contained in:
Brad Groux 2026-08-30 07:48:32 -05:00 • committed by GitHub
parent 3f83da36bc
commit 09714d5be3
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 44 additions and 23 deletions

View file

@ -69,6 +69,23 @@ describe('FileWorkflowRunRepository', () => {
).resolves.toContain('workflow-1');
});
it('retries when an atomic write replaces the path after open', async () => {
await repository.save(run('run_replaced'), 0);
const runPath = path.join(runsDir, 'run_replaced', 'run.json');
const actual = await vi.importActual<typeof import('node:fs/promises')>('node:fs/promises');
let pathChecks = 0;
vi.mocked(lstat).mockImplementation(async (filePath) => {
const stats = await actual.lstat(filePath);
if (path.resolve(String(filePath)) !== runPath || pathChecks++ > 0) return stats;
return Object.assign(Object.create(Object.getPrototypeOf(stats)), stats, {
ino: stats.ino + 1,
});
});
await expect(repository.get('run_replaced')).resolves.toMatchObject({ id: 'run_replaced' });
expect(pathChecks).toBe(2);
});
it('rejects symbolic links, changed files, and non-file state', async () => {
const runDir = path.join(runsDir, 'run_unsafe');
await mkdir(runDir, { recursive: true });

View file

@ -9,6 +9,7 @@ import { atomicWriteFile } from './fs-helpers.js';
const MAX_WORKFLOW_RUN_BYTES = 16 * 1024 * 1024;
const MAX_WORKFLOW_SNAPSHOT_BYTES = 4 * 1024 * 1024;
const MAX_WORKFLOW_RUN_READ_ATTEMPTS = 3;
export type WorkflowRunFilters = { taskId?: string; workflowId?: string; status?: string };
export type WorkflowRunMetadata = Pick<
@ -177,31 +178,34 @@ export class FileWorkflowRunRepository implements WorkflowRunRepository {
}
private async readRun(runPath: string): Promise<WorkflowRun | null> {
let handle: Awaited<ReturnType<typeof open>> | undefined;
try {
handle = await open(runPath, constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0));
const [pathStats, stats] = await Promise.all([lstat(runPath), handle.stat()]);
if (
pathStats.isSymbolicLink() ||
pathStats.dev !== stats.dev ||
pathStats.ino !== stats.ino
) {
throw new Error('Workflow run must not use a symbolic link or changed file');
for (let attempt = 0; attempt < MAX_WORKFLOW_RUN_READ_ATTEMPTS; attempt++) {
let handle: Awaited<ReturnType<typeof open>> | undefined;
try {
handle = await open(runPath, constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0));
const [pathStats, stats] = await Promise.all([lstat(runPath), handle.stat()]);
if (pathStats.isSymbolicLink()) {
throw new Error('Workflow run must not use a symbolic link');
}
if (pathStats.dev !== stats.dev || pathStats.ino !== stats.ino) {
if (attempt + 1 < MAX_WORKFLOW_RUN_READ_ATTEMPTS) continue;
throw new Error('Workflow run must not use a persistently changed file');
}
if (!stats.isFile() || stats.size > MAX_WORKFLOW_RUN_BYTES) {
throw new Error('Workflow run must use a bounded regular file');
}
return JSON.parse(await handle.readFile({ encoding: 'utf8' })) as WorkflowRun;
} catch (error) {
const errorCode = (error as NodeJS.ErrnoException).code;
if (errorCode === 'ENOENT') return null;
if (errorCode === 'ELOOP') {
throw new Error('Workflow run must not use a symbolic link', { cause: error });
}
throw error;
} finally {
await handle?.close();
}
if (!stats.isFile() || stats.size > MAX_WORKFLOW_RUN_BYTES) {
throw new Error('Workflow run must use a bounded regular file');
}
return JSON.parse(await handle.readFile({ encoding: 'utf8' })) as WorkflowRun;
} catch (error) {
const errorCode = (error as NodeJS.ErrnoException).code;
if (errorCode === 'ENOENT') return null;
if (errorCode === 'ELOOP') {
throw new Error('Workflow run must not use a symbolic link', { cause: error });
}
throw error;
} finally {
await handle?.close();
}
return null;
}
}