Skip to content
5 changes: 5 additions & 0 deletions .changeset/world-local-per-run-dirs.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@workflow/world-local": patch
---

Store event and step files in one directory per run for faster reads. Local run data written by earlier releases is deleted on upgrade.
4 changes: 4 additions & 0 deletions docs/content/worlds/v5/local.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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/<runId>/`, `steps/<runId>/`), 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.

## Limitations

The Local World is designed for development, not production:
Expand Down
50 changes: 50 additions & 0 deletions packages/world-local/src/fs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,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 `<entityDir>/<runId>/`. 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
Expand Down Expand Up @@ -575,6 +603,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<string[]> {
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<T> {
directory: string;
schema: z.ZodType<T>;
Expand Down
21 changes: 18 additions & 3 deletions packages/world-local/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -162,11 +162,19 @@ export function createWorld(args?: Partial<Config>): 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/*.<tag>.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',
Expand All @@ -181,6 +189,13 @@ export function createWorld(args?: Partial<Config>): 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);
Expand Down
65 changes: 65 additions & 0 deletions packages/world-local/src/init.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -443,6 +443,71 @@ 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(() => {});
vi.spyOn(console, 'log').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(() => {});
vi.spyOn(console, 'log').mockImplementation(() => {});

await initDataDir(dataDir);

expect(warn).not.toHaveBeenCalled();
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);
});
});

describe('DataDirVersionError', () => {
it('should store version information', () => {
const oldVersion = parseVersion('1.0.0');
Expand Down
42 changes: 42 additions & 0 deletions packages/world-local/src/init.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -306,6 +310,29 @@ async function writeVersionFile(
await writeFile(versionFilePath, content);
}

/**
* Whether the data directory was written by a release that kept every event
* and step file directly in `events/` and `steps/`, rather than one
* subdirectory per run. Stops at the first such file, so a large legacy
* directory costs one partial listing.
*/
async function hasFlatLayout(dataDir: string): Promise<boolean> {
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;
}
// Leaving the loop early closes `dir`.
for await (const entry of dir) {
if (entry.isFile() && entry.name.endsWith('.json')) return true;
}
}
return false;
}

/**
* Gets the suggested downgrade version based on the old version.
* If a specific version is suggested in the error, use that.
Expand Down Expand Up @@ -335,6 +362,21 @@ export async function initDataDir(dataDir: string): Promise<void> {
// First ensure the directory exists and is accessible
await ensureDataDir(dataDir);

// Local run data from the old flat layout is not migrated: wipe it and
// start over as a new data directory.
if (await hasFlatLayout(dataDir)) {
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);

Expand Down
21 changes: 17 additions & 4 deletions packages/world-local/src/storage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -730,6 +730,7 @@ describe('Storage', () => {
const filePath = path.join(
testDir,
'steps',
testRunId,
`${testRunId}-step_123.json`
);
const fileExists = await fs
Expand Down Expand Up @@ -1318,6 +1319,7 @@ describe('Storage', () => {
const filePath = path.join(
testDir,
'events',
testRunId,
`${testRunId}-${event.eventId}.json`
);
const fileExists = await fs
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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`
),
'{'
);

Expand Down Expand Up @@ -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`),
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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 {
Expand Down
Loading
Loading