From d8ddd062a050bf6a6c6edc8bb36bdec0fd736a51 Mon Sep 17 00:00:00 2001 From: Rui Conti Date: Tue, 6 Oct 2026 08:57:18 -0400 Subject: [PATCH 01/14] [world-local] Store events and steps in one directory per run Signed-off-by: Rui Conti --- .changeset/world-local-per-run-dirs.md | 5 + packages/world-local/src/fs.ts | 47 +++ packages/world-local/src/index.ts | 22 +- packages/world-local/src/storage.test.ts | 21 +- .../world-local/src/storage/events-storage.ts | 76 +++-- packages/world-local/src/storage/helpers.ts | 10 +- .../world-local/src/storage/hook-index.ts | 52 ++-- packages/world-local/src/storage/index.ts | 20 +- .../world-local/src/storage/layout.test.ts | 257 +++++++++++++++++ packages/world-local/src/storage/layout.ts | 269 ++++++++++++++++++ packages/world-local/src/storage/legacy.ts | 3 +- .../world-local/src/storage/run-retention.ts | 50 ++-- .../src/storage/slot-identity.test.ts | 2 +- .../world-local/src/storage/steps-storage.ts | 5 +- packages/world-local/src/tag.test.ts | 12 +- 15 files changed, 760 insertions(+), 91 deletions(-) create mode 100644 .changeset/world-local-per-run-dirs.md create mode 100644 packages/world-local/src/storage/layout.test.ts create mode 100644 packages/world-local/src/storage/layout.ts diff --git a/.changeset/world-local-per-run-dirs.md b/.changeset/world-local-per-run-dirs.md new file mode 100644 index 0000000000..6af34c5fe2 --- /dev/null +++ b/.changeset/world-local-per-run-dirs.md @@ -0,0 +1,5 @@ +--- +"@workflow/world-local": patch +--- + +Store event and step files in one directory per run (`events//`, `steps//`), so creating an event and listing a run's events or steps costs time proportional to that run instead of to every file in the data directory. Existing data directories are converted in place on first use; downgrading afterwards requires moving the files back. diff --git a/packages/world-local/src/fs.ts b/packages/world-local/src/fs.ts index af18b50b6d..eb20ecd92d 100644 --- a/packages/world-local/src/fs.ts +++ b/packages/world-local/src/fs.ts @@ -219,6 +219,31 @@ 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, so a file moved from the old flat layout keeps its name. + */ +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 @@ -575,6 +600,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..002867e282 100644 --- a/packages/world-local/src/index.ts +++ b/packages/world-local/src/index.ts @@ -11,6 +11,7 @@ import { deleteJSON, hasTag, isUntagged, + listRunScopedDirs, listTaggedFiles, listTaggedFilesByExtension, readJSON, @@ -162,11 +163,19 @@ export function createWorld(args?: Partial): LocalWorld { }) ); - // Delete tagged entity files across all directories + // Delete tagged entity files across all directories. Steps and + // events are stored one subdirectory per run. + const runScopedDirs = ( + await Promise.all([ + listRunScopedDirs(basedir, 'steps'), + listRunScopedDirs(basedir, 'events'), + ]) + ) + .flat() + .map((dir) => path.relative(basedir, dir)); const entityDirs = [ 'runs', - 'steps', - 'events', + ...runScopedDirs, 'hooks', 'hooks/by-run', 'waits', @@ -181,6 +190,13 @@ export function createWorld(args?: Partial): LocalWorld { ); }) ); + // Drop the run directories that clearing left empty. `rmdir` refuses + // a non-empty one, so another tag's (or untagged) files keep theirs. + await Promise.all( + 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); 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..109a28d0a2 100644 --- a/packages/world-local/src/storage/helpers.ts +++ b/packages/world-local/src/storage/helpers.ts @@ -13,6 +13,7 @@ import { isUntagged, readJSON, resolveWithinBase, + runEntityDir, stripTag, ulidToDate, withWindowsRetry, @@ -335,11 +336,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 +348,7 @@ export async function scanRunEventIds( ): Promise { let files: string[] = []; try { - files = await fs.readdir(path.join(basedir, 'events')); + files = await fs.readdir(path.join(basedir, runEntityDir('events', runId))); } 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..614cfdcf6e 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,32 +266,32 @@ 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 eventFiles: string[] = []; + for (const runDir of await listRunScopedDirs(basedir, 'events')) { + for (const fileId of await listJSONFiles(runDir)) { + eventFiles.push(path.join(runDir, `${fileId}.json`)); + } + } + await forEachConcurrent(eventFiles, 32, async (eventFile) => { + 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. } - ); + }); const hooksDir = path.join(basedir, 'hooks'); await forEachConcurrent(await listJSONFiles(hooksDir), 32, async (fileId) => { @@ -373,7 +375,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/index.ts b/packages/world-local/src/storage/index.ts index 6db7b29a9e..e24e264a0b 100644 --- a/packages/world-local/src/storage/index.ts +++ b/packages/world-local/src/storage/index.ts @@ -2,6 +2,7 @@ import type { Storage } from '@workflow/world'; import { instrumentObject } from '../instrumentObject.js'; import { createEventsStorage } from './events-storage.js'; import { createHooksStorage } from './hooks-storage.js'; +import { gateOnRunScopedLayout } from './layout.js'; import { createRunsStorage, type LocalRunsStorage } from './runs-storage.js'; import { createSnapshotsStorage } from './snapshots-storage.js'; import { createStepsStorage } from './steps-storage.js'; @@ -28,9 +29,22 @@ export type LocalStorage = Omit & { export function createStorage(basedir: string, tag?: string): LocalStorage { // Create raw storage implementations const runs = createRunsStorage(basedir, tag); - const steps = createStepsStorage(basedir, tag); - const events = createEventsStorage(basedir, tag); - const hooks = createHooksStorage(basedir, tag); + // Steps, events and hooks (which resolve through the event log) read the + // per-run layout, so each call first waits for any flat files left by an + // older version to be moved into it. + const steps = gateOnRunScopedLayout( + basedir, + createStepsStorage(basedir, tag) + ); + const events = gateOnRunScopedLayout( + basedir, + createEventsStorage(basedir, tag), + ['clearCache'] + ); + const hooks = gateOnRunScopedLayout( + basedir, + createHooksStorage(basedir, tag) + ); const snapshots = createSnapshotsStorage(basedir); // Instrument all storage methods with tracing diff --git a/packages/world-local/src/storage/layout.test.ts b/packages/world-local/src/storage/layout.test.ts new file mode 100644 index 0000000000..7e797ac041 --- /dev/null +++ b/packages/world-local/src/storage/layout.test.ts @@ -0,0 +1,257 @@ +import { promises as fs } from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { createStorage } from '../storage.js'; +import { + createRun, + createStep, + updateRun, + updateStep, +} from '../test-helpers.js'; +import { + migrateFlatRunScopedFiles, + resetRunScopedLayoutCache, +} from './layout.js'; + +/** + * Rewrite a run-scoped data directory into the flat layout every version + * before it used: every file of `events//` and `steps//` + * directly in `events/` and `steps/`, under the same name. + */ +async function flatten(dataDir: string): Promise { + const moved: string[] = []; + for (const entityDir of ['events', 'steps']) { + const root = path.join(dataDir, entityDir); + for (const runDir of await fs.readdir(root)) { + const full = path.join(root, runDir); + if (!(await fs.stat(full)).isDirectory()) continue; + for (const name of await fs.readdir(full)) { + await fs.rename(path.join(full, name), path.join(root, name)); + moved.push(path.join(entityDir, name)); + } + await fs.rmdir(full); + } + } + return moved.sort(); +} + +/** Every file under `dir`, relative to it. */ +async function walk(dir: string, prefix = ''): Promise { + const out: string[] = []; + for (const entry of await fs.readdir(dir, { withFileTypes: true })) { + const rel = path.join(prefix, entry.name); + if (entry.isDirectory()) { + out.push(...(await walk(path.join(dir, entry.name), rel))); + } else { + out.push(rel); + } + } + return out.sort(); +} + +async function seedRun(dataDir: string, tag?: string) { + const storage = createStorage(dataDir, tag); + const run = await createRun(storage, { + deploymentId: 'dep-1', + workflowName: 'wf', + input: new Uint8Array([1, 2, 3]), + }); + await updateRun(storage, run.runId, 'run_started'); + for (const i of [0, 1, 2]) { + const stepId = `step_${i}`; + await createStep(storage, run.runId, { + stepId, + stepName: 'my-step', + input: new Uint8Array([i]), + }); + await updateStep(storage, run.runId, stepId, 'step_started'); + await updateStep(storage, run.runId, stepId, 'step_completed', { + result: new Uint8Array([i, i]), + }); + } + return run.runId; +} + +async function snapshot(dataDir: string, runId: string, tag?: string) { + const storage = createStorage(dataDir, tag); + const events = await storage.events.list({ + runId, + pagination: { limit: 1000, sortOrder: 'asc' }, + resolveData: 'all', + }); + const steps = await storage.steps.list({ + runId, + pagination: { limit: 1000, sortOrder: 'asc' }, + resolveData: 'all', + }); + const step = await storage.steps.get(runId, 'step_1', { + resolveData: 'all', + }); + return { events: events.data, steps: steps.data, step }; +} + +describe('run-scoped layout migration', () => { + let dataDir: string; + + beforeEach(async () => { + dataDir = await fs.mkdtemp(path.join(os.tmpdir(), 'layout-test-')); + resetRunScopedLayoutCache(); + }); + + afterEach(async () => { + resetRunScopedLayoutCache(); + await fs.rm(dataDir, { recursive: true, force: true }); + }); + + it('stores events and steps in one directory per run', async () => { + const runId = await seedRun(dataDir); + const files = await walk(dataDir); + const eventFiles = files.filter((f) => f.startsWith(`events${path.sep}`)); + const stepFiles = files.filter((f) => f.startsWith(`steps${path.sep}`)); + expect(eventFiles.length).toBeGreaterThan(0); + expect(stepFiles).toHaveLength(3); + for (const f of [...eventFiles, ...stepFiles]) { + expect(f.split(path.sep)[1]).toBe(runId); + } + }); + + it('moves a flat store into per-run directories on first use, losslessly', async () => { + const runA = await seedRun(dataDir); + const runB = await seedRun(dataDir, 'vitest-0'); + const before = { + a: await snapshot(dataDir, runA), + b: await snapshot(dataDir, runB, 'vitest-0'), + }; + const flatFiles = await flatten(dataDir); + expect(flatFiles.length).toBeGreaterThan(20); + // A write that crashed before its rename: not an entity file, stays put. + const tmp = path.join(dataDir, 'events', `${runA}-evnt_x.json.tmp.01ABC`); + await fs.writeFile(tmp, '{'); + + // A fresh process: the first storage call converts the directory. + resetRunScopedLayoutCache(); + expect(await snapshot(dataDir, runA)).toEqual(before.a); + expect(await snapshot(dataDir, runB, 'vitest-0')).toEqual(before.b); + + const after = await walk(dataDir); + for (const f of flatFiles) { + const [entityDir, name] = f.split(path.sep); + const runId = name.startsWith(runA) ? runA : runB; + expect(after).toContain(path.join(entityDir, runId, name)); + expect(after).not.toContain(f); + } + await expect(fs.access(tmp)).resolves.toBeUndefined(); + + // Writes after the move land in the run's directory and read back. + const storage = createStorage(dataDir); + await updateRun(storage, runA, 'run_completed', { + output: new Uint8Array([9]), + }); + const events = await storage.events.list({ + runId: runA, + pagination: { limit: 1000 }, + }); + expect(events.data).toHaveLength(before.a.events.length + 1); + }); + + it('finishes a half-migrated store and is idempotent', async () => { + const runId = await seedRun(dataDir); + const before = await snapshot(dataDir, runId); + const flatFiles = await flatten(dataDir); + + // Simulate a pass interrupted part-way: some files already moved, one + // linked at both paths but not yet unlinked from the flat directory. + const half = flatFiles.slice(0, Math.floor(flatFiles.length / 2)); + for (const f of half) { + const [entityDir, name] = f.split(path.sep); + await fs.mkdir(path.join(dataDir, entityDir, runId), { recursive: true }); + await fs.rename( + path.join(dataDir, f), + path.join(dataDir, entityDir, runId, name) + ); + } + const linked = flatFiles[half.length]; + { + const [entityDir, name] = linked.split(path.sep); + await fs.mkdir(path.join(dataDir, entityDir, runId), { recursive: true }); + await fs.link( + path.join(dataDir, linked), + path.join(dataDir, entityDir, runId, name) + ); + } + + const first = await migrateFlatRunScopedFiles(dataDir); + expect(first).toEqual({ + moved: flatFiles.length - half.length, + skipped: 0, + }); + expect(await migrateFlatRunScopedFiles(dataDir)).toEqual({ + moved: 0, + skipped: 0, + }); + resetRunScopedLayoutCache(); + expect(await snapshot(dataDir, runId)).toEqual(before); + }); + + it('never overwrites a different file already at the destination', async () => { + const runId = await seedRun(dataDir); + const flatFiles = await flatten(dataDir); + const stepFile = flatFiles.find((f) => f.startsWith(`steps${path.sep}`))!; + const name = path.basename(stepFile); + const nested = path.join(dataDir, 'steps', runId, name); + await fs.mkdir(path.dirname(nested), { recursive: true }); + await fs.writeFile(nested, '{"newer": true}'); + + const result = await migrateFlatRunScopedFiles(dataDir); + expect(result).toEqual({ moved: flatFiles.length - 1, skipped: 1 }); + expect(await fs.readFile(nested, 'utf8')).toBe('{"newer": true}'); + // The flat copy is left alone rather than deleted. + await expect( + fs.access(path.join(dataDir, stepFile)) + ).resolves.toBeUndefined(); + }); + + it('is safe to run concurrently with itself', async () => { + const runs = [await seedRun(dataDir), await seedRun(dataDir)]; + const before = await Promise.all(runs.map((r) => snapshot(dataDir, r))); + const flatFiles = await flatten(dataDir); + + const passes = await Promise.all([ + migrateFlatRunScopedFiles(dataDir), + migrateFlatRunScopedFiles(dataDir), + migrateFlatRunScopedFiles(dataDir), + ]); + // Racing passes may each count a file they both finished moving. + expect(passes.reduce((n, p) => n + p.moved, 0)).toBeGreaterThanOrEqual( + flatFiles.length + ); + expect(passes.every((p) => p.skipped === 0)).toBe(true); + for (const entityDir of ['events', 'steps']) { + const flatLeft = ( + await fs.readdir(path.join(dataDir, entityDir), { withFileTypes: true }) + ).filter((e) => e.isFile()); + expect(flatLeft).toEqual([]); + } + + resetRunScopedLayoutCache(); + const after = await Promise.all(runs.map((r) => snapshot(dataDir, r))); + expect(after).toEqual(before); + }); + + it('routes run ids containing dashes by their last entity-id prefix', async () => { + const runId = 'wrun_custom-id-with-dashes'; + const eventsDir = path.join(dataDir, 'events'); + await fs.mkdir(eventsDir, { recursive: true }); + const name = `${runId}-evnt_000000000000000000000001.json`; + await fs.writeFile(path.join(eventsDir, name), '{}'); + + expect(await migrateFlatRunScopedFiles(dataDir)).toEqual({ + moved: 1, + skipped: 0, + }); + await expect( + fs.access(path.join(eventsDir, runId, name)) + ).resolves.toBeUndefined(); + }); +}); diff --git a/packages/world-local/src/storage/layout.ts b/packages/world-local/src/storage/layout.ts new file mode 100644 index 0000000000..8bae04e3be --- /dev/null +++ b/packages/world-local/src/storage/layout.ts @@ -0,0 +1,269 @@ +import fs from 'node:fs/promises'; +import path from 'node:path'; +import { globalSingleton } from '@workflow/utils'; +import { z } from 'zod'; +import { + RUN_SCOPED_ENTITY_DIRS, + type RunScopedEntityDir, + readJSON, + runEntityDir, + withWindowsRetry, +} from '../fs.js'; + +/** + * Data directories written before run-scoped storage kept every event and + * step file directly in `events/` and `steps/`. Each per-run read then had to + * list the whole directory, so its cost grew with every run the directory had + * ever held. This module moves those flat files into their run's + * subdirectory (`events//`, `steps//`) under their existing + * names. + * + * The move runs once per data directory per process, before the first storage + * call. It is cheap to repeat: once a directory has been converted, a pass is + * one `readdir` of `events/` and of `steps/`, which by then hold one entry + * per run. Repeating it is what makes the conversion safe to interrupt, and + * lets it pick up files an older version of this package wrote to the flat + * layout after an earlier pass. + * + * A file already at the destination is never overwritten: if it is the same + * inode (hard-linked at both paths) the flat name is dropped, otherwise the + * flat file is left in place and reported as skipped. + */ + +/** Per-entity id prefix of the second half of a file id, `${runId}-${id}`. */ +const ENTITY_ID_PREFIX: Record = { + events: '-evnt_', + steps: '-step_', +}; + +const RunIdSchema = z.object({ runId: z.string() }); + +export interface FlatLayoutMigrationResult { + /** Files moved into their run's subdirectory. */ + moved: number; + /** Flat files left in place: unparseable, or a different file was already at the destination. */ + skipped: number; +} + +/** + * The run a flat file belongs to, from its name: `${runId}-${entityId}` plus + * an optional `.${tag}` and the `.json` extension. Run ids may contain `-` + * (custom ids), so split at the last `-evnt_` / `-step_`. Falls back to the + * file's own `runId` for names that do not carry the prefix. + */ +async function runIdOfFlatFile( + entityDir: RunScopedEntityDir, + filePath: string +): Promise { + const name = path.basename(filePath); + const split = name.lastIndexOf(ENTITY_ID_PREFIX[entityDir]); + if (split > 0) { + return name.slice(0, split); + } + try { + return (await readJSON(filePath, RunIdSchema))?.runId ?? null; + } catch { + return null; + } +} + +/** + * How the flat file `from` relates to the file already at `to`: the same + * inode (hard-linked at both paths), already gone from the flat directory + * (a concurrent pass finished the move), or a different file. + */ +async function compareWithDestination( + from: string, + to: string +): Promise<'same' | 'gone' | 'different'> { + let sa: import('node:fs').Stats; + try { + sa = await fs.stat(from); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return 'gone'; + throw error; + } + const sb = await fs.stat(to); + return sa.ino === sb.ino && sa.dev === sb.dev ? 'same' : 'different'; +} + +async function exists(filePath: string): Promise { + try { + await fs.lstat(filePath); + return true; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return false; + throw error; + } +} + +async function moveFlatFile( + basedir: string, + entityDir: RunScopedEntityDir, + name: string, + madeDirs: Set +): Promise<'moved' | 'skipped' | 'gone'> { + const from = path.join(basedir, entityDir, name); + const runId = await runIdOfFlatFile(entityDir, from); + let toDir: string; + try { + if (!runId) throw new Error('no run id'); + toDir = path.join(basedir, runEntityDir(entityDir, runId)); + } catch { + return 'skipped'; + } + if (!madeDirs.has(toDir)) { + await fs.mkdir(toDir, { recursive: true }); + madeDirs.add(toDir); + } + const to = path.join(toDir, name); + // `rename` would replace a file already at the destination, so look first. + // Only a concurrent pass can put one there between this check and the + // rename, and it would be moving this very file (event files are never + // rewritten, and nothing writes a run-scoped path before the pass that + // gates it has finished). + if (await exists(to)) { + const relation = await compareWithDestination(from, to); + if (relation === 'gone') return 'gone'; + if (relation === 'different') return 'skipped'; + // Hard-linked at both paths (an interrupted copy): drop the flat name. + try { + await withWindowsRetry(() => fs.unlink(from)); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + } + return 'moved'; + } + try { + await withWindowsRetry(() => fs.rename(from, to)); + } catch (error) { + // A concurrent pass moved it first. + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return 'gone'; + throw error; + } + return 'moved'; +} + +/** + * Move every flat event and step file in `basedir` into its run's + * subdirectory. Safe to run concurrently with itself, and to interrupt. + */ +export async function migrateFlatRunScopedFiles( + basedir: string, + options: { onStart?: (fileCount: number) => void } = {} +): Promise { + const result: FlatLayoutMigrationResult = { moved: 0, skipped: 0 }; + const madeDirs = new Set(); + const pending: [RunScopedEntityDir, string[]][] = []; + for (const entityDir of RUN_SCOPED_ENTITY_DIRS) { + let entries: import('node:fs').Dirent[]; + try { + entries = await fs.readdir(path.join(basedir, entityDir), { + withFileTypes: true, + }); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') continue; + throw error; + } + // Only finished entity files. A `*.json.tmp.` is a write that has not + // been renamed into place yet (or never will be); its writer still + // expects the flat path, and a crashed one leaves debris either way. + const names = entries + .filter((entry) => entry.isFile() && entry.name.endsWith('.json')) + .map((entry) => entry.name); + if (names.length > 0) pending.push([entityDir, names]); + } + const total = pending.reduce((n, [, names]) => n + names.length, 0); + if (total > 0) options.onStart?.(total); + for (const [entityDir, names] of pending) { + const concurrency = 32; + for (let start = 0; start < names.length; start += concurrency) { + const outcomes = await Promise.all( + names + .slice(start, start + concurrency) + .map((name) => moveFlatFile(basedir, entityDir, name, madeDirs)) + ); + for (const outcome of outcomes) { + if (outcome === 'moved') result.moved++; + else if (outcome === 'skipped') result.skipped++; + } + } + } + return result; +} + +// Per-process memo, on `globalThis` so several bundled copies of this module +// share it (see `globalSingleton`). Failed passes are forgotten and retried +// by the next storage call. +const layoutState = globalSingleton( + '@workflow/world-local//runScopedLayout', + 1, + () => ({ passes: new Map>() }) +); + +/** + * Resolve once `basedir` holds no flat event or step files, converting it on + * the first call in this process. Every storage entry point awaits this, so + * nothing reads the run-scoped layout before the conversion has finished. + */ +export function ensureRunScopedLayout(basedir: string): Promise { + const key = path.resolve(basedir); + let pass = layoutState.passes.get(key); + if (!pass) { + pass = (async () => { + const startedAt = Date.now(); + const { moved, skipped } = await migrateFlatRunScopedFiles(key, { + onStart: (fileCount) => { + if (fileCount >= 1000) { + console.log( + `[world-local] Moving ${fileCount} event and step files in ` + + `${key} into per-run directories (one-time, this can take ` + + `a while)...` + ); + } + }, + }); + if (moved > 0 || skipped > 0) { + console.log( + `[world-local] Moved ${moved} event and step files into per-run ` + + `directories in ${Date.now() - startedAt}ms` + + (skipped > 0 ? ` (left ${skipped} in place)` : '') + ); + } + })().catch((error) => { + layoutState.passes.delete(key); + throw error; + }); + layoutState.passes.set(key, pass); + } + return pass; +} + +/** Forget completed passes (tests). */ +export function resetRunScopedLayoutCache(): void { + layoutState.passes.clear(); +} + +/** + * Wrap every async method of a storage object so it first awaits + * {@link ensureRunScopedLayout}. Methods listed in `syncMethods` are passed + * through untouched. + */ +export function gateOnRunScopedLayout( + basedir: string, + target: T, + syncMethods: readonly string[] = [] +): T { + const gated: Record = {}; + for (const [name, value] of Object.entries(target)) { + if (typeof value !== 'function' || syncMethods.includes(name)) { + gated[name] = value; + continue; + } + gated[name] = async (...args: unknown[]) => { + await ensureRunScopedLayout(basedir); + return (value as (...a: unknown[]) => unknown).apply(target, args); + }; + } + return gated as T; +} 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$/); From 535a79f8824d3773f34c3afde5554d1179b22c9e Mon Sep 17 00:00:00 2001 From: Rui Conti Date: Tue, 6 Oct 2026 14:09:29 -0400 Subject: [PATCH 02/14] [world-local] Route ambiguous flat files by stored run id; migrate before tagged clear Step ids may contain -step_ (e.g. step_a-step_b), so a flat file's run is taken from its name only when the separator occurs once; otherwise from the runId stored in the file. Tagged clear() now awaits the layout migration and also clears the flat directories. Document that older flat-layout writers must be stopped before upgrading, and how to roll back. Signed-off-by: Rui Conti --- packages/world-local/src/index.ts | 9 +- .../world-local/src/storage/layout.test.ts | 89 ++++++++++++++++++- packages/world-local/src/storage/layout.ts | 30 +++++-- 3 files changed, 118 insertions(+), 10 deletions(-) diff --git a/packages/world-local/src/index.ts b/packages/world-local/src/index.ts index 002867e282..3d1c157c05 100644 --- a/packages/world-local/src/index.ts +++ b/packages/world-local/src/index.ts @@ -21,6 +21,7 @@ import { instrumentObject } from './instrumentObject.js'; import { createQueue, type DirectHandler } from './queue.js'; import { hashToken, hookRecoveryMarkerPath } from './storage/helpers.js'; import { resetHookIndexEnsureCache } from './storage/hook-index.js'; +import { ensureRunScopedLayout } from './storage/layout.js'; import { createStorage } from './storage.js'; import { createStreamer } from './streamer.js'; @@ -164,7 +165,11 @@ export function createWorld(args?: Partial): LocalWorld { ); // Delete tagged entity files across all directories. Steps and - // events are stored one subdirectory per run. + // events are stored one subdirectory per run; finish converting a + // flat store first, so this pass sees every one of them. The flat + // `events/` and `steps/` directories are cleared too, for files the + // conversion left in place. + await ensureRunScopedLayout(basedir); const runScopedDirs = ( await Promise.all([ listRunScopedDirs(basedir, 'steps'), @@ -175,6 +180,8 @@ export function createWorld(args?: Partial): LocalWorld { .map((dir) => path.relative(basedir, dir)); const entityDirs = [ 'runs', + 'steps', + 'events', ...runScopedDirs, 'hooks', 'hooks/by-run', diff --git a/packages/world-local/src/storage/layout.test.ts b/packages/world-local/src/storage/layout.test.ts index 7e797ac041..5316319d7f 100644 --- a/packages/world-local/src/storage/layout.test.ts +++ b/packages/world-local/src/storage/layout.test.ts @@ -239,7 +239,7 @@ describe('run-scoped layout migration', () => { expect(after).toEqual(before); }); - it('routes run ids containing dashes by their last entity-id prefix', async () => { + it('routes run ids containing dashes by their entity-id prefix', async () => { const runId = 'wrun_custom-id-with-dashes'; const eventsDir = path.join(dataDir, 'events'); await fs.mkdir(eventsDir, { recursive: true }); @@ -254,4 +254,91 @@ describe('run-scoped layout migration', () => { fs.access(path.join(eventsDir, runId, name)) ).resolves.toBeUndefined(); }); + + it('routes step ids containing the separator by the stored run id', async () => { + const storage = createStorage(dataDir); + const run = await createRun(storage, { + deploymentId: 'dep-1', + workflowName: 'wf', + input: new Uint8Array([1]), + }); + await updateRun(storage, run.runId, 'run_started'); + const stepId = 'step_a-step_b'; + await createStep(storage, run.runId, { + stepId, + stepName: 'my-step', + input: new Uint8Array([7]), + }); + const before = await storage.steps.get(run.runId, stepId, { + resolveData: 'all', + }); + + await flatten(dataDir); + resetRunScopedLayoutCache(); + const result = await migrateFlatRunScopedFiles(dataDir); + expect(result.skipped).toBe(0); + await expect( + fs.access( + path.join(dataDir, 'steps', run.runId, `${run.runId}-${stepId}.json`) + ) + ).resolves.toBeUndefined(); + const stepDirs = await fs.readdir(path.join(dataDir, 'steps')); + expect(stepDirs).toEqual([run.runId]); + + resetRunScopedLayoutCache(); + const after = await createStorage(dataDir).steps.get(run.runId, stepId, { + resolveData: 'all', + }); + expect(after).toEqual(before); + }); + + it('leaves an ambiguous file whose stored run id does not match its name', async () => { + const stepsDir = path.join(dataDir, 'steps'); + await fs.mkdir(stepsDir, { recursive: true }); + const name = 'wrun_x-step_a-step_b.json'; + await fs.writeFile( + path.join(stepsDir, name), + JSON.stringify({ runId: 'wrun_other' }) + ); + expect(await migrateFlatRunScopedFiles(dataDir)).toEqual({ + moved: 0, + skipped: 1, + }); + await expect(fs.access(path.join(stepsDir, name))).resolves.toBeUndefined(); + }); + + it('tagged clear() on an unmigrated flat store removes only that tag', async () => { + const { createWorld } = await import('../index.js'); + const keepUntagged = await seedRun(dataDir); + const keepOther = await seedRun(dataDir, 'vitest-1'); + const drop = await seedRun(dataDir, 'vitest-0'); + const before = { + untagged: await snapshot(dataDir, keepUntagged), + other: await snapshot(dataDir, keepOther, 'vitest-1'), + }; + await flatten(dataDir); + resetRunScopedLayoutCache(); + + // clear() is the first call in this process: nothing has migrated yet. + const world = createWorld({ dataDir, tag: 'vitest-0' }); + await world.clear(); + + const files = await walk(dataDir); + expect(files.filter((f) => f.includes('vitest-0'))).toEqual([]); + expect(files.filter((f) => f.includes(drop))).toEqual([]); + // Every remaining event and step file sits in its run's directory. + for (const f of files) { + const [entityDir, ...rest] = f.split(path.sep); + if (entityDir === 'events' || entityDir === 'steps') { + expect(rest).toHaveLength(2); + } + } + + resetRunScopedLayoutCache(); + expect({ + untagged: await snapshot(dataDir, keepUntagged), + other: await snapshot(dataDir, keepOther, 'vitest-1'), + }).toEqual(before); + await world.close?.(); + }); }); diff --git a/packages/world-local/src/storage/layout.ts b/packages/world-local/src/storage/layout.ts index 8bae04e3be..1cb53232c8 100644 --- a/packages/world-local/src/storage/layout.ts +++ b/packages/world-local/src/storage/layout.ts @@ -28,6 +28,14 @@ import { * A file already at the destination is never overwritten: if it is the same * inode (hard-linked at both paths) the flat name is dropped, otherwise the * flat file is left in place and reported as skipped. + * + * Stop every process running an older version of this package on the data + * directory before upgrading. The pass runs once per process, so a flat file + * an older writer adds afterwards stays invisible to this process (its run + * looks truncated) until the next process start moves it. There is no + * automatic downgrade: older versions only read the flat layout, so rolling + * back means moving each `events//` and `steps//` file back up + * one level, under the same name, with every process stopped. */ /** Per-entity id prefix of the second half of a file id, `${runId}-${id}`. */ @@ -46,25 +54,31 @@ export interface FlatLayoutMigrationResult { } /** - * The run a flat file belongs to, from its name: `${runId}-${entityId}` plus - * an optional `.${tag}` and the `.json` extension. Run ids may contain `-` - * (custom ids), so split at the last `-evnt_` / `-step_`. Falls back to the - * file's own `runId` for names that do not carry the prefix. + * The run a flat file belongs to. Its name is `${runId}-${entityId}` plus an + * optional `.${tag}` and the `.json` extension, but both halves may contain + * `-` (custom run ids, and step ids such as `step_a-step_b`), so the name + * alone is only conclusive when the `-evnt_` / `-step_` separator occurs + * exactly once. Otherwise the run id stored in the file decides, provided the + * name really starts with it; anything else is left in place. */ async function runIdOfFlatFile( entityDir: RunScopedEntityDir, filePath: string ): Promise { const name = path.basename(filePath); - const split = name.lastIndexOf(ENTITY_ID_PREFIX[entityDir]); - if (split > 0) { - return name.slice(0, split); + const separator = ENTITY_ID_PREFIX[entityDir]; + const first = name.indexOf(separator); + if (first > 0 && name.indexOf(separator, first + 1) === -1) { + return name.slice(0, first); } + let stored: string | undefined; try { - return (await readJSON(filePath, RunIdSchema))?.runId ?? null; + stored = (await readJSON(filePath, RunIdSchema))?.runId; } catch { return null; } + if (!stored || !name.startsWith(`${stored}-`)) return null; + return stored; } /** From dae30333b99d9c68035b512eb2172c2b67ac173c Mon Sep 17 00:00:00 2001 From: Rui Conti Date: Tue, 6 Oct 2026 15:00:19 -0400 Subject: [PATCH 03/14] [world-local] Add layout benchmark and flat-layout rollback script benchmark-layout.mjs measures event create/list latency on a padded copy of a real data directory; prepare refuses existing or overlapping destinations before writing. flatten-layout.mjs moves per-run files back to the flat layout for downgrades and is published with the package. Tests cover both, and the README documents the layout, migration and rollback. Signed-off-by: Rui Conti --- .changeset/world-local-per-run-dirs.md | 2 +- packages/world-local/README.md | 33 ++ packages/world-local/package.json | 5 +- .../world-local/scripts/benchmark-layout.mjs | 341 ++++++++++++++++++ .../world-local/scripts/flatten-layout.mjs | 50 +++ .../scripts/layout-scripts.test.ts | 192 ++++++++++ 6 files changed, 620 insertions(+), 3 deletions(-) create mode 100644 packages/world-local/scripts/benchmark-layout.mjs create mode 100644 packages/world-local/scripts/flatten-layout.mjs create mode 100644 packages/world-local/scripts/layout-scripts.test.ts diff --git a/.changeset/world-local-per-run-dirs.md b/.changeset/world-local-per-run-dirs.md index 6af34c5fe2..d5342016dc 100644 --- a/.changeset/world-local-per-run-dirs.md +++ b/.changeset/world-local-per-run-dirs.md @@ -2,4 +2,4 @@ "@workflow/world-local": patch --- -Store event and step files in one directory per run (`events//`, `steps//`), so creating an event and listing a run's events or steps costs time proportional to that run instead of to every file in the data directory. Existing data directories are converted in place on first use; downgrading afterwards requires moving the files back. +Store event and step files in one directory per run (`events//`, `steps//`), so creating an event and listing a run's events or steps costs time proportional to that run instead of to every file in the data directory. Existing data directories are converted in place on first use; stop processes running older versions on the same data directory first. Downgrading afterwards requires moving the files back with `scripts/flatten-layout.mjs`. diff --git a/packages/world-local/README.md b/packages/world-local/README.md index 09e33326c2..7e2efed585 100644 --- a/packages/world-local/README.md +++ b/packages/world-local/README.md @@ -30,3 +30,36 @@ const world = createWorld({ dataDir: './custom-workflow-data', }); ``` + +## Data directory layout + +Event and step files are stored in one directory per run: +`events//-.json` and +`steps//-.json`. Reading or appending to one run lists +only that run's directory, so the cost stays proportional to the run instead +of to every run the data directory has ever held. + +Data directories written by 5.0.1 and earlier keep every file directly in +`events/` and `steps/`. They are converted on first use: before the first +storage call in a process, each flat file is renamed into its run's +directory. Renames never overwrite, so the conversion is safe to interrupt +and to run from several processes at once. It takes roughly 0.35 ms per file +on APFS (about 20 s for 60k files) and logs a notice from 1,000 files up. + +- **Stop older writers first.** The conversion runs once per process. Files + that an older version keeps writing to the flat layout afterwards are only + picked up on the next process start. +- **Downgrading** requires moving the files back. Stop every process using the + data directory, then run + `node node_modules/@workflow/world-local/scripts/flatten-layout.mjs ` + (or `scripts/flatten-layout.mjs` from this repository). It never + overwrites a file already at the flat path. + +To compare layouts on a copy of a real data directory: + +```sh +# Copy-on-write clone where supported, padded to 50k event files. +# The destination must be a new path outside the source. +node scripts/benchmark-layout.mjs prepare 50000 +node scripts/benchmark-layout.mjs run /dist/index.js