diff --git a/.changeset/world-local-per-run-dirs.md b/.changeset/world-local-per-run-dirs.md new file mode 100644 index 0000000000..54e738c17c --- /dev/null +++ b/.changeset/world-local-per-run-dirs.md @@ -0,0 +1,5 @@ +--- +"@workflow/world-local": minor +--- + +Store event and step files in one directory per run for faster reads. Local run data written by earlier releases is deleted on upgrade. diff --git a/docs/content/worlds/v5/local.mdx b/docs/content/worlds/v5/local.mdx index 76ef47a457..78a352cdee 100644 --- a/docs/content/worlds/v5/local.mdx +++ b/docs/content/worlds/v5/local.mdx @@ -117,6 +117,10 @@ WORKFLOW_TARGET_WORLD="./my-world.ts" `createWorld()` also accepts `tag`, which scopes local storage files to a suffix. It is mainly used by test harnesses that share one `.workflow-data` directory. +## Data directory layout + +Each run's event and step files are stored in their own directory (`events//`, `steps//`), so reading or appending to a run costs time proportional to that run. Data directories written by earlier releases kept every file directly in `events/` and `steps/`. They are not migrated: when the World starts on one, it deletes the local run data and starts with an empty data directory. A directory that mixes both layouts, or holds `.json` files in `events/` or `steps/` that the World did not write, is never deleted: starting fails with a `DataDirLayoutError` that lists the files. Stop any older version still using the directory, then delete it or remove those files. + ## Limitations The Local World is designed for development, not production: diff --git a/packages/world-local/src/fs.ts b/packages/world-local/src/fs.ts index af18b50b6d..7f08019567 100644 --- a/packages/world-local/src/fs.ts +++ b/packages/world-local/src/fs.ts @@ -147,6 +147,22 @@ export async function withWindowsRetry( throw new Error('Retry loop exited unexpectedly'); } +/** + * `Promise.all` that waits for every promise to settle before rejecting with + * the first rejection (in input order). Use it where the caller finishing + * must mean none of its filesystem work is still running, e.g. `clear()`. + */ +export async function settleAll( + promises: Iterable> +): Promise { + const results = await Promise.allSettled(promises); + const failed = results.find( + (r): r is PromiseRejectedResult => r.status === 'rejected' + ); + if (failed) throw failed.reason; + return results.map((r) => (r as PromiseFulfilledResult).value); +} + /** * Clear write-path caches. Useful for testing or when files are deleted externally. */ @@ -219,6 +235,34 @@ export function taggedPath( return resolveWithinBase(basedir, entityDir, filename); } +/** + * Entity directories whose files are stored one subdirectory per run. + * + * Event and step files are only ever looked up or listed for one run at a + * time, so they live under `//`. That keeps every per-run + * read and listing proportional to that run's own files rather than to every + * file the data directory has accumulated. File names keep their + * `${runId}-` prefix. + * + * Releases before this one kept these files directly in `events/` and + * `steps/`; `initDataDir` wipes such a data directory rather than reading it. + */ +export const RUN_SCOPED_ENTITY_DIRS = ['events', 'steps'] as const; +export type RunScopedEntityDir = (typeof RUN_SCOPED_ENTITY_DIRS)[number]; + +/** + * The entity-relative directory holding one run's files for a run-scoped + * entity: `runEntityDir('events', 'wrun_ABC')` → `events/wrun_ABC`. + * Pass the result as the `entityDir` of {@link taggedPath} and friends. + */ +export function runEntityDir( + entityDir: RunScopedEntityDir, + runId: string +): string { + assertSafeEntityId('runId', runId); + return path.join(entityDir, runId); +} + /** * Read a JSON entity with tagged fallback. * When a tag is set, tries the tagged path first, then falls back to the @@ -256,6 +300,7 @@ export async function listTaggedFiles( ): Promise { const suffix = `.${tag}.json`; try { + await assertNotSymlinkedRunDir(dirPath); const files = await fs.readdir(dirPath); return files.filter((f) => f.endsWith(suffix)); } catch (error) { @@ -275,6 +320,7 @@ export async function listTaggedFilesByExtension( ): Promise { const suffix = `.${tag}${extension}`; try { + await assertNotSymlinkedRunDir(dirPath); const files = await fs.readdir(dirPath); return files.filter((f) => f.endsWith(suffix)); } catch (error) { @@ -286,12 +332,16 @@ export async function listTaggedFilesByExtension( export async function ensureDir(dirPath: string): Promise { const resolvedPath = path.resolve(dirPath); if (fsState.createdDirectoriesCache.has(resolvedPath)) { + // Checked on every use, not once: the directory can be replaced after + // it was cached. + await assertNotSymlinkedRunDir(resolvedPath); return; } + let mkdirError: unknown; try { await fs.mkdir(resolvedPath, { recursive: true }); - fsState.createdDirectoriesCache.add(resolvedPath); } catch (error) { + mkdirError = error; // A filesystem that refuses the directory outright will refuse every write // into it too, and the caller's write would surface as a confusing ENOENT // on the file rather than a missing directory. Report it here instead, @@ -307,6 +357,50 @@ export async function ensureDir(dirPath: string): Promise { } // Ignore if already exists } + if (mkdirError !== undefined) { + // The fallback above accepts an existing directory via `stat`, which + // follows symlinks, so a symlinked run directory still has to be refused + // here. Any other `lstat` failure keeps the historical "ignore" behavior. + await assertNotSymlinkedRunDir(resolvedPath).catch((error) => { + if (error instanceof SymlinkedRunDirError) throw error; + }); + return; + } + await assertNotSymlinkedRunDir(resolvedPath); + fsState.createdDirectoriesCache.add(resolvedPath); +} + +/** Thrown by {@link assertNotSymlinkedRunDir}. */ +export class SymlinkedRunDirError extends WorkflowWorldError { + constructor(dirPath: string) { + super( + `Refusing to use symlinked run directory ${dirPath}: ` + + `replace it with a real directory.` + ); + this.name = 'SymlinkedRunDirError'; + } +} + +/** + * A run's `events/` or `steps/` directory must be a real + * directory: `mkdir` accepts a symlink to one, and reading, writing or + * deleting through it would reach files outside the data directory. Every + * primitive below that touches a file in, or lists, such a directory calls + * this first (one `lstat`); a missing directory passes. + */ +export async function assertNotSymlinkedRunDir(dirPath: string): Promise { + const parent = path.basename(path.dirname(dirPath)); + if (!(RUN_SCOPED_ENTITY_DIRS as readonly string[]).includes(parent)) return; + let stats: import('node:fs').Stats; + try { + stats = await fs.lstat(dirPath); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return; + throw error; + } + if (stats.isSymbolicLink()) { + throw new SymlinkedRunDirError(dirPath); + } } async function isExistingDirectory(dirPath: string): Promise { @@ -440,6 +534,7 @@ export async function readJSON( filePath: string, decoder: z.ZodType ): Promise { + await assertNotSymlinkedRunDir(path.dirname(filePath)); try { const content = await withWindowsRetry(() => fs.readFile(filePath, 'utf-8') @@ -452,6 +547,7 @@ export async function readJSON( } export async function readBuffer(filePath: string): Promise { + await assertNotSymlinkedRunDir(path.dirname(filePath)); const content = await fs.readFile(filePath); return content; } @@ -459,6 +555,7 @@ export async function readBuffer(filePath: string): Promise { export async function readFirstByte( filePath: string ): Promise { + await assertNotSymlinkedRunDir(path.dirname(filePath)); const file = await fs.open(filePath, 'r'); try { const byte = Buffer.allocUnsafe(1); @@ -470,6 +567,7 @@ export async function readFirstByte( } export async function deleteJSON(filePath: string): Promise { + await assertNotSymlinkedRunDir(path.dirname(filePath)); try { // On Windows, a concurrent reader briefly holding the file open makes // unlink fail with EPERM (share violation), so retry like the other @@ -565,6 +663,7 @@ export async function listFilesByExtension( extension: string ): Promise { try { + await assertNotSymlinkedRunDir(dirPath); const files = await fs.readdir(dirPath); return files .filter((f) => f.endsWith(extension)) @@ -575,6 +674,28 @@ export async function listFilesByExtension( } } +/** + * Every run subdirectory of a run-scoped entity directory, as absolute paths. + * For the few whole-store walks (index backfills, tagged `clear()`); per-run + * reads go straight to {@link runEntityDir}. + */ +export async function listRunScopedDirs( + basedir: string, + entityDir: RunScopedEntityDir +): Promise { + const root = path.join(basedir, entityDir); + let entries: import('node:fs').Dirent[]; + try { + entries = await fs.readdir(root, { withFileTypes: true }); + } catch (error) { + if ((error as any).code === 'ENOENT') return []; + throw error; + } + return entries + .filter((entry) => entry.isDirectory()) + .map((entry) => path.join(root, entry.name)); +} + interface PaginatedFileSystemQueryConfig { directory: string; schema: z.ZodType; diff --git a/packages/world-local/src/index.ts b/packages/world-local/src/index.ts index 195b2c0071..11591dafe1 100644 --- a/packages/world-local/src/index.ts +++ b/packages/world-local/src/index.ts @@ -14,6 +14,7 @@ import { listTaggedFiles, listTaggedFilesByExtension, readJSON, + settleAll, } from './fs.js'; import { initDataDir } from './init.js'; import { instrumentObject } from './instrumentObject.js'; @@ -27,6 +28,7 @@ export { UnwritableDataDirError } from './build-target-mismatch.js'; // Re-export init types and utilities for consumers export { DataDirAccessError, + DataDirLayoutError, DataDirVersionError, ensureDataDir, initDataDir, @@ -140,7 +142,7 @@ export function createWorld(args?: Partial): LocalWorld { const hooksDir = path.join(basedir, 'hooks'); const taggedHookFiles = await listTaggedFiles(hooksDir, tag); const { HookSchema } = await import('@workflow/world'); - await Promise.all( + await settleAll( taggedHookFiles.map(async (hookFile) => { const hook = await readJSON( path.join(hooksDir, hookFile), @@ -162,25 +164,40 @@ export function createWorld(args?: Partial): LocalWorld { }) ); - // Delete tagged entity files across all directories + // Delete tagged entity files across all directories. Steps and + // events live one directory per run, so walk only this tag's runs + // (found from `runs/*..json`) rather than every run the shared + // data directory holds. + const runScopedDirs = ( + await listTaggedFiles(path.join(basedir, 'runs'), tag) + ).flatMap((file) => { + const runId = file.slice(0, -`.${tag}.json`.length); + return [path.join('steps', runId), path.join('events', runId)]; + }); const entityDirs = [ 'runs', - 'steps', - 'events', + ...runScopedDirs, 'hooks', 'hooks/by-run', 'waits', 'streams/runs', ]; - await Promise.all( + await settleAll( entityDirs.map(async (dir) => { const fullDir = path.join(basedir, dir); const files = await listTaggedFiles(fullDir, tag); - await Promise.all( + await settleAll( files.map((f) => deleteJSON(path.join(fullDir, f))) ); }) ); + // Drop the run directories that clearing left empty. `rmdir` refuses + // a non-empty one, so another tag's (or untagged) files keep theirs. + await settleAll( + runScopedDirs.map((dir) => + fs.rmdir(path.join(basedir, dir)).catch(() => {}) + ) + ); // Delete tagged hook-index entries (nested per-key directories) for (const indexDir of ['token-index', 'id-index']) { const fullIndexDir = path.join(basedir, 'hooks', indexDir); @@ -192,13 +209,13 @@ export function createWorld(args?: Partial): LocalWorld { } catch { keyDirEntries = []; } - await Promise.all( + await settleAll( keyDirEntries .filter((entry) => entry.isDirectory()) .map(async (entry) => { const keyDir = path.join(fullIndexDir, entry.name); const taggedEntryFiles = await listTaggedFiles(keyDir, tag); - await Promise.all( + await settleAll( taggedEntryFiles.map((f) => deleteJSON(path.join(keyDir, f))) ); }) @@ -222,7 +239,7 @@ export function createWorld(args?: Partial): LocalWorld { } catch { streamDirEntries = []; } - await Promise.all( + await settleAll( streamDirEntries .filter((entry) => entry.isDirectory()) .map(async (entry) => { @@ -232,7 +249,7 @@ export function createWorld(args?: Partial): LocalWorld { tag, '.bin' ); - await Promise.all( + await settleAll( taggedBinFiles.map((f) => fs.unlink(path.join(streamChunkDir, f)).catch(() => {}) ) diff --git a/packages/world-local/src/init.test.ts b/packages/world-local/src/init.test.ts index aa2ab7e303..752134f539 100644 --- a/packages/world-local/src/init.test.ts +++ b/packages/world-local/src/init.test.ts @@ -11,6 +11,7 @@ import path from 'node:path'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { DataDirAccessError, + DataDirLayoutError, DataDirVersionError, ensureDataDir, formatVersion, @@ -128,8 +129,8 @@ describe('formatVersionFile', () => { }); describe('upgradeVersion', () => { - it('should log upgrade message', () => { - const consoleSpy = vi.spyOn(console, 'log').mockImplementation(() => {}); + it('should log upgrade message on stderr', () => { + const consoleSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); const oldVersion = parseVersion('3.0.0'); const newVersion = parseVersion('4.0.1-beta.20'); @@ -386,7 +387,7 @@ describe('initDataDir', () => { const currentVersion = `${packageInfo.name}@${packageInfo.version}`; writeFileSync(versionPath, currentVersion); - const consoleSpy = vi.spyOn(console, 'log').mockImplementation(() => {}); + const consoleSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); await initDataDir(dataDir); @@ -406,7 +407,7 @@ describe('initDataDir', () => { const versionPath = path.join(dataDir, 'version.txt'); writeFileSync(versionPath, '@workflow/world-local@3.0.0'); - const consoleSpy = vi.spyOn(console, 'log').mockImplementation(() => {}); + const consoleSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); await initDataDir(dataDir); @@ -431,7 +432,7 @@ describe('initDataDir', () => { const versionPath = path.join(dataDir, 'version.txt'); writeFileSync(versionPath, `${packageInfo.name}@${newerVersion}`); - const consoleSpy = vi.spyOn(console, 'log').mockImplementation(() => {}); + const consoleSpy = vi.spyOn(console, 'warn').mockImplementation(() => {}); // This will call upgradeVersion which just logs for now await initDataDir(dataDir); @@ -443,6 +444,110 @@ describe('initDataDir', () => { }); }); +describe('initDataDir legacy flat layout', () => { + let dataDir: string; + + beforeEach(() => { + dataDir = path.join( + tmpdir(), + `workflow-init-layout-test-${Date.now()}-${Math.random().toString(36).slice(2)}` + ); + mkdirSync(dataDir, { recursive: true }); + writeFileSync( + path.join(dataDir, 'version.txt'), + '@workflow/world-local@4.1.0' + ); + mkdirSync(path.join(dataDir, 'runs')); + writeFileSync(path.join(dataDir, 'runs', 'wrun_A.json'), '{}'); + }); + + afterEach(() => { + rmSync(dataDir, { recursive: true, force: true }); + vi.restoreAllMocks(); + }); + + it.each([ + 'events', + 'steps', + ])('wipes a data directory with flat %s files', async (entityDir) => { + mkdirSync(path.join(dataDir, entityDir)); + writeFileSync(path.join(dataDir, entityDir, 'wrun_A-x_1.json'), '{}'); + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + + await initDataDir(dataDir); + + expect(warn).toHaveBeenCalledWith( + expect.stringContaining('Deleting local workflow data') + ); + expect(existsSync(path.join(dataDir, 'runs'))).toBe(false); + expect(existsSync(path.join(dataDir, entityDir))).toBe(false); + const packageInfo = await getPackageInfo(); + expect(readFileSync(path.join(dataDir, 'version.txt'), 'utf-8')).toBe( + `${packageInfo.name}@${packageInfo.version}` + ); + }); + + it('keeps a data directory in the per-run layout', async () => { + mkdirSync(path.join(dataDir, 'events', 'wrun_A'), { recursive: true }); + writeFileSync( + path.join(dataDir, 'events', 'wrun_A', 'wrun_A-evnt_1.json'), + '{}' + ); + // Not an entity file, so not a sign of the old layout. + writeFileSync(path.join(dataDir, 'events', '.DS_Store'), ''); + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + + await initDataDir(dataDir); + + expect(warn).not.toHaveBeenCalledWith( + expect.stringContaining('Deleting local workflow data') + ); + expect(existsSync(path.join(dataDir, 'runs', 'wrun_A.json'))).toBe(true); + expect( + existsSync(path.join(dataDir, 'events', 'wrun_A', 'wrun_A-evnt_1.json')) + ).toBe(true); + }); + it.each([ + ['a flat event file next to per-run directories', 'wrun_A-evnt_2.json'], + ['a .json file this package did not write', 'notes.json'], + ])('refuses, without deleting, %s', async (_, stray) => { + mkdirSync(path.join(dataDir, 'events', 'wrun_A'), { recursive: true }); + writeFileSync( + path.join(dataDir, 'events', 'wrun_A', 'wrun_A-evnt_1.json'), + '{}' + ); + writeFileSync(path.join(dataDir, 'events', stray), '{}'); + vi.spyOn(console, 'warn').mockImplementation(() => {}); + + const error = await initDataDir(dataDir).catch((e) => e); + + expect(error).toBeInstanceOf(DataDirLayoutError); + expect(error.message).toContain(path.join('events', stray)); + expect(error.message).toContain(path.resolve(dataDir)); + for (const kept of [ + path.join('runs', 'wrun_A.json'), + path.join('events', 'wrun_A', 'wrun_A-evnt_1.json'), + path.join('events', stray), + ]) { + expect(existsSync(path.join(dataDir, kept))).toBe(true); + } + }); + + it('refuses an unrecognized file in an otherwise flat directory', async () => { + mkdirSync(path.join(dataDir, 'steps')); + writeFileSync(path.join(dataDir, 'steps', 'wrun_A-step_1.json'), '{}'); + writeFileSync(path.join(dataDir, 'steps', 'backup.json'), '{}'); + vi.spyOn(console, 'warn').mockImplementation(() => {}); + + await expect(initDataDir(dataDir)).rejects.toBeInstanceOf( + DataDirLayoutError + ); + expect(existsSync(path.join(dataDir, 'steps', 'wrun_A-step_1.json'))).toBe( + true + ); + }); +}); + describe('DataDirVersionError', () => { it('should store version information', () => { const oldVersion = parseVersion('1.0.0'); diff --git a/packages/world-local/src/init.ts b/packages/world-local/src/init.ts index b4a35e65e8..c56ed34190 100644 --- a/packages/world-local/src/init.ts +++ b/packages/world-local/src/init.ts @@ -2,13 +2,17 @@ import { access, constants, mkdir, + opendir, readFile, + rm, unlink, writeFile, } from 'node:fs/promises'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; import { globalSingleton } from '@workflow/utils'; +import { clearCreatedFilesCache, RUN_SCOPED_ENTITY_DIRS } from './fs.js'; +import { resetHookIndexEnsureCache } from './storage/hook-index.js'; /** Package name - hardcoded since it doesn't change */ const PACKAGE_NAME = '@workflow/world-local'; @@ -202,7 +206,8 @@ export function upgradeVersion( oldVersion: ParsedVersion, newVersion: ParsedVersion ): void { - console.log( + // stderr: CLI commands print JSON on stdout. + console.warn( `[world-local] Upgrading from version ${formatVersion(oldVersion)} to ${formatVersion(newVersion)}` ); } @@ -306,6 +311,87 @@ async function writeVersionFile( await writeFile(versionFilePath, content); } +/** + * Name of an event or step file in the old flat layout: + * `wrun_-[.].json`. + */ +const FLAT_ENTITY_FILE = /^wrun_[0-9A-Za-z]+-[^/\\]+\.json$/; + +interface EntityDirsLayout { + /** Root-level files named like old flat-layout event/step files. */ + flat: string[]; + /** Root-level `.json` entries that are neither of the above. */ + unrecognized: string[]; + /** Whether any per-run subdirectory exists. */ + runScoped: boolean; +} + +function classifyEntityDirEntry( + layout: EntityDirsLayout, + entityDir: string, + entry: import('node:fs').Dirent +): void { + if (entry.isDirectory()) { + layout.runScoped = true; + return; + } + if (!entry.name.endsWith('.json')) return; + const relative = path.join(entityDir, entry.name); + if (entry.isFile() && FLAT_ENTITY_FILE.test(entry.name)) { + layout.flat.push(relative); + } else { + layout.unrecognized.push(relative); + } +} + +/** + * What `events/` and `steps/` hold at their root: per-run subdirectories, + * old flat-layout event/step files, or `.json` entries this package did not + * write (non-JSON entries such as `.DS_Store` or temp files are ignored). + */ +async function inspectEntityDirs(dataDir: string): Promise { + const layout: EntityDirsLayout = { + flat: [], + unrecognized: [], + runScoped: false, + }; + for (const entityDir of RUN_SCOPED_ENTITY_DIRS) { + let dir: import('node:fs').Dir; + try { + dir = await opendir(path.join(dataDir, entityDir)); + } catch (error: unknown) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') continue; + throw error; + } + for await (const entry of dir) { + classifyEntityDirEntry(layout, entityDir, entry); + } + } + return layout; +} + +/** + * Thrown by {@link initDataDir} for a data directory it can neither use nor + * safely identify as old flat-layout data: per-run directories mixed with + * flat event/step files (an older release, or another implementation, is + * still writing to it), or `.json` entries it did not write. + */ +export class DataDirLayoutError extends Error { + constructor( + public readonly dataDir: string, + public readonly files: string[] + ) { + const shown = files.slice(0, 5).join(', '); + super( + `[world-local] ${path.resolve(dataDir)} holds ${files.length} event/step ` + + `file(s) this version does not read (${shown}${files.length > 5 ? ', …' : ''}). ` + + `If an older version of ${PACKAGE_NAME} is using this directory, stop it. ` + + `Then delete the directory, or move or remove those files.` + ); + this.name = 'DataDirLayoutError'; + } +} + /** * Gets the suggested downgrade version based on the old version. * If a specific version is suggested in the error, use that. @@ -335,6 +421,32 @@ export async function initDataDir(dataDir: string): Promise { // First ensure the directory exists and is accessible await ensureDataDir(dataDir); + // Local run data from the old flat layout is not migrated: a directory that + // holds only such data is wiped and started over. Anything else outside + // the per-run layout is refused rather than deleted. + const layout = await inspectEntityDirs(dataDir); + if ( + layout.unrecognized.length > 0 || + (layout.runScoped && layout.flat.length > 0) + ) { + throw new DataDirLayoutError(dataDir, [ + ...layout.unrecognized, + ...layout.flat, + ]); + } + if (layout.flat.length > 0) { + console.warn( + `[world-local] Deleting local workflow data in "${path.resolve(dataDir)}": ` + + `it was written by an older version of ${PACKAGE_NAME} with an ` + + `incompatible storage layout.` + ); + clearCreatedFilesCache(); + resetHookIndexEnsureCache(); + // Retries ENOTEMPTY from another process writing into it meanwhile. + await rm(dataDir, { recursive: true, force: true, maxRetries: 3 }); + await ensureDataDir(dataDir); + } + const packageInfo = await getPackageInfo(); const currentVersion = parseVersion(packageInfo.version); diff --git a/packages/world-local/src/run-dir-guards.test.ts b/packages/world-local/src/run-dir-guards.test.ts new file mode 100644 index 0000000000..e3935610c6 --- /dev/null +++ b/packages/world-local/src/run-dir-guards.test.ts @@ -0,0 +1,167 @@ +import { promises as fs } from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { + clearCreatedFilesCache, + SymlinkedRunDirError, + writeJSON, +} from './fs.js'; +import { createWorld } from './index.js'; +import { createStorage } from './storage.js'; +import { createRun, createStep } from './test-helpers.js'; + +let root: string; +let dataDir: string; + +beforeEach(async () => { + root = await fs.mkdtemp(path.join(os.tmpdir(), 'run-dir-guards-')); + dataDir = path.join(root, 'data'); + await fs.mkdir(dataDir); +}); + +afterEach(async () => { + vi.restoreAllMocks(); + clearCreatedFilesCache(); + await fs.rm(root, { recursive: true, force: true }); +}); + +async function seedRun(tag?: string): Promise { + const storage = createStorage(dataDir, tag); + const run = await createRun(storage, { + deploymentId: 'dep-1', + workflowName: 'wf', + input: new Uint8Array(), + }); + for (const stepId of ['step_1', 'step_2', 'step_3']) { + await createStep(storage, run.runId, { + stepId, + stepName: 'step', + input: new Uint8Array(), + }); + } + return run.runId; +} + +function deferred() { + let resolve!: () => void; + const promise = new Promise((r) => { + resolve = r; + }); + return { promise, resolve }; +} + +describe('clear()', () => { + it('a failed deletion rejects only after every started deletion finished', async () => { + const kept = await seedRun(); + const cleared = await seedRun('vitest-0'); + const clearedDir = path.join(dataDir, 'events', cleared); + const order: string[] = []; + const paused = deferred(); + const resume = deferred(); + const unlink = fs.unlink.bind(fs); + let calls = 0; + vi.spyOn(fs, 'unlink').mockImplementation(async (p) => { + if (typeof p === 'string' && path.dirname(p) === clearedDir) { + const call = ++calls; + if (call === 1) { + paused.resolve(); + await resume.promise; + await unlink(p); + order.push('sibling-done'); + return; + } + if (call === 2) { + order.push('failed'); + throw Object.assign(new Error('i/o error'), { code: 'EIO' }); + } + } + return unlink(p); + }); + + const clear = createWorld({ dataDir, tag: 'vitest-0' }) + .clear() + .then( + () => order.push('cleared'), + (error) => order.push(`rejected:${error.code}`) + ); + await paused.promise; + // A bounded wait for something that must not happen: the rejection + // surfacing while the first deletion is still held. + for (let i = 0; i < 20 && order.length < 2; i++) { + await new Promise((r) => setTimeout(r, 10)); + } + expect(order).toEqual(['failed']); + resume.resolve(); + await clear; + + expect(order).toEqual(['failed', 'sibling-done', 'rejected:EIO']); + expect(await fs.readdir(path.join(dataDir, 'events', kept))).toHaveLength( + 4 + ); + }); +}); + +describe('run directories', () => { + async function symlinkRunDir(entityDir: 'events' | 'steps', runId: string) { + const outside = path.join(root, `outside-${entityDir}`); + const runDir = path.join(dataDir, entityDir, runId); + await fs.rename(runDir, outside); + await fs.symlink(outside, runDir, 'dir'); + return outside; + } + + // No cache reset between creating the run and swapping its directory: the + // write-path caches already hold it, as in a long-running dev server. + it('refuses writes, reads and listings through a symlinked steps dir', async () => { + const runId = await seedRun(); + const storage = createStorage(dataDir); + const outside = await symlinkRunDir('steps', runId); + const before = (await fs.readdir(outside)).sort(); + + await expect( + createStep(storage, runId, { + stepId: 'step_4', + stepName: 'step', + input: new Uint8Array(), + }) + ).rejects.toThrow(/symlinked run directory/); + await expect(storage.steps.get(runId, 'step_1')).rejects.toThrow( + /symlinked run directory/ + ); + await expect(storage.steps.list({ runId })).rejects.toThrow( + /symlinked run directory/ + ); + expect((await fs.readdir(outside)).sort()).toEqual(before); + }); + + it('refuses a symlinked run dir when mkdir fails on a cold cache', async () => { + const runId = await seedRun(); + const outside = await symlinkRunDir('steps', runId); + const before = (await fs.readdir(outside)).sort(); + clearCreatedFilesCache(); + const runDir = path.join(dataDir, 'steps', runId); + const mkdir = fs.mkdir.bind(fs); + vi.spyOn(fs, 'mkdir').mockImplementation(async (p, options) => { + if (p === runDir) { + throw Object.assign(new Error('permission denied'), { code: 'EACCES' }); + } + return mkdir(p, options); + }); + + await expect( + writeJSON(path.join(runDir, `${runId}-step_9.json`), {}) + ).rejects.toThrow(SymlinkedRunDirError); + expect((await fs.readdir(outside)).sort()).toEqual(before); + }); + + it('refuses reading the event log through a symlinked events dir', async () => { + const runId = await seedRun(); + const storage = createStorage(dataDir); + await symlinkRunDir('events', runId); + + await expect(storage.events.list({ runId })).rejects.toThrow( + /symlinked run directory/ + ); + }); +}); diff --git a/packages/world-local/src/storage.test.ts b/packages/world-local/src/storage.test.ts index ec4546fa5a..251e3c2a05 100644 --- a/packages/world-local/src/storage.test.ts +++ b/packages/world-local/src/storage.test.ts @@ -730,6 +730,7 @@ describe('Storage', () => { const filePath = path.join( testDir, 'steps', + testRunId, `${testRunId}-step_123.json` ); const fileExists = await fs @@ -1318,6 +1319,7 @@ describe('Storage', () => { const filePath = path.join( testDir, 'events', + testRunId, `${testRunId}-${event.eventId}.json` ); const fileExists = await fs @@ -4002,7 +4004,12 @@ describe('Storage', () => { // Simulate a crash after the hook entity write but before the // event write by deleting the just-written event from disk. await fs.unlink( - path.join(testDir, 'events', `${testRunId}-${first.event.eventId}.json`) + path.join( + testDir, + 'events', + testRunId, + `${testRunId}-${first.event.eventId}.json` + ) ); // Sanity: the hook entity is still durable but the @@ -4200,7 +4207,12 @@ describe('Storage', () => { await fs.unlink(hookPath); await fs.unlink(tokenClaimPath); await fs.writeFile( - path.join(testDir, 'events', 'wrun_malformed-event.json'), + path.join( + testDir, + 'events', + testRunId, + `${testRunId}-evnt_malformed.json` + ), '{' ); @@ -4571,7 +4583,7 @@ describe('Storage', () => { JSON.stringify({ token, hookId, runId: run.runId }) ); const preExistingEventId = 'evnt_pre_upgrade_existing'; - const eventsDir = path.join(testDir, 'events'); + const eventsDir = path.join(testDir, 'events', run.runId); await fs.mkdir(eventsDir, { recursive: true }); await fs.writeFile( path.join(eventsDir, `${run.runId}-${preExistingEventId}.json`), @@ -5482,6 +5494,7 @@ describe('Storage', () => { const eventPath = path.join( testDir, 'events', + run.runId, `${run.runId}-${stalledEventId}.json` ); await expect(promoteExclusive(stagedPath, eventPath)).resolves.toBe( @@ -5629,7 +5642,7 @@ describe('Storage', () => { workflowName: 'test-workflow', input: new Uint8Array(), }); - const eventsDir = path.join(testDir, 'events'); + const eventsDir = path.join(testDir, 'events', run.runId); await fs.chmod(eventsDir, 0o000); try { diff --git a/packages/world-local/src/storage/events-storage.ts b/packages/world-local/src/storage/events-storage.ts index 04fff86645..d4a429001b 100644 --- a/packages/world-local/src/storage/events-storage.ts +++ b/packages/world-local/src/storage/events-storage.ts @@ -69,6 +69,7 @@ import { readJSON, readJSONWithFallback, resolveWithinBase, + runEntityDir, SORT_KEY_CURSOR_PREFIX, stripTag, taggedPath, @@ -244,7 +245,7 @@ async function findCommittedResumeEvent( for (const eventId of scan.ids) { const event = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', runId), `${runId}-${eventId}`, ReadEventSchema, tag @@ -369,7 +370,7 @@ async function findExistingHookCreatedEventId( correlationId: string ): Promise { const result = await paginatedFileSystemQuery({ - directory: path.join(basedir, 'events'), + directory: path.join(basedir, runEntityDir('events', runId)), schema: ReadEventSchema, filePrefix: `${runId}-`, filter: (event) => @@ -412,7 +413,7 @@ async function repairHookEntityFromPersistedEvent( const compositeKey = `${runId}-${persistedEventId}`; const persistedEvent = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', runId), compositeKey, ReadEventSchema, tag @@ -696,10 +697,10 @@ export function createEventsStorage( const fileId = `${runId}-${slotToEventId(slot)}`; for (const candidate of tag ? [ - taggedPath(basedir, 'events', fileId, tag), - taggedPath(basedir, 'events', fileId), + taggedPath(basedir, runEntityDir('events', runId), fileId, tag), + taggedPath(basedir, runEntityDir('events', runId), fileId), ] - : [taggedPath(basedir, 'events', fileId)]) { + : [taggedPath(basedir, runEntityDir('events', runId), fileId)]) { try { await fs.stat(candidate); return true; @@ -884,7 +885,7 @@ export function createEventsStorage( for (let attempt = 0; ; attempt++) { const eventPath = taggedPath( basedir, - 'events', + runEntityDir('events', current.runId), `${current.runId}-${current.eventId}`, tag ); @@ -936,7 +937,7 @@ export function createEventsStorage( const queryRunEvents = (runId: string, pagination: PaginationOptions) => paginatedFileSystemQuery({ - directory: path.join(basedir, 'events'), + directory: path.join(basedir, runEntityDir('events', runId)), schema: ReadEventSchema, cachedItems: eventCache, filePrefix: `${runId}-`, @@ -1341,7 +1342,7 @@ export function createEventsStorage( const stepCompositeKey = `${effectiveRunId}-${data.correlationId}`; validatedStep = await readJSONWithFallback( basedir, - 'steps', + runEntityDir('steps', effectiveRunId), stepCompositeKey, StepSchema, tag @@ -1433,7 +1434,7 @@ export function createEventsStorage( ) { const atClaimedId = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${committedClaim.eventId}`, ReadEventSchema, tag @@ -1526,7 +1527,7 @@ export function createEventsStorage( } const atClaimedId = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${claim.eventId}`, ReadEventSchema, tag @@ -2099,7 +2100,12 @@ export function createEventsStorage( }; const stepCompositeKey = `${effectiveRunId}-${data.correlationId}`; await writeJSON( - taggedPath(basedir, 'steps', stepCompositeKey, tag), + taggedPath( + basedir, + runEntityDir('steps', effectiveRunId), + stepCompositeKey, + tag + ), step ); } else if (data.eventType === 'step_started') { @@ -2158,7 +2164,7 @@ export function createEventsStorage( await writeJSON( taggedPath( basedir, - 'steps', + runEntityDir('steps', effectiveRunId), `${effectiveRunId}-${data.correlationId}`, tag ), @@ -2219,7 +2225,7 @@ export function createEventsStorage( const stepCompositeKey = `${effectiveRunId}-${data.correlationId}`; const freshStep = await readJSONWithFallback( basedir, - 'steps', + runEntityDir('steps', effectiveRunId), stepCompositeKey, StepSchema, tag @@ -2242,7 +2248,12 @@ export function createEventsStorage( updatedAt: now, }; await writeJSON( - taggedPath(basedir, 'steps', stepCompositeKey, tag), + taggedPath( + basedir, + runEntityDir('steps', effectiveRunId), + stepCompositeKey, + tag + ), step, { overwrite: true } ); @@ -2277,7 +2288,12 @@ export function createEventsStorage( updatedAt: now, }; await writeJSON( - taggedPath(basedir, 'steps', stepCompositeKey, tag), + taggedPath( + basedir, + runEntityDir('steps', effectiveRunId), + stepCompositeKey, + tag + ), step, { overwrite: true } ); @@ -2317,7 +2333,12 @@ export function createEventsStorage( updatedAt: now, }; await writeJSON( - taggedPath(basedir, 'steps', stepCompositeKey, tag), + taggedPath( + basedir, + runEntityDir('steps', effectiveRunId), + stepCompositeKey, + tag + ), step, { overwrite: true } ); @@ -2339,7 +2360,12 @@ export function createEventsStorage( updatedAt: now, }; await writeJSON( - taggedPath(basedir, 'steps', stepCompositeKey, tag), + taggedPath( + basedir, + runEntityDir('steps', effectiveRunId), + stepCompositeKey, + tag + ), step, { overwrite: true } ); @@ -3008,7 +3034,7 @@ export function createEventsStorage( let eventPath = taggedPath( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${eventId}`, tag ); @@ -3060,7 +3086,7 @@ export function createEventsStorage( event = { ...event, eventId }; eventPath = taggedPath( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${eventId}`, tag ); @@ -3107,7 +3133,7 @@ export function createEventsStorage( event = prePublishedEvent; eventPath = taggedPath( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${eventId}`, tag ); @@ -3220,7 +3246,7 @@ export function createEventsStorage( if (data.eventType === 'hook_received' && params?.resumeId) { const occupant = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${eventId}`, ReadEventSchema, tag @@ -3248,7 +3274,7 @@ export function createEventsStorage( if (data.eventType === 'hook_created' && data.correlationId) { const occupant = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', effectiveRunId), `${effectiveRunId}-${eventId}`, ReadEventSchema, tag @@ -3453,7 +3479,7 @@ export function createEventsStorage( const compositeKey = `${runId}-${eventId}`; const event = await readJSONWithFallback( basedir, - 'events', + runEntityDir('events', runId), compositeKey, ReadEventSchema, tag @@ -3493,7 +3519,7 @@ export function createEventsStorage( assertSafeEntityId('runId', params.runId); const resolveData = params.resolveData ?? DEFAULT_RESOLVE_DATA_OPTION; const result = await paginatedFileSystemQuery({ - directory: path.join(basedir, 'events'), + directory: path.join(basedir, runEntityDir('events', params.runId)), schema: ReadEventSchema, cachedItems: eventCache, // Scoped to the run's own event files, since a correlation id diff --git a/packages/world-local/src/storage/helpers.ts b/packages/world-local/src/storage/helpers.ts index c4e70e5c4b..7264038785 100644 --- a/packages/world-local/src/storage/helpers.ts +++ b/packages/world-local/src/storage/helpers.ts @@ -8,11 +8,13 @@ import { lock } from 'proper-lockfile'; import { decodeTime, monotonicFactory } from 'ulid'; import { z } from 'zod'; import { + assertNotSymlinkedRunDir, deleteJSON, hasTag, isUntagged, readJSON, resolveWithinBase, + runEntityDir, stripTag, ulidToDate, withWindowsRetry, @@ -335,11 +337,10 @@ export interface RunEventIdScan { } /** - * Scans the events directory for one run's ids, honoring tag visibility. + * Scans one run's events directory for its ids, honoring tag visibility. * - * O(all event files), like every other directory-walking read in this - * backend. Callers that run it per write memoize the result and use the - * publish itself to detect when the memo has fallen behind. + * O(the run's event files). Callers that run it per write memoize the result + * and use the publish itself to detect when the memo has fallen behind. */ export async function scanRunEventIds( basedir: string, @@ -348,7 +349,9 @@ export async function scanRunEventIds( ): Promise { let files: string[] = []; try { - files = await fs.readdir(path.join(basedir, 'events')); + const eventsDir = path.join(basedir, runEntityDir('events', runId)); + await assertNotSymlinkedRunDir(eventsDir); + files = await fs.readdir(eventsDir); } catch (error) { // Only ENOENT ("no events directory yet") means there is provably // nothing visible. Any other failure would silently report an empty run, diff --git a/packages/world-local/src/storage/hook-index.ts b/packages/world-local/src/storage/hook-index.ts index c379868f6d..14b989e4e2 100644 --- a/packages/world-local/src/storage/hook-index.ts +++ b/packages/world-local/src/storage/hook-index.ts @@ -10,8 +10,10 @@ import { hasTag, isUntagged, listJSONFiles, + listRunScopedDirs, readJSON, resolveWithinBase, + runEntityDir, stripTag, taggedPath, writeExclusive, @@ -264,31 +266,35 @@ async function ensureHookIndexesImpl(basedir: string): Promise { // Marker absent, so backfill below. } - const eventsDir = path.join(basedir, 'events'); - await forEachConcurrent( - await listJSONFiles(eventsDir), - 32, - async (fileId) => { - const event = await readEventLenient( - path.join(eventsDir, `${fileId}.json`) + const indexEventFile = async (eventFile: string) => { + const fileId = path.basename(eventFile, '.json'); + const event = await readEventLenient(eventFile); + if (!event || event.eventType !== 'hook_created') return; + if (typeof event.correlationId !== 'string') return; + const token = (event.eventData as { token?: unknown } | undefined)?.token; + if (typeof token !== 'string') return; + try { + await writeHookCreatedIndexEntries( + basedir, + token, + event.runId, + event.correlationId, + event.eventId, + tagOf(fileId) ); - if (!event || event.eventType !== 'hook_created') return; - if (typeof event.correlationId !== 'string') return; - const token = (event.eventData as { token?: unknown } | undefined)?.token; - if (typeof token !== 'string') return; - try { - await writeHookCreatedIndexEntries( - basedir, - token, - event.runId, - event.correlationId, - event.eventId, - tagOf(fileId) - ); - } catch { - // Unsafe ids cannot have been written by this storage layer; skip. - } + } catch { + // Unsafe ids cannot have been written by this storage layer; skip. } + }; + // One run at a time per worker, so only one run's file names are held in + // memory per worker rather than every event path in the store. + await forEachConcurrent( + await listRunScopedDirs(basedir, 'events'), + 8, + async (runDir) => + forEachConcurrent(await listJSONFiles(runDir), 32, (fileId) => + indexEventFile(path.join(runDir, `${fileId}.json`)) + ) ); const hooksDir = path.join(basedir, 'hooks'); @@ -373,7 +379,7 @@ export async function findIndexedHookCreatedEvent( try { eventPath = taggedPath( basedir, - 'events', + runEntityDir('events', entry.runId), `${entry.runId}-${eventId}`, tagOf(entryId) ); diff --git a/packages/world-local/src/storage/legacy.ts b/packages/world-local/src/storage/legacy.ts index f20c3ad4da..e9be943d96 100644 --- a/packages/world-local/src/storage/legacy.ts +++ b/packages/world-local/src/storage/legacy.ts @@ -11,6 +11,7 @@ import { jsonReplacer, promoteExclusive, resolveWithinBase, + runEntityDir, writeExclusive, writeJSON, } from '../fs.js'; @@ -185,7 +186,7 @@ export async function handleLegacyEvent( const compositeKey = `${runId}-${eventId}`; const eventPath = resolveWithinBase( basedir, - 'events', + runEntityDir('events', runId), `${compositeKey}.json` ); if (data.eventType === 'hook_received') { diff --git a/packages/world-local/src/storage/run-retention.ts b/packages/world-local/src/storage/run-retention.ts index 3127f5c060..c440ed7e00 100644 --- a/packages/world-local/src/storage/run-retention.ts +++ b/packages/world-local/src/storage/run-retention.ts @@ -7,7 +7,13 @@ import { readRunRetention, } from '@workflow/world'; import { z } from 'zod'; -import { listJSONFiles, readJSON, taggedPath, writeJSON } from '../fs.js'; +import { + listJSONFiles, + readJSON, + runEntityDir, + taggedPath, + writeJSON, +} from '../fs.js'; import { purgeRunStreamData } from '../streamer.js'; import { ensureHookIndexes, listHookByRunMarkers } from './hook-index.js'; @@ -89,25 +95,33 @@ export async function purgeRunEntityData( tag: string | undefined ): Promise { await Promise.all([ - scrubEntityFiles(path.join(basedir, 'steps'), runId, (step) => { - step.input = undefined; - step.output = undefined; - step.error = undefined; - }), - scrubEntityFiles(path.join(basedir, 'events'), runId, (event) => { - const eventData = event.eventData; - if (!eventData || typeof eventData !== 'object') return; - for (const field of getEventDataRefFields(String(event.eventType))) { - delete (eventData as Record)[field]; + scrubEntityFiles( + path.join(basedir, runEntityDir('steps', runId)), + runId, + (step) => { + step.input = undefined; + step.output = undefined; + step.error = undefined; } - if ( - event.eventType === 'run_created' || - event.eventType === 'run_started' - ) { - delete (eventData as Record).dynamicWorkflowCode; - delete (eventData as Record).dynamicWorkflowCodeRef; + ), + scrubEntityFiles( + path.join(basedir, runEntityDir('events', runId)), + runId, + (event) => { + const eventData = event.eventData; + if (!eventData || typeof eventData !== 'object') return; + for (const field of getEventDataRefFields(String(event.eventType))) { + delete (eventData as Record)[field]; + } + if ( + event.eventType === 'run_created' || + event.eventType === 'run_started' + ) { + delete (eventData as Record).dynamicWorkflowCode; + delete (eventData as Record).dynamicWorkflowCodeRef; + } } - }), + ), scrubHookMetadata(basedir, runId), purgeRunStreamData(basedir, runId, tag), ]); diff --git a/packages/world-local/src/storage/slot-identity.test.ts b/packages/world-local/src/storage/slot-identity.test.ts index 38d8afc099..7ab33fcdc4 100644 --- a/packages/world-local/src/storage/slot-identity.test.ts +++ b/packages/world-local/src/storage/slot-identity.test.ts @@ -287,7 +287,7 @@ describe('slot event ids', () => { // Rewrite the run's log the way it would look had it been created before // slot ids existed. The scheme is pinned by what is on disk, not by a // stored flag, so this is the whole of the upgrade path. - const eventsDir = path.join(testDir, 'events'); + const eventsDir = path.join(testDir, 'events', runId); const files = (await fs.readdir(eventsDir)).filter((file) => file.startsWith(`${runId}-`) ); diff --git a/packages/world-local/src/storage/steps-storage.ts b/packages/world-local/src/storage/steps-storage.ts index 1102f4300b..e391a587b9 100644 --- a/packages/world-local/src/storage/steps-storage.ts +++ b/packages/world-local/src/storage/steps-storage.ts @@ -6,6 +6,7 @@ import { assertSafeEntityId, paginatedFileSystemQuery, readJSONWithFallback, + runEntityDir, } from '../fs.js'; import { filterStepData } from './filters.js'; import { getObjectCreatedAt } from './helpers.js'; @@ -25,7 +26,7 @@ export function createStepsStorage( const compositeKey = `${runId}-${stepId}`; const step = await readJSONWithFallback( basedir, - 'steps', + runEntityDir('steps', runId), compositeKey, StepSchema, tag @@ -41,7 +42,7 @@ export function createStepsStorage( assertSafeEntityId('runId', params.runId); const resolveData = params.resolveData ?? DEFAULT_RESOLVE_DATA_OPTION; const result = await paginatedFileSystemQuery({ - directory: path.join(basedir, 'steps'), + directory: path.join(basedir, runEntityDir('steps', params.runId)), schema: StepSchema, filePrefix: `${params.runId}-`, sortOrder: params.pagination?.sortOrder ?? 'desc', diff --git a/packages/world-local/src/tag.test.ts b/packages/world-local/src/tag.test.ts index fe898453ca..6226e80769 100644 --- a/packages/world-local/src/tag.test.ts +++ b/packages/world-local/src/tag.test.ts @@ -62,13 +62,13 @@ describe('File tagging', () => { it('should write event files with tag suffix', async () => { const storage = createStorage(testDir, 'vitest-0'); - await createRun(storage, { + const run = await createRun(storage, { deploymentId: 'dep-1', workflowName: 'test-wf', input: new Uint8Array(), }); - const eventsDir = path.join(testDir, 'events'); + const eventsDir = path.join(testDir, 'events', run.runId); const files = await fs.readdir(eventsDir); expect(files).toHaveLength(1); expect(files[0]).toMatch(/\.vitest-0\.json$/); @@ -88,7 +88,7 @@ describe('File tagging', () => { input: new Uint8Array(), }); - const stepsDir = path.join(testDir, 'steps'); + const stepsDir = path.join(testDir, 'steps', run.runId); const files = await fs.readdir(stepsDir); expect(files).toHaveLength(1); expect(files[0]).toMatch(/\.vitest-0\.json$/); @@ -460,7 +460,11 @@ describe('File tagging', () => { const runsDir = path.join(testDir, 'runs'); const eventsDir = path.join(testDir, 'events'); const stepsDir = path.join(testDir, 'steps'); - for (const dir of [runsDir, eventsDir, stepsDir]) { + for (const dir of [ + runsDir, + path.join(eventsDir, run.runId), + path.join(stepsDir, run.runId), + ]) { const files = await fs.readdir(dir); for (const file of files) { expect(file).toMatch(/\.vitest-0\.json$/); diff --git a/workbench/python/pyproject.toml b/workbench/python/pyproject.toml index c46ed90d86..fec6a10d36 100644 --- a/workbench/python/pyproject.toml +++ b/workbench/python/pyproject.toml @@ -9,7 +9,8 @@ dependencies = [ ] [tool.uv.sources] -vercel-workflow = { git = "https://github.com/vercel/vercel-py", branch = "main", subdirectory = "src/vercel-workflow" } +# TEMPORARY: matching per-run layout in vercel-py; revert to vercel/vercel-py main before merge. +vercel-workflow = { git = "https://github.com/ruiconti/vercel-py", branch = "rui/local-world-per-run-dirs-on-cc65153", subdirectory = "src/vercel-workflow" } [tool.vercel] entrypoint = "app:app" diff --git a/workbench/python/uv.lock b/workbench/python/uv.lock index 4d926d0456..3564693bd6 100644 --- a/workbench/python/uv.lock +++ b/workbench/python/uv.lock @@ -601,12 +601,12 @@ wheels = [ [[package]] name = "vercel-headers" version = "0.7.2" -source = { git = "https://github.com/vercel/vercel-py?subdirectory=src%2Fvercel-headers&branch=main#cc65153ae8417f27ca9569470ef1cb0904ad0735" } +source = { git = "https://github.com/ruiconti/vercel-py?subdirectory=src%2Fvercel-headers&branch=rui%2Flocal-world-per-run-dirs-on-cc65153#45e4a47e18a447bca29190f2e94f398f8b8e9985" } [[package]] name = "vercel-internal-core" version = "0.1.3" -source = { git = "https://github.com/vercel/vercel-py?subdirectory=src%2Fvercel-internal-core&branch=main#cc65153ae8417f27ca9569470ef1cb0904ad0735" } +source = { git = "https://github.com/ruiconti/vercel-py?subdirectory=src%2Fvercel-internal-core&branch=rui%2Flocal-world-per-run-dirs-on-cc65153#45e4a47e18a447bca29190f2e94f398f8b8e9985" } dependencies = [ { name = "anyio" }, { name = "httpx" }, @@ -616,7 +616,7 @@ dependencies = [ [[package]] name = "vercel-oidc" version = "0.8.1" -source = { git = "https://github.com/vercel/vercel-py?subdirectory=src%2Fvercel-oidc&branch=main#cc65153ae8417f27ca9569470ef1cb0904ad0735" } +source = { git = "https://github.com/ruiconti/vercel-py?subdirectory=src%2Fvercel-oidc&branch=rui%2Flocal-world-per-run-dirs-on-cc65153#45e4a47e18a447bca29190f2e94f398f8b8e9985" } dependencies = [ { name = "anyio" }, { name = "httpx" }, @@ -626,7 +626,7 @@ dependencies = [ [[package]] name = "vercel-queue" version = "0.8.1" -source = { git = "https://github.com/vercel/vercel-py?subdirectory=src%2Fvercel-queue&branch=main#cc65153ae8417f27ca9569470ef1cb0904ad0735" } +source = { git = "https://github.com/ruiconti/vercel-py?subdirectory=src%2Fvercel-queue&branch=rui%2Flocal-world-per-run-dirs-on-cc65153#45e4a47e18a447bca29190f2e94f398f8b8e9985" } dependencies = [ { name = "anyio" }, { name = "httpx", extra = ["http2"] }, @@ -640,7 +640,7 @@ dependencies = [ [[package]] name = "vercel-workflow" version = "0.10.0" -source = { git = "https://github.com/vercel/vercel-py?subdirectory=src%2Fvercel-workflow&branch=main#cc65153ae8417f27ca9569470ef1cb0904ad0735" } +source = { git = "https://github.com/ruiconti/vercel-py?subdirectory=src%2Fvercel-workflow&branch=rui%2Flocal-world-per-run-dirs-on-cc65153#45e4a47e18a447bca29190f2e94f398f8b8e9985" } dependencies = [ { name = "anyio" }, { name = "cbor2" }, @@ -666,5 +666,5 @@ dependencies = [ [package.metadata] requires-dist = [ { name = "uvicorn", specifier = ">=0.30" }, - { name = "vercel-workflow", git = "https://github.com/vercel/vercel-py?subdirectory=src%2Fvercel-workflow&branch=main" }, + { name = "vercel-workflow", git = "https://github.com/ruiconti/vercel-py?subdirectory=src%2Fvercel-workflow&branch=rui%2Flocal-world-per-run-dirs-on-cc65153" }, ]