Repository navigation
Conversation
Signed-off-by: Rui Conti <ruiconti@gmail.com>
…fore 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 <ruiconti@gmail.com>
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 <ruiconti@gmail.com>
🦋 Changeset detectedLatest commit: f0745a8 The changes in this PR will be included in the next version bump. This PR includes changesets to release 19 packages
Not sure what this means? Click here to learn what changesets are. Click here if you're a maintainer who wants to add another changeset to this PR |
|
@ruiconti is attempting to deploy a commit to the Vercel Labs Team on Vercel. A member of the Team first needs to authorize it. |
| 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 |
There was a problem hiding this comment.
package.json is already at 5.0.2, and the 5.0.2 release (top of CHANGELOG.md) still writes the flat layout. So "5.0.1 and earlier" leaves out the version most people are on right now, and the downgrade note sends 5.0.2 users the wrong way. The same text is at scripts/flatten-layout.mjs:2. Something like "5.0.2 and earlier" (or "versions before run-scoped storage", so it doesn't go stale) would cover it.
| mode: fs.constants.COPYFILE_FICLONE, | ||
| }); | ||
| const dir = path.join(dst, 'events'); | ||
| const files = fs.readdirSync(dir).filter((f) => f.endsWith('.json')); |
There was a problem hiding this comment.
prepare only picks up .json files sitting directly in events/. Once a data directory has been opened by this PR's build (or any later one), every entry there is a <runId>/ directory, so files comes back empty. The padding loop then reads real[pad % 0], which is undefined, and path.join throws. By that point cpSync has already created dst, so retrying with the same destination gets refused with "already exists". The eventFiles() helper above already handles both layouts. Could prepare flatten (or enumerate) the copied run directories before padding, and fail before copying when there's nothing to pad from? layout-scripts.test.ts only seeds a flat store, so a run-scoped source case would catch this.
A related bug in run at lines 327-329: in the new layout, steps/ contains the <runId> directory itself. It passes f.startsWith(runId), so its directory stat size gets added to newRunBytes, on top of the per-file sizes counted from stepDir.
VaguelySerious
left a comment
There was a problem hiding this comment.
AI review: blocking issues found
| * 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<void> { |
There was a problem hiding this comment.
AI Review: Blocking
Any storage call from a newer process converts the whole store, including calls from read-only tools, and an older process that is still running loses every in-flight run's event log.
The gate wraps every steps/events/hooks method. So workflow inspect events, the web UI, or a vitest worker moves every flat file the first time it reads anything. (Vitest workers share .workflow-data with the dev server and call clear() first.) The CLI ships its own copy of @workflow/world-local, and world-local isn't in the fixed changeset group, so a newer CLI running next to an older dev server isn't exotic. A common way to hit it: run pnpm up with the dev server still running (externalized server packages stay loaded until restart), then open npx workflow web.
Repro: I built main's world-local and this branch side by side. An older process creates a run and starts a step. A separate process on this branch runs one events.list. Back in the older process:
[old server] before: run_created,run_started,step_created,step_started
[world-local] Moved 5 event and step files into per-run directories in 3ms
[new reader] saw 4 events
[old server] replay after new read: []
[old server] step_completed threw: WorkflowWorldError Step "step_…" not found
[old server] step_created -> pending
[old server] replay now: ["step_created"]
The next newer process moves that orphaned step_created into the run's directory. The run's log then ends with a step_created that was written against an empty replay.
The README and changeset say "stop older writers first". But a user reading dev-server data with the CLI doesn't think of it as a writer, and nothing warns at runtime. Options:
- Convert only from the process that owns the store (
world.start(), which already doesinitDataDir/version handling), and record completion in a layout marker. Processes that find no marker read the flat layout as before. The CLI and web UI then stay non-destructive; a slow read beats a corrupted run. - At minimum, detect an active older writer (new flat
*.jsonshowing up after the pass) and log loudly instead of silently splitting the log.
| ); | ||
| } | ||
| })().catch((error) => { | ||
| layoutState.passes.delete(key); |
There was a problem hiding this comment.
AI Review: Note
When the pass throws, every events/steps/hooks call in the process fails, and each call retries the full flat readdir. mkdir/rename/stat errors aren't handled per file, so a single EACCES/EROFS/EPERM rejects the whole pass.
On main, a store the process can read but not write still reads fine: a root-owned .workflow-data from a container, a read-only mount, or a copied CI artifact opened with workflow inspect/web. On this branch every event/step read fails. I simulated it by making mkdir/rename throw EROFS (this sandbox doesn't enforce permissions):
old events.list ok 9 (x3)
new events.list THREW EROFS | EROFS: read-only file system, mkdir '…/events/wrun_…' (x3)
new readdir(events/) calls: 3
Suggest: on permission/read-only errors, log once and fall back to flat reads for that process rather than failing every storage call.
| // `events/` and `steps/` directories are cleared too, for files the | ||
| // conversion left in place. | ||
| await ensureRunScopedLayout(basedir); | ||
| const runScopedDirs = ( |
There was a problem hiding this comment.
AI Review: Note
Tagged clear() now does one readdir (listTaggedFiles) and one rmdir for every run directory in the store, including runs that belong to the untagged dev server. @workflow/vitest calls this in every worker's setup against the shared .workflow-data.
Synthetic store with 10k runs × 7 files: main takes ~35 ms per call. This branch takes ~600 ms per call (steady state, after migration), per worker.
Fixes: derive the tag's own runs from runs/*.<tag>.json (one readdir) and walk only those run directories, and rmdir only directories something was actually deleted from.
| 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 <dataDir>` |
There was a problem hiding this comment.
AI Review: Note
node node_modules/@workflow/world-local/scripts/flatten-layout.mjs only works when @workflow/world-local is hoisted to the project root. Most apps depend on workflow, which reaches world-local through @workflow/core, so under pnpm (or any non-hoisting install) that path doesn't exist.
The script also has to run before the package is downgraded, since the older version doesn't ship it. The README doesn't say so.
A downgrade without the script is silent. The older version sees every run with an empty event log and re-enqueues the active ones on start(). Its next write reuses slot 1, which this branch then leaves in place, with a log line on every start:
[old after migrate] events of runA: 0 | run status: running
[world-local] Moved 0 event and step files into per-run directories in 3ms (left 1 in place)
Consider exposing the rollback through the CLI (or a package bin) and stating the order explicitly. Also, the user-facing docs (docs/content/worlds/v5/local.mdx) don't mention the layout change or the upgrade/downgrade procedure; only this package README does, and most users never see it.
| @@ -0,0 +1,5 @@ | |||
| --- | |||
| "@workflow/world-local": patch | |||
There was a problem hiding this comment.
AI Review: Note
This is an on-disk format change that older versions can't read, with no automatic way back. As a patch it reaches every workflow user through @workflow/core/@workflow/cli on a routine update, often with a dev server still running (see the other comments). A minor with an explicit upgrade note in the release notes fits better.
| } | ||
| }); | ||
|
|
||
| it('moves a flat store into per-run directories on first use, losslessly', async () => { |
There was a problem hiding this comment.
AI Review: Note
Every migration test builds its "flat" store by flattening one this branch wrote, so they can't catch a difference in what main actually writes.
I ran a cross-version check locally (not committed):
- Setup:
main's build seeded an in-flight run (step running, hook pending, wait pending), a completed run, and avitest-0-tagged run. This branch then read and continued it. - Reads:
events.list,steps.list,listByCorrelationId, andhooks.getByTokenreturned identical results. A cursor from before the upgrade still paged correctly. - Writes:
sinceCursorreturned the correct delta, and the run completed with unique, monotonic slots. - Tags and rollback: tag isolation and tagged
clear()held. Afterflatten-layout.mjs,main's build read everything back.
So the stopped-server upgrade path is fine. A checked-in fixture (a small data directory written by the previous release) would guard it, and would be the natural home for a test of the mixed-version case.
| }, | ||
| }); | ||
| if (moved > 0 || skipped > 0) { | ||
| console.log( |
There was a problem hiding this comment.
AI Review: Nit
Once any file is skipped, this logs Moved 0 event and step files … (left 1 in place) on every process start, with no path and no hint about what to do. Listing the skipped paths (or the first few) and the reason (conflicting destination vs. unparseable) would make it actionable.
| const event = await readEventLenient( | ||
| path.join(eventsDir, `${fileId}.json`) | ||
| const eventFiles: string[] = []; | ||
| for (const runDir of await listRunScopedDirs(basedir, 'events')) { |
There was a problem hiding this comment.
AI Review: Nit
The backfill now calls readdir on each run directory one at a time (for … await), so a store without the index marker pays one serial readdir per run. Running forEachConcurrent over the run directories keeps it parallel. It's a one-time cost.
Existing data directories keep the flat layout and are read and written in place; nothing is moved implicitly. A layout.json marker (written via temp file + rename, schema-validated) selects the run-scoped layout, and only the initializer of a new, empty directory writes it. Conversion in either direction (workflow-local-layout migrate|flatten, or migrateLayout: true) takes a cross-process lock, refuses while another live process has the store open, and publishes a migrating/flattening state that closes the store until the conversion completes. Identical leftovers and hard links are dropped; differing files are reported, keep the store closed, and make the command exit nonzero. Also: diagnostics go to stderr, tagged clear() no longer touches other runs, the hook-index backfill is incremental, and the benchmark accepts both layouts. Fixtures written by 4.1.0, 5.0.0-beta.34, 5.0.0-beta.48 and 5.0.2 plus a mixed-version test against the published 5.0.2 cover compatibility. Signed-off-by: Rui Conti <ruiconti@gmail.com>
…ess and across lock recovery - convertLayout() counts the calling process's own registration, and refuses while one of its storage calls is in flight, instead of moving files under that process's pinned layout. Owner-side conversion stays with migrateLayout, which drains and reopens. - Storage calls (and clear()) are admitted and counted in the same synchronous step that checks the conversion gate, and open the store only after admission; an in-process conversion reads the count and closes the gate in one step, and conversions of one process run one at a time. Storage objects drop caches built for the previous layout. - Stale conversion locks are broken only under a second exclusive directory, after re-reading that the owner (pid, host, random token) is still the dead one observed; release removes the lock only while it still carries its own token. Signed-off-by: Rui Conti <ruiconti@gmail.com>
…ses exclusion A failed rename no longer short-circuits its batch: moves run through Promise.allSettled, so the conversion lock is released (and an owner's storage gate reopened) only once every started move has finished. Filesystem errors are reported against the file they hit, which stays in place; the transitional marker remains and the command exits 1. Anything else is rethrown after its batch settles, with no further batch started. Signed-off-by: Rui Conti <ruiconti@gmail.com>
|
Thanks for the thorough review. The latest push reworks the migration contract; point by point: layout.ts:223 (blocking): any storage call converts the store; an old live process loses its history. layout.ts:248: a failed pass breaks every call and retries the full readdir; EROFS/EACCES reject everything. index.ts:173: tagged README.md:54: changeset: patch → minor. layout.test.ts:119: fixtures are only this branch's output moved upward. layout.ts:241 (nit): "left 1 in place" with no path. hook-index.ts:270 (nit): serial readdir per run. Bot comments. The README version wording now refers to "releases before run-scoped storage" rather than 5.0.1. Python conformance failures. Migration diagnostics were going to stdout and corrupted the CLI's JSON output ( |
|
Review the following changes in direct dependencies. Learn more about Socket for GitHub.
|
Summary
@workflow/world-localkeeps every event and step file of every run in two flat directories,events/andsteps/. Every read scoped to one run lists the whole directory and filters by file-name prefix, so its cost grows with every run the data directory has ever held, not with the run being read.On a long-lived local data directory (~57k event files) this dominated the runtime: one
events.listtook 918 ms, and libuv's four threadpool threads sat near 99% inscandir, holding the worker process at ~350% CPU.This PR stores each run's files in its own subdirectory,
events/<runId>/andsteps/<runId>/, under unchanged file names. Run-scoped reads list only that run. New data directories use this layout. Existing flat data directories keep working as they are and are converted only by an explicit, offline command.Where the cost was
Each path below did one
readdirof the full flat directory per call:events.list({ runId })(paginatedFileSystemQuery)events.create(..., { sinceCursor }), which returns the log deltaevents.createper run per process (latest-seq scan)events.listByCorrelationId, hook lookupssteps.list({ runId })The per-run latest-seq memo already makes plain creates cheap; the
sinceCursorcreates and the listings were still whole-store scans.Design
events/<runId>/<runId>-<eventId>[.<tag>].jsonandsteps/<runId>/<runId>-<stepId>[.<tag>].json. File names are unchanged, so tag filtering, cursors and cache keys behave as before. This mirrors how world-postgres and world-vercel scope storage by run.<dataDir>/layout.json(written via temp file +link, schema-validated) selects the run-scoped layout. No marker means flat. Only the owner initializing a new, empty data directory writes it. Reads never create the directory or the marker, and a malformed or newer-schema marker is rejected instead of guessed at. An existing data directory, including an empty one with aversion.txt, stays flat and is read and written in place, including on read-only mounts.npx -p @workflow/world-local workflow-local-layout migrate <dataDir>(orflatten, orstatus), ormigrateLayout: true/WORKFLOW_LOCAL_MIGRATE_LAYOUT=1on the owning process..layout/holders/before reading the marker. A conversion takes an exclusive lock, publishes amigrating/flatteningmarker, then refuses if any live holder exists, including the calling process for the programmaticconvertLayout(). A process opening the store during or after an interrupted conversion is refused until the same command is run again. So no reader sees files mid-move.migrateLayoutconverts from inside the owner: it closes that process's storage gate, waits for admitted calls (includingclear()) to finish, converts, and reopens in the new layout with layout-specific caches dropped.--quarantinemoves such files to.layout/quarantine/.-evnt_/-step_separator occurs exactly once. Otherwise therunIdstored in the file decides, and only if the file name starts with it.migrateabsorbs them.clear()removes only that tag's files and never converts. On a run-scoped store it walks only that tag's runs.Benchmark
packages/world-local/scripts/benchmark-layout.mjs:Data: a copy of a long-running local app's data directory (206 real runs, 9.4k steps) padded to 50k event files with cloned files under synthetic run ids. Apple Silicon, APFS. Median of three runs' p50, in ms (range in brackets). Listings read all pages with
resolveData: 'all'.mainevents.createwithsinceCursor(61 writes/run)events.createrun_createdevents.createrun_completedevents.createstep events, no cursorevents.list, 2-event runevents.list, 117-event runevents.list, 1,672-event runThe remaining cost is proportional to the run being read. The 1,672-event listing is dominated by reading and parsing its own files (multi-MB step inputs), not by scanning, so it barely moves.
On a quieter machine against 5.0.1 the same harness measured
sinceCursorcreates at 79.0 → 2.6 ms and the 2-event listing at 98.4 → 0.3 ms.One-time conversion: ~22.6 s for ~59k files (≈0.35 ms per APFS rename). A flat store with 1,000+ event files logs a hint pointing at the command.
Upgrading and rolling back
This is a
minorrelease. Nothing changes for an existing data directory until it is converted. Releases before this one can't read a run-scoped directory (new ones included), so before downgrading, stop every process on it and runworkflow-local-layout flatten <dataDir>with this release installed. It exits nonzero and keeps the store closed while any file can't be placed.Tests
storage/layout.test.ts: layout selection (new, empty-legacy, nonexistent, competing initializers, malformed/newer markers); flat stores read and written in place, including read-only; migrate and flatten round trips; refusal with live holders (other processes and the caller); interrupted and resumed conversions; identical, hard-linked and conflicting leftovers; unroutable files; lock serialization and stale-lock recovery by two contenders; in-process admission and draining, includingclear(); failed moves settling before the lock is released; CLI exit codes and stderr-only diagnostics.storage/legacy-fixtures.test.ts: data directories written by the published 4.1.0, 5.0.0-beta.34, 5.0.0-beta.48 and 5.0.2 packages (generators committed), read in place and converted both ways: replay, tags, binary stream chunks, hooks, waits, pagination cursors.storage/mixed-version.test.ts: the published 5.0.2 as a live writer next to this reader.scripts/layout-scripts.test.ts:benchmark-layout.mjs preparesafety (same-path, ancestor, descendant including<src>/..bench, symlink-aliased and existing destinations), flat and run-scoped sources.pnpm testinpackages/world-local: 693 tests, plus 678 rerun with the flat layout forced (vitest.flat.config.ts).tsc --noEmitclean.Not in this PR
steps/and to itsstep_createdevent (175 KB–2 MB each on the store above). world-postgres duplicates it too, so this would be a local disk saving only.