Files
Nathan Rajlich f2a7bdeb0a fix(world-local,world-postgres): make duplicate hook_created idempotent (#2295)
* fix(world-local): make duplicate hook_created idempotent

Duplicate processing of the same hook_created — same runId, hookId, and
token, e.g. cross-process replay or queue redelivery — was being recorded
as a hook_conflict in the event log, which then replayed as a self-
conflict HookConflictError.

The fix mirrors the existing step_created duplicate-correlation path:
when the exclusive token claim fails and the existing claim has the same
(runId, hookId), throw EntityConflictError so the runtime's existing
concurrent-replay catch path swallows it. Different runId or hookId
reusing the same token still produces a real hook_conflict.

The persisted token claim already carried hookId; only the read schema
was dropping it. The schema now preserves hookId (marked optional for
backward compatibility with older claim files).

Fixes #2283

* fix(world-postgres): make duplicate hook_created idempotent

world-postgres has the same gap as world-local was just fixed for: the
duplicate-token check in events.create unconditionally writes a
hook_conflict event when an existing hook with the same token is found,
even when the existing hook has the same (runId, hookId) as the
incoming event. The unique partial index on workflow_events does not
catch this because the duplicate path inserts hook_conflict, not
hook_created.

Mirror the world-local fix: when the existing hook's (runId, hookId)
matches the incoming event, throw EntityConflictError so the runtime's
existing concurrent-replay catch path swallows it. Different runId or
hookId reusing the same token still produces a real hook_conflict.

Refs #2283

* test(e2e): add regression test for hook_conflict from same-tick replay race

Regression test for #1665 / #2283. A parent workflow awaits 6 child
workflows with Promise.all; each child does a tiny step and creates one
webhook. Awaited children flatten into the parent run, so all webhook
creations land on the same workflow body. When their step resolutions
align in the same tick the workflow body is re-walked and each pass
submits hook_created with the same deterministic (correlationId, token).

Before the world-side idempotency fix, the world wrote hook_conflict
events for the duplicates and the workflow failed with
HookConflictError. With the fix, duplicates throw EntityConflictError
(swallowed by the suspension handler), no hook_conflict events appear in
the log, and the webhooks resolve normally.

Verified locally against world-local: the test fails reliably (3/3) on
the unfixed code and passes reliably (5/5) on the fixed code.

* test(e2e): rewrite parallelStepsThenWebhookWorkflow to match the actual #1665 repro

The earlier version invoked another 'use workflow' function directly
from inside the parent workflow, which is not a valid child-workflow
invocation (child workflows must be spawned via start()) and didn't
mirror the bug shape on #1665 anyway.

Rewrite the workflow as a single 'use workflow' function that exactly
mirrors Paolo's minimal repro:

  await Promise.all([stepA(), stepB()]);
  using webhook = createWebhook();
  await webhook;

The for-loop runs N independent iterations of that sequence in series,
each disposing its webhook via 'using' before the next, to give the
timing-sensitive race multiple chances to fire.

The race is hard to force deterministically on fast local dev — but
the same (runId, hookId) idempotency invariant is covered
deterministically by the new unit tests in world-local and
world-postgres. This e2e test serves as a higher-level regression net:
its assertions (no hook_conflict event in the log, no
HookConflictError-failed run) are correct whether the race fires or
not, and will catch any future regression on a run that does hit
it.

* fix(world-local,world-postgres): recover crash-orphaned hook claims/rows instead of suppressing the retry

Addresses review feedback on PR #2295.

The original idempotency fix made duplicate same-(runId, hookId)
hook_created submissions throw EntityConflictError so the suspension
handler's concurrent-replay catch path swallows them. But the claim
file (world-local) and hook row (world-postgres) are written before
the durable hook_created event, and the writes are not atomic. A
process / DB interruption between the claim/hook write and the event
write leaves an orphaned claim/hook row; the retry then matched the
same (runId, hookId), threw EntityConflictError, got swallowed, and
the run was permanently left with no hook_created event in the log.

world-local:

- Add a per-(runId, hookId) in-process mutex (withHookLock) mirroring
  the existing withStepLock, so two same-tick concurrent calls
  serialize on the entity write and the dedup branch never observes
  an in-flight winner mid-write.
- In the dedup branch, when the existing claim is for the same
  (runId, hookId) we are trying to create, check whether the durable
  hook entity actually exists on disk:
    - exists  → real duplicate: throw EntityConflictError as before.
    - missing → orphaned claim from a prior crash: fall through and
      complete the partial write (write the hook entity with
      overwrite, then emit hook_created via the outer code path).

world-postgres:

- In the dedup branch, when the existing hook row matches the
  incoming (runId, hookId), check whether a hook_created event for
  this (runId, correlationId) already exists in the event log:
    - exists  → real duplicate: throw EntityConflictError as before.
    - missing → orphaned hook row from a prior crash between hook
      INSERT and events INSERT: skip the hook insert (the row is
      already there) and let the outer code path emit hook_created,
      completing the partial write.

Tests:

- world-local: pre-seed an orphaned token claim with no matching hook
  entity, retry hook_created, assert hook entity and hook_created
  event both land (no hook_conflict, no EntityConflictError).
- world-postgres: pre-seed an orphaned hook row with no matching
  hook_created event, retry, assert hook_created event lands (no
  hook_conflict, no EntityConflictError).

Both tests fail on the prior implementation (EntityConflictError
thrown on retry, exact symptom from the review).

* fix(world-local): probe the event log (not the hook entity) to detect duplicate hook_created

Addresses follow-up review on PR #2295.

The previous dedup branch checked whether the durable hook entity
existed on disk. But the hook entity is written before the
`hook_created` event, and the two writes are not atomic, so a crash
between them leaves both the claim file and the hook entity on disk
with no event in the log. The dedup branch then matched on
`(runId, hookId)`, found the hook entity, threw EntityConflictError,
and the suspension handler swallowed the retry — permanently losing
`hook_created` from the event log.

The fix mirrors what the world-postgres branch already does: probe
the run's event log for an existing `hook_created` event for the
same `(runId, correlationId)`. The event is the durable record of a
successful hook creation; the claim file and hook entity are partial-
write artifacts that may exist without the event.

- exists  → real duplicate: throw EntityConflictError so the
  runtime's concurrent-replay catch path swallows it.
- missing → orphaned partial write (crash at any point before the
  event landed): re-write the hook entity (with overwrite: true, in
  case a stale partial copy exists) and let the outer code path emit
  the hook_created event.

Added a new helper findHookCreatedEvent that runs a filtered
paginatedFileSystemQuery with limit:1 over the run's events.

Regression test "should recover an orphaned hook entity with no
matching hook_created event" added — pre-creates a hook, deletes
just the hook_created event from disk to simulate a crash between
the entity write and the event write, asserts the retry emits a
fresh hook_created event (no hook_conflict, no swallowed
EntityConflictError). I verified this test fails on the prior fix
(throws `EntityConflictError: Hook "hook_orphan_entity_1" already
created`, exactly as pranaygp reported) and passes on this commit.

The previous test ("should recover an orphaned hook token claim
with no matching hook entity") continues to pass — the event-log
probe is a strict superset of the entity probe, since a missing
entity always also implies a missing event.

* fix(world-local): converge same-hook creation across workers via canonical eventId

Addresses follow-up review on PR #2295.

The previous fix made the dedup branch probe the event log to decide
real-duplicate vs orphan-recovery, but the probe and the recovery
write are not a single atomic operation. Two workers sharing a data
directory (or two retries that lose `writeExclusive(constraintPath)`
back to back) could both pass the probe (each observing no
hook_created event yet), both fall through to the recovery write,
and both append a hook_created event with a different eventId —
producing two events in the log for the same (runId, hookId). The
in-process `withHookLock` mutex does not help here because it is
process-local and tag-specific.

The fix persists `eventId` in the durable token claim file (written
by the original `writeExclusive(constraintPath)`). On a same-(runId,
hookId) dedup match, retries adopt that canonical eventId and
rebuild the event with a deterministic createdAt derived from the
eventId (a ULID). The outer event write switches from `writeJSON`
(check-then-write, TOCTOU) to `writeExclusive` (O_CREAT|O_EXCL via
temp-file + hard-link, atomic across processes). Either worker may
win the publish; the other throws EntityConflictError which the
runtime's existing concurrent-replay catch path swallows. Net
result: exactly one hook_created event per logical creation.

Backward compatibility: a claim file written before this commit lacks
`eventId`. Retries that read such a claim fall back to the
event-log probe + fresh-eventId recovery — the legacy behavior that
does not converge across workers but cannot regress for freshly-
written claims after upgrade.

world-postgres already converges across workers via the partial
unique index on workflow_events_entity_creation_unique
(runId+correlationId+eventType for hook/step/wait_created): the
loser's INSERT raises 23505 which is already translated to
EntityConflictError.

Regression tests:

- world-local: `converges same-hook creation across workers to one
  event` uses two tagged storage instances sharing one data
  directory and fires 25 paired Promise.allSettled hook_created
  calls. Expected 25 hook_created events total; before this fix
  yielded 50.

- world-postgres: `converges same-hook creation across concurrent
  calls to one event` exercises the same shape against the real
  Postgres unique index. Already converges; the test is a guard
  against future regressions to the catch path.

Verified the world-local test fails on c7b23e1b5 with exactly the
shape pranaygp reported (50 events for 25 logical creations) and
passes on this commit. The earlier orphaned-claim and orphaned-
entity recovery tests also continue to pass.

* fix(world-local): converge legacy hook claims via recovery-marker sidecar; replace tag-proxy test with real subprocess workers

Addresses follow-up review on PR #2295.

Two distinct issues, both flagged by pranaygp as P1:

1. The fallback path for token claims written by versions before
   eventId was persisted inline (legacy claims after upgrade) still
   permitted the same cross-process corruption the inline fast path
   was fixed to prevent. Two processes both reading a legacy claim
   each generated their own eventId, landed their
   writeExclusive(eventPath) calls at different paths, and appended
   two hook_created events for the same (runId, hookId). Existing
   persisted claims after a real upgrade are exactly the state the
   crash-recovery branch needs to repair, so leaving the legacy path
   non-convergent is silent corruption, not backward compatibility.

2. The committed cross-worker convergence test used two tagged
   storage instances sharing one directory as a proxy for separate
   processes. But tags change the destination filename
   (events/wrun_X-evnt_Y.worker-a.json vs ...worker-b.json), so two
   tagged workers can each writeExclusive their own event at
   different paths and both fulfill. The Map-by-eventId
   deduplication in the assertion then masked the duplicate
   publication, so the test passed for the wrong reason.

Implementation:

- New HookRecoveryMarkerSchema (`{ eventId, hookId, runId }`) and
  HookRecoveryMarkerPath helper. The marker is a sidecar at
  hooks/tokens/<hash>.recovery.json, written via writeExclusive so
  the first cross-process retry pins its candidate eventId as
  canonical; subsequent retries read the marker and adopt that
  eventId. Together with the existing writeExclusive(eventPath) in
  the outer publish, this gives the legacy-fallback path the same
  single-event convergence guarantee as the inline-eventId fast
  path.

- pinCanonicalEventIdForLegacyClaim() encapsulates the marker
  write-or-read. A stale marker for a different (runId, hookId)
  (token-reuse with leaked state) is overwritten best-effort — the
  common cross-worker race for the same hook still converges; only
  the narrow stale-token-reuse case loses convergence.

- hook_disposed now also deletes the recovery marker when it
  deletes the token constraint file, preventing a future legacy
  recovery for a recycled token from latching onto a stale eventId.

- The dedup branch unified: existingClaim.eventId for new claims,
  pinCanonicalEventIdForLegacyClaim() for legacy ones. Removed the
  now-redundant findHookCreatedEvent helper — the
  writeExclusive(eventPath) in the outer publish is the
  authoritative duplicate-vs-orphan detector.

Tests:

- New test fixture test-fixtures/hook-race-worker.ts (TypeScript,
  run via child_process.fork with tsx as execPath — tsx is a
  transitive dev dep via vitest). Each subprocess gets its own
  createStorage(testDir) so the in-process hookLocks Map cannot
  serialize across workers.

- Replaced the tag-proxy test with
  "converges same-hook creation across separate OS processes to one
  event". Spawns workerCount subprocesses, releases them from a
  barrier into the same hook_created, asserts exactly one fulfilled
  + (N-1) rejected with EntityConflictError, and asserts directly
  on the raw events.list() result (no Map dedup) that the number of
  hook_created entries equals the number of logical creations.

- Added "converges same-hook creation across processes when only a
  legacy token claim exists". Same shape, but pre-seeds the legacy
  claim format (`{ token, hookId, runId }` with no eventId) before
  each race. Verified to FAIL on 7ce66551b (both subprocesses
  fulfill, no convergence) and pass on this commit.

- Also verified the new-eventId subprocess test FAILS when the
  event write is reverted to writeJSON (TOCTOU), confirming it
  exercises the writeExclusive-based cross-process arbitration.
  Both prior orphaned-claim / orphaned-entity recovery tests also
  continue to pass.

* fix(world-local): per-lifetime recovery markers, restore event-log probe, fix CI tsx resolution

Addresses three P1 review comments on PR #2295.

1. Stale recovery marker leaking across token-reuse lifetimes
   (pranaygp):

   The previous marker path used `hashToken(token)` so a stale
   marker for run A could leak into run B's recovery when the same
   token was reused after run A terminated through normal lifecycle.
   `deleteAllHooksForRun()` and tagged `world.clear()` deleted the
   token constraint and hook entity but NOT the marker sidecar, so
   the next legacy claim on the same token entered the stale-marker
   overwrite branch and the workers overwrote it non-atomically,
   yielding divergent publication.

   Fix:
   - Marker path now hashes `(token, runId, hookId)` together
     (`hookRecoveryMarkerPath` in storage/helpers.ts). Different
     lifetimes can never share a marker, so the stale-marker
     overwrite branch is removed entirely.
   - `hookRecoveryMarkerPath` is moved to helpers.ts and shared
     across events-storage.ts, hooks-storage.ts, and index.ts.
   - `deleteAllHooksForRun()` and tagged `world.clear()` now also
     delete the recovery marker for each hook (disk hygiene; per-
     lifetime identity makes leaks no longer corrupting).
   - `hook_disposed` now uses the new per-lifetime marker path too.

2. Duplicate `hook_created` event when a legacy claim's event was
   already published (VADE bot, also implied by pranaygp's analysis):

   Removing the event-log probe from the legacy fallback let a post-
   upgrade retry pin a new canonical eventId via the marker and
   publish a duplicate event at that path, even when the original
   pre-upgrade writer had already successfully published the event
   with its own (different) eventId.

   Fix:
   - Restore `findExistingHookCreatedEventId()` (renamed and made
     to return the eventId for clearer semantics).
   - Legacy fallback now probes the event log BEFORE pinning the
     marker; if a matching `hook_created` event already exists,
     throw `EntityConflictError` so the runtime's concurrent-replay
     catch path swallows the retry.
   - Inline-`eventId` fast path does NOT need the probe — the claim
     itself is the durable convergence key.

3. CI failure: tsx not resolvable under pnpm isolated linking
   (pranaygp; confirmed by ubuntu/windows unit test 60s timeouts):

   The previous test hard-coded `node_modules/.bin/tsx` assuming
   tsx would be hoisted there. But tsx was only a transitive peer
   dep via vitest, and pnpm's isolated linking does NOT link
   transitive peer deps into the workspace bin after a fresh
   install — so neither root nor package-local `.bin/tsx` existed
   in CI, the subprocess fork never started, and the barrier hung
   until vitest killed the test.

   Fix:
   - Add `tsx` as a direct `devDependency` of `@workflow/world-
     local` (pinned to 4.20.6 to match the existing transitive
     resolution).
   - Resolve via `import.meta.resolve('tsx/package.json')` and read
     the `bin` field dynamically, so we adapt to wherever pnpm
     links tsx for this package — not a hard-coded layout.
   - Lazy-init the resolver (no module-load IIFE) so an absent tsx
     fails only the convergence tests, not all 376 tests in the
     file.
   - Surface a clear error message if resolution fails, calling
     out the cause (transitive vs direct deps) for future readers.

   Also: harden the barrier helper so `error` events and pre-ready
   exits resolve BOTH `readyPromises` and `donePromises`, then
   `SIGKILL` siblings. Previously a broken child only resolved
   `donePromises`, leaving `Promise.all(readyPromises)` pending
   until the per-test timeout (60s in CI).

Regression tests added:

- `legacy claim whose hook_created event was already published does
  not append a duplicate event` — pre-seeds a legacy claim AND a
  pre-existing `hook_created` event with a different eventId,
  asserts the retry throws EntityConflictError and the log still
  has exactly the original event.

- `converges legacy claim recovery across run lifetimes after token
  reuse` — runs pranaygp's full lifecycle path: race subprocess
  workers on run A's legacy claim, terminate run A via
  `run_completed` (triggers `deleteAllHooksForRun`), reuse the
  token in a legacy claim for run B, race subprocess workers again,
  asserts exactly one fulfillment + one `EntityConflictError` per
  race and exactly one `hook_created` event per run.

Both new tests verified to fail on 2c673e436 (after rebuilding):
the published-event test throws via duplicate publish instead of
EntityConflictError, the token-reuse test sees both run B workers
fulfill (2 events instead of 1).

The existing orphaned-claim and orphaned-entity recovery tests also
continue to pass.

CI loop confirmed to be repaired locally by spawning subprocesses
via the new resolver and intentionally breaking the worker fixture
to verify the helper fails fast (~500ms) instead of hanging at the
barrier.

* fix(world-local): defer hook entity write until event publish commits

Addresses karthikscale3's P1 review comment on PR #2295.

The dedup-recovery path used to write the hook entity BEFORE the
outer event publish proved whether the attempt was repairing a
missing event or just colliding with an already-published
`hook_created`. For already-committed duplicates, the event write
then throws `EntityConflictError`, but the hook entity had
already been overwritten with the retry's payload — leaving the
durable hook entity and the event log inconsistent (e.g. the
entity reflects the retry's metadata while the event still
carries the original).

karthikscale3 reproduced this on the prior head by creating
`hook_created` with metadata `{ v: "a" }`, then retrying the
same `(runId, hookId, token)` with metadata `{ v: "b" }` and
`isWebhook: false`: the retry threw `EntityConflictError` but
`hooks.get()` returned the retry's payload.

Fix: defer the hook entity write until AFTER the outer
`writeExclusive(eventPath)` commits. The branch now only
captures the entity-to-write and its overwrite options; the
actual write happens immediately after the event publish in the
shared trailing block. A retry that ends in
`EntityConflictError` (the event was already published) now
leaves the entity untouched.

The first-writer happy path and all recovery paths (orphaned-
claim, orphaned-entity, cross-worker convergence, legacy claim,
token-reuse across lifetimes) are unaffected — they all reach
the event publish successfully, then the entity write runs as
before.

Regression test `does not mutate an already-committed hook
entity when a duplicate hook_created retry collides` added to
world-local: runs karthikscale3's exact scenario and asserts the
persisted entity still carries the original metadata and
isWebhook. Verified to fail on the prior commit (persisted
metadata = 0xbb instead of 0xaa) and pass on this commit after
rebuilding.

Parallel guard test `does not mutate an already-committed hook
entity when a duplicate hook_created retry collides` added to
world-postgres. Postgres already protected this via
`onConflictDoNothing()` on the hook INSERT, but the test guards
against a future regression that adds an UPDATE/UPSERT to the
dedup path.

* refactor(world-local): per-instance in-process locks; drop tsx subprocess test plumbing

You were right that the tsx subprocess machinery was overkill for a
storage-level convergence test. Replaced with a simple two-instance
in-process test that exercises the same cross-process semantics
without spawning anything.

The trick: `stepLocks` and `hookLocks` were module-level Maps shared
by all `createEventsStorage` calls in the same process. Move them
inside the function so each `createStorage(dir)` call gets its own
lock map. Two storage instances sharing one data directory then
behave exactly like two separate OS processes:

  - independent in-process `hookLocks` Maps (no in-process
    serialization between them), and
  - a shared filesystem (so the on-disk `writeExclusive` claim /
    marker / event publish primitives are the only thing arbitrating
    convergence).

This is also a real architectural improvement — the global lock map
was always a leaky abstraction that made unit-test simulation of
the cross-process path awkward.

Changes:

- `stepLocks` and `hookLocks` moved from module scope into
  `createEventsStorage`. `withStepLock` and `withHookLock` wrappers
  collapsed into direct `withInProcessLock(map, key, fn)` calls at
  the two call sites that need them.

- The three convergence regression tests in `storage.test.ts` now
  use `const workerA = createStorage(testDir); const workerB =
  createStorage(testDir);` and race `Promise.allSettled` of
  `events.create` from both — no subprocess, no IPC, no barrier
  helper, no `raceHookCreatedAcrossProcesses`. Same assertions
  (exactly one fulfillment + N-1 `EntityConflictError` per race,
  raw `events.list()` shows exactly one `hook_created` per
  logical creation — no Map dedup) so the regression catches are
  identical.

- Removed: `tsx` devDep, `test-fixtures/hook-race-worker.ts`,
  `HOOK_RACE_WORKER` / `resolveTsxLoaderUrl` / `TSX_BIN` /
  `raceHookCreatedAcrossProcesses` and the
  `fork`/`fileURLToPath` imports they pulled in.

Verified (after rebuilding world-local):

- All 379 tests pass on macOS in ~1s (was ~6.7s with subprocesses).
- Convergence tests confirmed to still catch the bugs: temporarily
  reverted the `eventId = canonicalEventId` adoption → both workers
  fulfilled (2 events instead of 1). Temporarily reverted the
  legacy-claim marker pin → same: both workers fulfilled.
- No subprocess machinery means no Windows-specific quirks
  (cli.mjs shebang, .cmd wrappers, .bin hoisting under pnpm
  isolated linking, etc.) that produced the Windows CI 60s
  timeouts.
- World-postgres still has its own parallel guard test for the
  karthikscale3 "no-mutate-on-duplicate" regression; that one
  exercises real DB concurrency and is unaffected by this change.

Full repo `pnpm test` (43 packages) and the
`parallelStepsThenWebhookWorkflow` e2e test against world-local
both green.

* fix(world-local): repair event-first hook orphans from the persisted event; skip #1665 e2e on world-postgres

- A crash between the hook_created event publish and the deferred hook
  entity write left the event committed with the entity missing and
  unrepairable (retries threw EntityConflictError without materializing
  the entity). Retries now rebuild the entity from the PERSISTED event's
  payload — never the retry's eventData — via a race-safe writeExclusive,
  on both the canonical-eventId collision path and the legacy-claim
  probe path.
- Skip parallelStepsThenWebhookWorkflow e2e on world-postgres: the
  same-tick replay pattern surfaces a separate pre-existing step_started
  ordering bug there (#2331).

---------

Co-authored-by: Peter Wielander <peter.wielander@vercel.com>
2026-06-11 14:04:49 -07:00

2867 lines
94 KiB
TypeScript

import { execSync } from 'node:child_process';
import { PostgreSqlContainer } from '@testcontainers/postgresql';
import type { Hook, Step, WorkflowRun } from '@workflow/world';
import { SPEC_VERSION_CURRENT } from '@workflow/world';
import { encode } from 'cbor-x';
import { Pool } from 'pg';
import {
afterAll,
beforeAll,
beforeEach,
describe,
expect,
it,
test,
} from 'vitest';
import { createClient } from '../src/drizzle/index.js';
import * as DrizzleSchema from '../src/drizzle/schema.js';
import {
createEventsStorage,
createHooksStorage,
createRunsStorage,
createStepsStorage,
} from '../src/storage.js';
// Helper types for events storage
type EventsStorage = ReturnType<typeof createEventsStorage>;
// Helper functions to create entities through events.create
async function createRun(
events: EventsStorage,
data: {
deploymentId: string;
workflowName: string;
input: Uint8Array;
executionContext?: Record<string, unknown>;
attributes?: Record<string, string>;
}
): Promise<WorkflowRun> {
const result = await events.create(null, {
eventType: 'run_created',
eventData: data,
});
if (!result.run) {
throw new Error('Expected run to be created');
}
return result.run;
}
async function updateRun(
events: EventsStorage,
runId: string,
eventType: 'run_started' | 'run_completed' | 'run_failed',
eventData?: Record<string, unknown>
): Promise<WorkflowRun> {
const result = await events.create(runId, {
eventType,
eventData,
});
if (!result.run) {
throw new Error('Expected run to be updated');
}
return result.run;
}
async function createStep(
events: EventsStorage,
runId: string,
data: {
stepId: string;
stepName: string;
input: Uint8Array;
}
): Promise<Step> {
const result = await events.create(runId, {
eventType: 'step_created',
correlationId: data.stepId,
eventData: { stepName: data.stepName, input: data.input },
});
if (!result.step) {
throw new Error('Expected step to be created');
}
return result.step;
}
async function updateStep(
events: EventsStorage,
runId: string,
stepId: string,
eventType: 'step_started' | 'step_completed' | 'step_failed',
eventData?: Record<string, unknown>
): Promise<Step> {
const result = await events.create(runId, {
eventType,
correlationId: stepId,
eventData,
});
if (!result.step) {
throw new Error('Expected step to be updated');
}
return result.step;
}
async function createHook(
events: EventsStorage,
runId: string,
data: {
hookId: string;
token: string;
metadata?: unknown;
}
): Promise<Hook> {
const result = await events.create(runId, {
eventType: 'hook_created',
correlationId: data.hookId,
eventData: { token: data.token, metadata: data.metadata },
});
if (!result.hook) {
throw new Error('Expected hook to be created');
}
return result.hook;
}
describe('Storage (Postgres integration)', () => {
if (process.platform === 'win32') {
test.skip('skipped on Windows since it relies on a docker container', () => {});
return;
}
let container: Awaited<ReturnType<PostgreSqlContainer['start']>>;
let pool: Pool;
let drizzle: ReturnType<typeof createClient>;
let runs: ReturnType<typeof createRunsStorage>;
let steps: ReturnType<typeof createStepsStorage>;
let events: ReturnType<typeof createEventsStorage>;
let hooks: ReturnType<typeof createHooksStorage>;
async function truncateTables() {
await pool.query(
'TRUNCATE TABLE workflow.workflow_events, workflow.workflow_steps, workflow.workflow_hooks, workflow.workflow_runs RESTART IDENTITY CASCADE'
);
}
beforeAll(async () => {
// Start PostgreSQL container
container = await new PostgreSqlContainer('postgres:15-alpine').start();
const dbUrl = container.getConnectionUri();
process.env.DATABASE_URL = dbUrl;
process.env.WORKFLOW_POSTGRES_URL = dbUrl;
// Apply schema
execSync('pnpm db:push', {
stdio: 'inherit',
cwd: process.cwd(),
env: process.env,
});
// Initialize database clients and storage
pool = new Pool({ connectionString: dbUrl, max: 1 });
drizzle = createClient(pool);
runs = createRunsStorage(drizzle);
steps = createStepsStorage(drizzle);
events = createEventsStorage(drizzle);
hooks = createHooksStorage(drizzle);
}, 120_000);
beforeEach(async () => {
await truncateTables();
});
afterAll(async () => {
await pool.end();
await container.stop();
});
describe('runs', () => {
describe('create', () => {
it('should create a new workflow run', async () => {
const runData = {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
executionContext: { userId: 'user-1' },
input: new Uint8Array([1, 2]),
};
const run = await createRun(events, runData);
expect(run.runId).toMatch(/^wrun_/);
expect(run.deploymentId).toBe('deployment-123');
expect(run.status).toBe('pending');
expect(run.workflowName).toBe('test-workflow');
expect(run.executionContext).toEqual({ userId: 'user-1' });
expect(run.input).toEqual(new Uint8Array([1, 2]));
expect(run.output).toBeUndefined();
expect(run.error).toBeUndefined();
expect(run.startedAt).toBeUndefined();
expect(run.completedAt).toBeUndefined();
expect(run.createdAt).toBeInstanceOf(Date);
expect(run.updatedAt).toBeInstanceOf(Date);
});
it('should handle minimal run data', async () => {
const runData = {
deploymentId: 'deployment-123',
workflowName: 'minimal-workflow',
input: new Uint8Array(),
};
const run = await createRun(events, runData);
expect(run.executionContext).toBeUndefined();
expect(run.input).toEqual(new Uint8Array());
});
it('should seed initial attributes from run_created', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'attributed-workflow',
input: new Uint8Array(),
attributes: { tenant: 't1', phase: 'created' },
});
expect(run.attributes).toEqual({ tenant: 't1', phase: 'created' });
});
it('treats SQL-looking initial attribute keys as literal JSON keys', async () => {
const key = "tenant'); DROP TABLE workflow_runs; --";
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'attributed-workflow',
input: new Uint8Array(),
attributes: { [key]: 'literal' },
});
expect(run.attributes).toEqual({ [key]: 'literal' });
});
});
describe('get', () => {
it('should retrieve an existing run', async () => {
const created = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array([1]),
});
const retrieved = await runs.get(created.runId);
expect(retrieved.runId).toBe(created.runId);
expect(retrieved.workflowName).toBe('test-workflow');
expect(retrieved.input).toEqual(new Uint8Array([1]));
});
it('should throw error for non-existent run', async () => {
await expect(runs.get('missing')).rejects.toMatchObject({
name: 'WorkflowRunNotFoundError',
});
});
});
describe('update via events', () => {
it('should update run status to running via run_started event', async () => {
const created = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
const updated = await updateRun(events, created.runId, 'run_started');
expect(updated.status).toBe('running');
expect(updated.startedAt).toBeInstanceOf(Date);
});
it('should update run status to completed via run_completed event', async () => {
const created = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
const updated = await updateRun(
events,
created.runId,
'run_completed',
{
output: new Uint8Array([42]),
}
);
expect(updated.status).toBe('completed');
expect(updated.completedAt).toBeInstanceOf(Date);
expect(updated.output).toEqual(new Uint8Array([42]));
});
it('should update run status to failed via run_failed event', async () => {
const created = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// The `error` field is opaque SerializedData (Uint8Array) produced by
// dehydrateRunError. The storage layer persists it verbatim.
const serializedError = new Uint8Array([1, 2, 3]);
const updated = await updateRun(events, created.runId, 'run_failed', {
error: serializedError,
});
expect(updated.status).toBe('failed');
expect(updated.error).toEqual(serializedError);
expect(updated.completedAt).toBeInstanceOf(Date);
});
});
describe('list', () => {
it('should list all runs', async () => {
const run1 = await createRun(events, {
deploymentId: 'deployment-1',
workflowName: 'workflow-1',
input: new Uint8Array(),
});
// Small delay to ensure different timestamps in createdAt
await new Promise((resolve) => setTimeout(resolve, 2));
const run2 = await createRun(events, {
deploymentId: 'deployment-2',
workflowName: 'workflow-2',
input: new Uint8Array(),
});
const result = await runs.list();
expect(result.data).toHaveLength(2);
// Should be in descending order (most recent first)
expect(result.data[0].runId).toBe(run2.runId);
expect(result.data[1].runId).toBe(run1.runId);
expect(result.data[0].createdAt.getTime()).toBeGreaterThan(
result.data[1].createdAt.getTime()
);
});
it('should filter runs by workflowName', async () => {
await createRun(events, {
deploymentId: 'deployment-1',
workflowName: 'workflow-1',
input: new Uint8Array(),
});
const run2 = await createRun(events, {
deploymentId: 'deployment-2',
workflowName: 'workflow-2',
input: new Uint8Array(),
});
const result = await runs.list({ workflowName: 'workflow-2' });
expect(result.data).toHaveLength(1);
expect(result.data[0].runId).toBe(run2.runId);
});
it('should support pagination', async () => {
// Create multiple runs
for (let i = 0; i < 5; i++) {
await createRun(events, {
deploymentId: `deployment-${i}`,
workflowName: `workflow-${i}`,
input: new Uint8Array(),
});
}
const page1 = await runs.list({
pagination: { limit: 2 },
});
expect(page1.data).toHaveLength(2);
expect(page1.cursor).not.toBeNull();
const page2 = await runs.list({
pagination: { limit: 2, cursor: page1.cursor || undefined },
});
expect(page2.data).toHaveLength(2);
expect(page2.data[0].runId).not.toBe(page1.data[0].runId);
});
});
describe('experimentalSetAttributes', () => {
it('upserts new keys', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
});
const result = await runs.experimentalSetAttributes!(run.runId, [
{ key: 'phase', value: 'init' },
{ key: 'tenant', value: 't1' },
]);
expect(result.attributes).toEqual({ phase: 'init', tenant: 't1' });
const fresh = await runs.get(run.runId);
expect(fresh.attributes).toEqual({ phase: 'init', tenant: 't1' });
});
it('merges across calls without clobbering prior keys', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
});
await runs.experimentalSetAttributes!(run.runId, [
{ key: 'a', value: '1' },
]);
const result = await runs.experimentalSetAttributes!(run.runId, [
{ key: 'b', value: '2' },
]);
expect(result.attributes).toEqual({ a: '1', b: '2' });
});
it('removes keys when value is null', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
});
await runs.experimentalSetAttributes!(run.runId, [
{ key: 'a', value: '1' },
{ key: 'b', value: '2' },
]);
const result = await runs.experimentalSetAttributes!(run.runId, [
{ key: 'a', value: null },
]);
expect(result.attributes).toEqual({ b: '2' });
});
});
describe('native attr_set events', () => {
it('materializes writes and removals on the run', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
attributes: { stale: 'remove' },
});
const result = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
correlationId: 'attr_1',
eventData: {
changes: [
{ key: 'phase', value: 'ready' },
{ key: 'stale', value: null },
],
writer: { type: 'workflow' },
},
});
expect(result.event?.eventType).toBe('attr_set');
expect(result.run?.attributes).toEqual({ phase: 'ready' });
expect((await runs.get(run.runId)).attributes).toEqual({
phase: 'ready',
});
});
it('requires reserved-key opt-in on native events', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
});
await expect(
events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key: '$system', value: 'nope' }],
writer: { type: 'workflow' },
},
})
).rejects.toThrow(/reserved prefix/);
const result = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key: '$system', value: 'ok' }],
writer: { type: 'workflow' },
allowReservedAttributes: true,
},
});
expect(result.run?.attributes).toEqual({ $system: 'ok' });
});
it('treats SQL-looking attribute keys as literal JSON keys', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
});
const key = "phase'); DROP TABLE workflow_runs; --";
const written = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key, value: 'literal' }],
writer: { type: 'workflow' },
},
});
expect(written.run?.attributes).toEqual({ [key]: 'literal' });
const removed = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key, value: null }],
writer: { type: 'workflow' },
},
});
expect(removed.run?.attributes).toEqual({});
});
it('enforces the per-run cap against existing attributes', async () => {
const initial: Record<string, string> = {};
for (let i = 0; i < 63; i++) initial[`a${i}`] = 'v';
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
attributes: initial,
});
// 64th attribute fits exactly at the cap.
const atCap = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key: 'a63', value: 'v' }],
writer: { type: 'workflow' },
},
});
expect(Object.keys(atCap.run?.attributes ?? {})).toHaveLength(64);
// A 65th attribute exceeds the cap with a clear error.
await expect(
events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key: 'a64', value: 'v' }],
writer: { type: 'workflow' },
},
})
).rejects.toThrow(/exceed limit 64/);
// Upserting an existing key at the cap is a zero-net change.
const upserted = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key: 'a0', value: 'updated' }],
writer: { type: 'step', stepId: 'step_1', attempt: 1 },
},
});
expect(upserted.run?.attributes?.a0).toBe('updated');
// Removing a key frees room for a new one in the same batch.
const swapped = await events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [
{ key: 'a1', value: null },
{ key: 'replacement', value: 'v' },
],
writer: { type: 'workflow' },
},
});
expect(swapped.run?.attributes?.replacement).toBe('v');
expect(swapped.run?.attributes).not.toHaveProperty('a1');
expect(Object.keys(swapped.run?.attributes ?? {})).toHaveLength(64);
});
it('rejects oversized attribute values on attr_set', async () => {
const run = await createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
});
await expect(
events.create(run.runId, {
eventType: 'attr_set',
specVersion: SPEC_VERSION_CURRENT,
eventData: {
changes: [{ key: 'note', value: 'v'.repeat(257) }],
writer: { type: 'workflow' },
},
})
).rejects.toThrow(/byte length 257 exceeds limit 256/);
});
it('rejects invalid initial attributes on run_created', async () => {
const overCap: Record<string, string> = {};
for (let i = 0; i <= 64; i++) overCap[`a${i}`] = 'v';
await expect(
createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
attributes: overCap,
})
).rejects.toThrow(/exceed limit 64/);
await expect(
createRun(events, {
deploymentId: 'd',
workflowName: 'w',
input: new Uint8Array(),
attributes: { $reserved: 'nope' },
})
).rejects.toThrow(/reserved prefix/);
});
});
});
describe('steps', () => {
let testRunId: string;
beforeEach(async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
testRunId = run.runId;
});
describe('create', () => {
it('should create a new step', async () => {
const stepData = {
stepId: 'step-123',
stepName: 'test-step',
input: new Uint8Array([1, 2]),
};
const step = await createStep(events, testRunId, stepData);
expect(step).toMatchObject({
runId: testRunId,
stepId: 'step-123',
stepName: 'test-step',
status: 'pending',
input: new Uint8Array([1, 2]),
output: undefined,
error: undefined,
attempt: 0, // steps are created with attempt 0
startedAt: undefined,
completedAt: undefined,
createdAt: expect.any(Date),
updatedAt: expect.any(Date),
specVersion: SPEC_VERSION_CURRENT,
});
});
});
describe('get', () => {
it('should retrieve a step with runId and stepId', async () => {
const created = await createStep(events, testRunId, {
stepId: 'step-123',
stepName: 'test-step',
input: new Uint8Array([1]),
});
const retrieved = await steps.get(testRunId, 'step-123');
expect(retrieved.stepId).toBe(created.stepId);
});
it('should throw error for non-existent step', async () => {
await expect(steps.get(testRunId, 'missing-step')).rejects.toThrow(
'Step not found'
);
});
});
describe('update via events', () => {
it('should update step status to running via step_started event', async () => {
await createStep(events, testRunId, {
stepId: 'step-123',
stepName: 'test-step',
input: new Uint8Array([1]),
});
const updated = await updateStep(
events,
testRunId,
'step-123',
'step_started',
{} // step_started no longer needs attempt in eventData - World increments it
);
expect(updated.status).toBe('running');
expect(updated.startedAt).toBeInstanceOf(Date);
expect(updated.attempt).toBe(1); // Incremented by step_started
});
it('should update step status to completed via step_completed event', async () => {
await createStep(events, testRunId, {
stepId: 'step-123',
stepName: 'test-step',
input: new Uint8Array([1]),
});
const updated = await updateStep(
events,
testRunId,
'step-123',
'step_completed',
{ result: new Uint8Array([1]) }
);
expect(updated.status).toBe('completed');
expect(updated.completedAt).toBeInstanceOf(Date);
expect(updated.output).toEqual(new Uint8Array([1]));
});
it('should update step status to failed via step_failed event', async () => {
await createStep(events, testRunId, {
stepId: 'step-123',
stepName: 'test-step',
input: new Uint8Array([1]),
});
// The `error` field is opaque SerializedData (Uint8Array) produced by
// dehydrateStepError. The storage layer persists it verbatim.
const serializedError = new Uint8Array([1, 2, 3]);
const updated = await updateStep(
events,
testRunId,
'step-123',
'step_failed',
{ error: serializedError }
);
expect(updated.status).toBe('failed');
expect(updated.error).toEqual(serializedError);
expect(updated.completedAt).toBeInstanceOf(Date);
});
});
describe('list', () => {
it('should list all steps for a run', async () => {
const step1 = await createStep(events, testRunId, {
stepId: 'step-1',
stepName: 'first-step',
input: new Uint8Array(),
});
const step2 = await createStep(events, testRunId, {
stepId: 'step-2',
stepName: 'second-step',
input: new Uint8Array(),
});
const result = await steps.list({
runId: testRunId,
});
expect(result.data).toHaveLength(2);
// Should be in descending order
expect(result.data[0].stepId).toBe(step2.stepId);
expect(result.data[1].stepId).toBe(step1.stepId);
expect(result.data[0].createdAt.getTime()).toBeGreaterThanOrEqual(
result.data[1].createdAt.getTime()
);
});
it('should support pagination', async () => {
// Create multiple steps
for (let i = 0; i < 5; i++) {
await createStep(events, testRunId, {
stepId: `step-${i}`,
stepName: `step-name-${i}`,
input: new Uint8Array(),
});
}
const page1 = await steps.list({
runId: testRunId,
pagination: { limit: 2 },
});
expect(page1.data).toHaveLength(2);
expect(page1.cursor).not.toBeNull();
const page2 = await steps.list({
runId: testRunId,
pagination: { limit: 2, cursor: page1.cursor || undefined },
});
expect(page2.data).toHaveLength(2);
expect(page2.data[0].stepId).not.toBe(page1.data[0].stepId);
});
});
});
describe('events', () => {
let testRunId: string;
beforeEach(async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
testRunId = run.runId;
});
describe('create', () => {
it('should create a new event', async () => {
// Create step before step_started event
await createStep(events, testRunId, {
stepId: 'corr_123',
stepName: 'test-step',
input: new Uint8Array(),
});
const eventData = {
eventType: 'step_started' as const,
correlationId: 'corr_123',
};
const result = await events.create(testRunId, eventData);
expect(result.event.runId).toBe(testRunId);
expect(result.event.eventId).toMatch(/^wevt_/);
expect(result.event.eventType).toBe('step_started');
expect(result.event.correlationId).toBe('corr_123');
expect(result.event.createdAt).toBeInstanceOf(Date);
});
it('should create a new event with null byte in payload', async () => {
// Create step before step_failed event
await createStep(events, testRunId, {
stepId: 'corr_123_null',
stepName: 'test-step-null',
input: new Uint8Array(),
});
await events.create(testRunId, {
eventType: 'step_started',
correlationId: 'corr_123_null',
});
const result = await events.create(testRunId, {
eventType: 'step_failed',
correlationId: 'corr_123_null',
eventData: { error: 'Error with null byte \u0000 in message' },
});
expect(result.event.runId).toBe(testRunId);
expect(result.event.eventId).toMatch(/^wevt_/);
expect(result.event.eventType).toBe('step_failed');
expect(result.event.correlationId).toBe('corr_123_null');
expect(result.event.createdAt).toBeInstanceOf(Date);
});
it('should handle run completed events', async () => {
const eventData = {
eventType: 'run_completed' as const,
eventData: { output: new Uint8Array([1]) },
};
const result = await events.create(testRunId, eventData);
expect(result.event.eventType).toBe('run_completed');
expect(result.event.correlationId).toBeUndefined();
});
});
describe('list', () => {
it('should list all events for a run', async () => {
const result1 = await events.create(testRunId, {
eventType: 'run_started' as const,
});
// Small delay to ensure different timestamps in event IDs
await new Promise((resolve) => setTimeout(resolve, 2));
// Create step before step_started event
await createStep(events, testRunId, {
stepId: 'corr-step-1',
stepName: 'test-step',
input: new Uint8Array(),
});
const result2 = await events.create(testRunId, {
eventType: 'step_started' as const,
correlationId: 'corr-step-1',
});
const result = await events.list({
runId: testRunId,
pagination: { sortOrder: 'asc' }, // Explicitly request ascending order
});
// 4 events: run_created (from createRun), run_started, step_created, step_started
expect(result.data).toHaveLength(4);
// Should be in chronological order (oldest first)
expect(result.data[0].eventType).toBe('run_created');
expect(result.data[1].eventId).toBe(result1.event.eventId);
expect(result.data[3].eventId).toBe(result2.event.eventId);
expect(result.data[3].createdAt.getTime()).toBeGreaterThanOrEqual(
result.data[1].createdAt.getTime()
);
});
it('should list events in descending order when explicitly requested (newest first)', async () => {
const result1 = await events.create(testRunId, {
eventType: 'run_started' as const,
});
// Small delay to ensure different timestamps in event IDs
await new Promise((resolve) => setTimeout(resolve, 2));
// Create step before step_started event
await createStep(events, testRunId, {
stepId: 'corr-step-1',
stepName: 'test-step',
input: new Uint8Array(),
});
const result2 = await events.create(testRunId, {
eventType: 'step_started' as const,
correlationId: 'corr-step-1',
});
const result = await events.list({
runId: testRunId,
pagination: { sortOrder: 'desc' },
});
// 4 events: run_created (from createRun), run_started, step_created, step_started
expect(result.data).toHaveLength(4);
// Should be in reverse chronological order (newest first)
expect(result.data[0].eventId).toBe(result2.event.eventId);
expect(result.data[1].eventType).toBe('step_created');
expect(result.data[2].eventId).toBe(result1.event.eventId);
expect(result.data[3].eventType).toBe('run_created');
expect(result.data[0].createdAt.getTime()).toBeGreaterThanOrEqual(
result.data[2].createdAt.getTime()
);
});
it('should support pagination', async () => {
// Create multiple events - must create steps first
for (let i = 0; i < 5; i++) {
await createStep(events, testRunId, {
stepId: `corr_${i}`,
stepName: `test-step-${i}`,
input: new Uint8Array(),
});
// Start the step before completing
await events.create(testRunId, {
eventType: 'step_started',
correlationId: `corr_${i}`,
});
await events.create(testRunId, {
eventType: 'step_completed',
correlationId: `corr_${i}`,
eventData: { result: new Uint8Array([i]) },
});
}
const page1 = await events.list({
runId: testRunId,
pagination: { limit: 2 },
});
expect(page1.data).toHaveLength(2);
expect(page1.cursor).not.toBeNull();
const page2 = await events.list({
runId: testRunId,
pagination: { limit: 2, cursor: page1.cursor || undefined },
});
expect(page2.data).toHaveLength(2);
expect(page2.data[0].eventId).not.toBe(page1.data[0].eventId);
});
});
describe('listByCorrelationId', () => {
it('should list all events with a specific correlation ID', async () => {
const correlationId = 'step-abc123';
// Create step before step events
await createStep(events, testRunId, {
stepId: correlationId,
stepName: 'test-step',
input: new Uint8Array(),
});
// Create events with the target correlation ID
const result1 = await events.create(testRunId, {
eventType: 'step_started',
correlationId,
});
await new Promise((resolve) => setTimeout(resolve, 2));
const result2 = await events.create(testRunId, {
eventType: 'step_completed',
correlationId,
eventData: { result: new Uint8Array([1]) },
});
// Create events with different correlation IDs (should be filtered out)
await createStep(events, testRunId, {
stepId: 'different-step',
stepName: 'different-step',
input: new Uint8Array(),
});
await events.create(testRunId, {
eventType: 'step_started',
correlationId: 'different-step',
});
await events.create(testRunId, {
eventType: 'run_completed',
eventData: { output: new Uint8Array([1]) },
});
const result = await events.listByCorrelationId({
correlationId,
pagination: {},
});
// 3 events: step_created, step_started, step_completed
expect(result.data).toHaveLength(3);
expect(result.data[0].eventType).toBe('step_created');
expect(result.data[1].eventId).toBe(result1.event.eventId);
expect(result.data[1].correlationId).toBe(correlationId);
expect(result.data[2].eventId).toBe(result2.event.eventId);
expect(result.data[2].correlationId).toBe(correlationId);
});
it('should list events across multiple runs with same correlation ID', async () => {
const correlationId = 'hook-xyz789';
// Create another run
const run2 = await createRun(events, {
deploymentId: 'deployment-456',
workflowName: 'test-workflow-2',
input: new Uint8Array(),
});
// Create events in both runs with same correlation ID
const result1 = await events.create(testRunId, {
eventType: 'hook_created',
correlationId,
eventData: { token: 'test-token-1' },
});
await new Promise((resolve) => setTimeout(resolve, 2));
const result2 = await events.create(run2.runId, {
eventType: 'hook_received',
correlationId,
eventData: { payload: new Uint8Array([1, 2, 3]) },
});
await new Promise((resolve) => setTimeout(resolve, 2));
const result3 = await events.create(testRunId, {
eventType: 'hook_disposed',
correlationId,
});
const result = await events.listByCorrelationId({
correlationId,
pagination: {},
});
expect(result.data).toHaveLength(3);
expect(result.data[0].eventId).toBe(result1.event.eventId);
expect(result.data[0].runId).toBe(testRunId);
expect(result.data[1].eventId).toBe(result2.event.eventId);
expect(result.data[1].runId).toBe(run2.runId);
expect(result.data[2].eventId).toBe(result3.event.eventId);
expect(result.data[2].runId).toBe(testRunId);
});
it('should return empty list for non-existent correlation ID', async () => {
// Create a step and start it
await createStep(events, testRunId, {
stepId: 'existing-step',
stepName: 'existing-step',
input: new Uint8Array(),
});
await events.create(testRunId, {
eventType: 'step_started',
correlationId: 'existing-step',
});
const result = await events.listByCorrelationId({
correlationId: 'non-existent-correlation-id',
pagination: {},
});
expect(result.data).toHaveLength(0);
expect(result.hasMore).toBe(false);
expect(result.cursor).toBeNull();
});
it('should respect pagination parameters', async () => {
const correlationId = 'step_paginated';
// Create step first
await createStep(events, testRunId, {
stepId: correlationId,
stepName: 'test-step',
input: new Uint8Array(),
});
// Create multiple events
await events.create(testRunId, {
eventType: 'step_started',
correlationId,
});
await new Promise((resolve) => setTimeout(resolve, 2));
await events.create(testRunId, {
eventType: 'step_retrying',
correlationId,
eventData: { error: 'retry error' },
});
await new Promise((resolve) => setTimeout(resolve, 2));
// Start again after retry
await events.create(testRunId, {
eventType: 'step_started',
correlationId,
});
await new Promise((resolve) => setTimeout(resolve, 2));
await events.create(testRunId, {
eventType: 'step_completed',
correlationId,
eventData: { result: new Uint8Array([1]) },
});
// Get first page (step_created, step_started, step_retrying)
const page1 = await events.listByCorrelationId({
correlationId,
pagination: { limit: 3 },
});
expect(page1.data).toHaveLength(3);
expect(page1.hasMore).toBe(true);
expect(page1.cursor).toBeDefined();
// Get second page (step_started, step_completed)
const page2 = await events.listByCorrelationId({
correlationId,
pagination: { limit: 3, cursor: page1.cursor || undefined },
});
expect(page2.data).toHaveLength(2);
expect(page2.hasMore).toBe(false);
});
it('should always return full event data', async () => {
// Create step first
await createStep(events, testRunId, {
stepId: 'step-with-data',
stepName: 'step-with-data',
input: new Uint8Array(),
});
// Start the step before completing
await events.create(testRunId, {
eventType: 'step_started',
correlationId: 'step-with-data',
});
await events.create(testRunId, {
eventType: 'step_completed',
correlationId: 'step-with-data',
eventData: { result: new Uint8Array([1]) },
});
// Note: resolveData parameter is ignored by the PG World storage implementation
const result = await events.listByCorrelationId({
correlationId: 'step-with-data',
pagination: {},
});
// 3 events: step_created, step_started, step_completed
expect(result.data).toHaveLength(3);
expect(result.data[2].correlationId).toBe('step-with-data');
});
it('should return events in ascending order by default', async () => {
const correlationId = 'step-ordering';
// Create step first
await createStep(events, testRunId, {
stepId: correlationId,
stepName: 'test-step',
input: new Uint8Array(),
});
// Create events with slight delays to ensure different timestamps
const result1 = await events.create(testRunId, {
eventType: 'step_started',
correlationId,
});
await new Promise((resolve) => setTimeout(resolve, 2));
const result2 = await events.create(testRunId, {
eventType: 'step_completed',
correlationId,
eventData: { result: new Uint8Array([1]) },
});
const result = await events.listByCorrelationId({
correlationId,
pagination: {},
});
// 3 events: step_created, step_started, step_completed
expect(result.data).toHaveLength(3);
expect(result.data[1].eventId).toBe(result1.event.eventId);
expect(result.data[2].eventId).toBe(result2.event.eventId);
expect(result.data[1].createdAt.getTime()).toBeLessThanOrEqual(
result.data[2].createdAt.getTime()
);
});
it('should support descending order', async () => {
const correlationId = 'step-desc-order';
// Create step first
await createStep(events, testRunId, {
stepId: correlationId,
stepName: 'test-step',
input: new Uint8Array(),
});
const result1 = await events.create(testRunId, {
eventType: 'step_started',
correlationId,
});
await new Promise((resolve) => setTimeout(resolve, 2));
const result2 = await events.create(testRunId, {
eventType: 'step_completed',
correlationId,
eventData: { result: new Uint8Array([1]) },
});
const result = await events.listByCorrelationId({
correlationId,
pagination: { sortOrder: 'desc' },
});
// 3 events in descending order: step_completed, step_started, step_created
expect(result.data).toHaveLength(3);
expect(result.data[0].eventId).toBe(result2.event.eventId);
expect(result.data[1].eventId).toBe(result1.event.eventId);
expect(result.data[0].createdAt.getTime()).toBeGreaterThanOrEqual(
result.data[1].createdAt.getTime()
);
});
it('should handle hook lifecycle events', async () => {
const hookId = 'hook_test123';
// Create a typical hook lifecycle
const createdResult = await events.create(testRunId, {
eventType: 'hook_created' as const,
correlationId: hookId,
eventData: { token: 'lifecycle-test-token' },
});
await new Promise((resolve) => setTimeout(resolve, 2));
const received1Result = await events.create(testRunId, {
eventType: 'hook_received' as const,
correlationId: hookId,
eventData: { payload: new Uint8Array([1]) },
});
await new Promise((resolve) => setTimeout(resolve, 2));
const received2Result = await events.create(testRunId, {
eventType: 'hook_received' as const,
correlationId: hookId,
eventData: { payload: new Uint8Array([2]) },
});
await new Promise((resolve) => setTimeout(resolve, 2));
const disposedResult = await events.create(testRunId, {
eventType: 'hook_disposed' as const,
correlationId: hookId,
});
const result = await events.listByCorrelationId({
correlationId: hookId,
pagination: {},
});
expect(result.data).toHaveLength(4);
expect(result.data[0].eventId).toBe(createdResult.event.eventId);
expect(result.data[0].eventType).toBe('hook_created');
expect(result.data[1].eventId).toBe(received1Result.event.eventId);
expect(result.data[1].eventType).toBe('hook_received');
expect(result.data[2].eventId).toBe(received2Result.event.eventId);
expect(result.data[2].eventType).toBe('hook_received');
expect(result.data[3].eventId).toBe(disposedResult.event.eventId);
expect(result.data[3].eventType).toBe('hook_disposed');
});
it('should enforce token uniqueness across different runs', async () => {
const token = 'unique-token-test';
// Create first hook with the token
await events.create(testRunId, {
eventType: 'hook_created' as const,
correlationId: 'hook_1',
eventData: { token },
});
// Create another run
const run2 = await createRun(events, {
deploymentId: 'deployment-456',
workflowName: 'test-workflow-2',
input: new Uint8Array(),
});
// Try to create another hook with the same token - should return hook_conflict event
const result = await events.create(run2.runId, {
eventType: 'hook_created' as const,
correlationId: 'hook_2',
eventData: { token },
});
// Should return a hook_conflict event instead of throwing
expect(result.event.eventType).toBe('hook_conflict');
expect(result.event.correlationId).toBe('hook_2');
expect((result.event as any).eventData.token).toBe(token);
expect((result.event as any).eventData.conflictingRunId).toBe(
testRunId
);
// No hook entity should be created
expect(result.hook).toBeUndefined();
});
it('should allow token reuse after hook is disposed', async () => {
const token = 'reusable-token-test';
// Create first hook with the token
await events.create(testRunId, {
eventType: 'hook_created' as const,
correlationId: 'hook_reuse_1',
eventData: { token },
});
// Dispose the first hook
await events.create(testRunId, {
eventType: 'hook_disposed' as const,
correlationId: 'hook_reuse_1',
});
// Create another run
const run2 = await createRun(events, {
deploymentId: 'deployment-789',
workflowName: 'test-workflow-3',
input: new Uint8Array(),
});
// Now creating a hook with the same token should succeed
const result = await events.create(run2.runId, {
eventType: 'hook_created' as const,
correlationId: 'hook_reuse_2',
eventData: { token },
});
expect(result.hook).toBeDefined();
expect(result.hook!.token).toBe(token);
});
});
});
describe('concurrent entity-creation races', () => {
let testRunId: string;
beforeEach(async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
testRunId = run.runId;
await updateRun(events, testRunId, 'run_started');
});
it('should reject concurrent step_created with the same correlationId', async () => {
// Two concurrent step_created calls with identical correlationIds
// (as produced by the snapshot runtime's deterministic ULIDs across
// concurrent VM invocations of the same resumption) must produce
// exactly one step_created event in the log. The unique partial
// index on workflow_events ensures the loser's INSERT raises a
// unique-violation, which storage translates to EntityConflictError
// for the runtime's existing dedup catch path.
const results = await Promise.allSettled([
createStep(events, testRunId, {
stepId: 'step_dup_1',
stepName: 'test-step',
input: new Uint8Array([1]),
}),
createStep(events, testRunId, {
stepId: 'step_dup_1',
stepName: 'test-step',
input: new Uint8Array([2]),
}),
]);
const fulfilled = results.filter((r) => r.status === 'fulfilled');
const rejected = results.filter((r) => r.status === 'rejected');
expect(fulfilled).toHaveLength(1);
expect(rejected).toHaveLength(1);
expect((rejected[0] as PromiseRejectedResult).reason).toMatchObject({
name: 'EntityConflictError',
});
// Verify only one step_created event exists in the log.
const evts = await events.list({
runId: testRunId,
pagination: {},
});
const stepCreated = evts.data.filter(
(e) =>
e.eventType === 'step_created' && e.correlationId === 'step_dup_1'
);
expect(stepCreated).toHaveLength(1);
});
it('should reject sequential duplicate step_created with EntityConflictError', async () => {
await createStep(events, testRunId, {
stepId: 'step_seq_dup',
stepName: 'test-step',
input: new Uint8Array(),
});
await expect(
createStep(events, testRunId, {
stepId: 'step_seq_dup',
stepName: 'test-step',
input: new Uint8Array(),
})
).rejects.toMatchObject({ name: 'EntityConflictError' });
});
it('should reject duplicate correlated workflow attr_set events', async () => {
await events.create(testRunId, {
eventType: 'attr_set',
correlationId: 'attr_dup_1',
eventData: {
changes: [{ key: 'phase', value: 'running' }],
writer: { type: 'workflow' },
},
});
await expect(
events.create(testRunId, {
eventType: 'attr_set',
correlationId: 'attr_dup_1',
eventData: {
changes: [{ key: 'phase', value: 'running' }],
writer: { type: 'workflow' },
},
})
).rejects.toMatchObject({ name: 'EntityConflictError' });
// A duplicate carrying *different* changes for the same correlationId
// must be rejected before touching the run snapshot — otherwise the
// materialized attributes would diverge from the event log.
await expect(
events.create(testRunId, {
eventType: 'attr_set',
correlationId: 'attr_dup_1',
eventData: {
changes: [{ key: 'phase', value: 'DIVERGED' }],
writer: { type: 'workflow' },
},
})
).rejects.toMatchObject({ name: 'EntityConflictError' });
expect((await runs.get(testRunId)).attributes?.phase).toBe('running');
const evts = await events.list({
runId: testRunId,
pagination: {},
});
expect(
evts.data.filter(
(event) =>
event.eventType === 'attr_set' &&
event.correlationId === 'attr_dup_1'
)
).toHaveLength(1);
});
it('should reject duplicate wait_created with EntityConflictError', async () => {
// Sequential duplicate wait_created — the wait_created insert path
// uses `INSERT ... onConflictDoNothing()` plus an existence check, so
// the second insert is silently dropped at the SQL level. The unique
// partial index on workflow_events still provides a stronger
// concurrent guarantee here, and the storage layer translates the
// resulting unique-violation into an EntityConflictError matching the
// step_created behavior.
await events.create(testRunId, {
eventType: 'wait_created',
correlationId: 'wait_seq_dup',
eventData: { resumeAt: new Date('2099-01-01') },
});
await expect(
events.create(testRunId, {
eventType: 'wait_created',
correlationId: 'wait_seq_dup',
eventData: { resumeAt: new Date('2099-01-02') },
})
).rejects.toMatchObject({ name: 'EntityConflictError' });
// Mirror the step_created test: assert exactly one wait_created
// event landed in the log, so a regression that allowed both
// inserts through would fail this test even if the second
// insert's translation to EntityConflictError still worked.
const evts = await events.list({
runId: testRunId,
pagination: {},
});
const waitCreated = evts.data.filter(
(e) =>
e.eventType === 'wait_created' && e.correlationId === 'wait_seq_dup'
);
expect(waitCreated).toHaveLength(1);
});
it('should reject duplicate same-hook hook_created with EntityConflictError, not hook_conflict', async () => {
// Regression test for https://github.com/vercel/workflow/issues/2283
//
// Duplicate processing of the *same* (runId, hookId, token) — e.g.
// queue redelivery or cross-process replay — must be idempotent.
// It must throw EntityConflictError (mirroring the step_created
// duplicate path) so the runtime's existing concurrent-replay catch
// path swallows it, and must NOT append a hook_conflict event that
// would later replay as a self-conflict HookConflictError.
const token = 'idempotent-token';
const hookId = 'hook_idem_1';
await createHook(events, testRunId, { hookId, token });
// Same runId, same hookId, same token — must be idempotent.
await expect(
events.create(testRunId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: { token },
})
).rejects.toMatchObject({ name: 'EntityConflictError' });
// No hook_conflict event should have been written to the log.
const evts = await events.list({
runId: testRunId,
pagination: {},
});
const hookCreatedEvents = evts.data.filter(
(e) => e.eventType === 'hook_created' && e.correlationId === hookId
);
const hookConflictEvents = evts.data.filter(
(e) => e.eventType === 'hook_conflict'
);
expect(hookCreatedEvents).toHaveLength(1);
expect(hookConflictEvents).toHaveLength(0);
});
it('should still emit hook_conflict for a different hookId reusing the same token in the same run', async () => {
// The idempotency guard must NOT mask genuine token conflicts — a
// different hookId reusing the same token (even in the same run)
// is still a real conflict.
const token = 'same-run-different-hook-token';
await createHook(events, testRunId, { hookId: 'hook_a', token });
const result = await events.create(testRunId, {
eventType: 'hook_created',
correlationId: 'hook_b',
eventData: { token },
});
expect(result.event.eventType).toBe('hook_conflict');
expect((result.event as any).eventData.conflictingRunId).toBe(testRunId);
expect(result.hook).toBeUndefined();
});
it('should still emit hook_conflict for the same hookId in a different run reusing the same token', async () => {
// The idempotency guard checks (runId, hookId) together — a
// different run reusing the same hookId (highly unlikely in
// practice, but a worthwhile boundary) must still produce a real
// hook_conflict.
const token = 'cross-run-same-hookid-token';
const hookId = 'hook_shared_id';
await createHook(events, testRunId, { hookId, token });
const otherRun = await createRun(events, {
deploymentId: 'deployment-other',
workflowName: 'other-workflow',
input: new Uint8Array(),
});
const result = await events.create(otherRun.runId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: { token },
});
expect(result.event.eventType).toBe('hook_conflict');
expect((result.event as any).eventData.conflictingRunId).toBe(testRunId);
expect(result.hook).toBeUndefined();
});
it('should recover an orphaned hook row that lacks a hook_created event', async () => {
// Crash-recovery regression: in `events.create`, the hook INSERT
// (line ~1185 of storage.ts) and the events INSERT (line ~1314)
// are not wrapped in a single transaction. If a process / DB
// interruption lands between them, the hook row exists but no
// `hook_created` event is in the log. The same-`(runId, hookId)`
// retry must not be treated as a "real duplicate" — that would
// throw EntityConflictError, which the runtime's concurrent-
// replay catch path would swallow, permanently leaving the run
// with a hook entity but no `hook_created` event in the log.
//
// The recovery path detects the missing event and completes the
// partial write: it skips re-inserting the hook row and lets the
// outer code path emit the `hook_created` event.
const token = 'orphaned-hook-row-token';
const hookId = 'hook_orphan_pg_1';
// Pre-seed an orphaned hook row that has no corresponding
// `hook_created` event in the events table.
await drizzle.insert(DrizzleSchema.hooks).values({
runId: testRunId,
hookId,
token,
ownerId: '',
projectId: '',
environment: '',
specVersion: SPEC_VERSION_CURRENT,
isWebhook: false,
isSystem: false,
});
// Sanity: the hook row exists but no hook_created event is in
// the log yet.
const preEvents = await events.list({
runId: testRunId,
pagination: {},
});
expect(
preEvents.data.filter((e) => e.eventType === 'hook_created').length
).toBe(0);
// Retry: must succeed and emit a hook_created event, NOT a
// hook_conflict event, and NOT throw EntityConflictError.
const result = await events.create(testRunId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: { token },
});
expect(result.event.eventType).toBe('hook_created');
expect(result.hook?.hookId).toBe(hookId);
const postEvents = await events.list({
runId: testRunId,
pagination: {},
});
const created = postEvents.data.filter(
(e) => e.eventType === 'hook_created' && e.correlationId === hookId
);
const conflicts = postEvents.data.filter(
(e) => e.eventType === 'hook_conflict'
);
expect(created).toHaveLength(1);
expect(conflicts).toHaveLength(0);
});
it('does not mutate an already-committed hook entity when a duplicate hook_created retry collides', async () => {
// Parallel to the world-local regression for karthikscale3's
// review on PR #2295. world-postgres uses
// `.insert(Schema.hooks).onConflictDoNothing()` so a duplicate
// hook_created retry'\''s hook INSERT is a no-op against an
// already-committed row — but this test guards against a
// future regression that adds an UPDATE/UPSERT or otherwise
// mutates the existing entity in the dedup path.
const token = 'no-mutate-on-duplicate-token-pg';
const hookId = 'hook_no_mutate_on_duplicate_pg';
const originalMetadata = encode({ v: 'a' }) as Uint8Array;
const retryMetadata = encode({ v: 'b' }) as Uint8Array;
// First write: original metadata + isWebhook: true.
const first = await events.create(testRunId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: {
token,
metadata: originalMetadata,
isWebhook: true,
},
});
expect(first.event.eventType).toBe('hook_created');
expect(first.hook?.isWebhook).toBe(true);
// Retry with DIFFERENT metadata and isWebhook.
await expect(
events.create(testRunId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: {
token,
metadata: retryMetadata,
isWebhook: false,
},
})
).rejects.toMatchObject({ name: 'EntityConflictError' });
// The hook entity still has the ORIGINAL metadata and
// isWebhook — the retry'\''s payload did NOT overwrite the
// already-committed entity.
const persisted = await hooks.get(hookId);
expect(persisted.isWebhook).toBe(true);
// Compare metadata as bytes since cbor round-trips through
// Buffer / Uint8Array.
expect(Buffer.from(persisted.metadata as Uint8Array)).toEqual(
Buffer.from(originalMetadata)
);
// Exactly one hook_created event in the log.
const evts = await events.list({
runId: testRunId,
pagination: { limit: 100 },
});
const hookCreated = evts.data.filter(
(e) => e.eventType === 'hook_created' && e.correlationId === hookId
);
expect(hookCreated).toHaveLength(1);
});
it('converges same-hook creation across concurrent calls to one event', async () => {
// Cross-worker convergence regression. The events table's
// partial unique index
// (workflow_events_entity_creation_unique on
// runId+correlationId+eventType for hook_created/step_created/
// wait_created) makes the events INSERT the durable
// convergence point — at most one `hook_created` event with
// the same `(runId, correlationId)` can land. The dedup branch
// can race with the original INSERT (both probe getHookByToken
// before the loser sees the event), but the outer events
// INSERT then raises 23505 (unique-violation) which is
// translated to EntityConflictError that the runtime's
// existing concurrent-replay catch path swallows. Net result:
// exactly one `hook_created` event per logical creation.
//
// This test is the world-postgres counterpart to the
// `converges same-hook creation across workers to one event`
// test in world-local, exercising true in-process concurrency
// since world-postgres has no per-process tag isolation.
const attempts = 25;
for (let i = 0; i < attempts; i++) {
const correlationId = `hook_pg_converge_${i}`;
const token = `token-pg-converge-${i}`;
await Promise.allSettled([
events.create(testRunId, {
eventType: 'hook_created',
correlationId,
eventData: { token },
}),
events.create(testRunId, {
eventType: 'hook_created',
correlationId,
eventData: { token },
}),
]);
}
const evts = await events.list({
runId: testRunId,
pagination: { limit: 1000 },
});
const created = evts.data.filter((e) => e.eventType === 'hook_created');
const conflicts = evts.data.filter(
(e) => e.eventType === 'hook_conflict'
);
expect(created).toHaveLength(attempts);
expect(conflicts).toHaveLength(0);
});
});
describe('step terminal state validation', () => {
let testRunId: string;
beforeEach(async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
testRunId = run.runId;
});
describe('completed step', () => {
it('should reject step_started on completed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_terminal_1',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
testRunId,
'step_terminal_1',
'step_completed',
{
result: new Uint8Array([1]),
}
);
await expect(
updateStep(events, testRunId, 'step_terminal_1', 'step_started')
).rejects.toThrow(/terminal/i);
});
it('should reject step_completed on already completed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_terminal_2',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
testRunId,
'step_terminal_2',
'step_completed',
{
result: new Uint8Array([1]),
}
);
await expect(
updateStep(events, testRunId, 'step_terminal_2', 'step_completed', {
result: new Uint8Array([2]),
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_failed on completed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_terminal_3',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
testRunId,
'step_terminal_3',
'step_completed',
{
result: new Uint8Array([1]),
}
);
await expect(
updateStep(events, testRunId, 'step_terminal_3', 'step_failed', {
error: 'Should not work',
})
).rejects.toThrow(/terminal/i);
});
});
describe('failed step', () => {
it('should reject step_started on failed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_failed_1',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, testRunId, 'step_failed_1', 'step_failed', {
error: 'Failed permanently',
});
await expect(
updateStep(events, testRunId, 'step_failed_1', 'step_started')
).rejects.toThrow(/terminal/i);
});
it('should reject step_completed on failed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_failed_2',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, testRunId, 'step_failed_2', 'step_failed', {
error: 'Failed permanently',
});
await expect(
updateStep(events, testRunId, 'step_failed_2', 'step_completed', {
result: new Uint8Array([3]),
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_failed on already failed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_failed_3',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, testRunId, 'step_failed_3', 'step_failed', {
error: 'Failed once',
});
await expect(
updateStep(events, testRunId, 'step_failed_3', 'step_failed', {
error: 'Failed again',
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_retrying on failed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_failed_retry',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
testRunId,
'step_failed_retry',
'step_failed',
{
error: 'Failed permanently',
}
);
await expect(
updateStep(events, testRunId, 'step_failed_retry', 'step_retrying', {
error: 'Retry attempt',
})
).rejects.toThrow(/terminal/i);
});
});
describe('step_retrying validation', () => {
it('should reject step_retrying on completed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_completed_retry',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
testRunId,
'step_completed_retry',
'step_completed',
{
result: new Uint8Array([1]),
}
);
await expect(
updateStep(
events,
testRunId,
'step_completed_retry',
'step_retrying',
{
error: 'Retry attempt',
}
)
).rejects.toThrow(/terminal/i);
});
});
});
describe('run terminal state validation', () => {
describe('completed run', () => {
it('should reject run_started on completed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
await expect(
updateRun(events, run.runId, 'run_started')
).rejects.toThrow(/terminal/i);
});
it('should reject run_failed on completed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
await expect(
updateRun(events, run.runId, 'run_failed', {
error: 'Should not work',
})
).rejects.toThrow(/terminal/i);
});
it('should reject run_cancelled on completed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
await expect(
events.create(run.runId, { eventType: 'run_cancelled' })
).rejects.toThrow(/terminal/i);
});
});
describe('failed run', () => {
it('should reject run_started on failed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_failed', { error: 'Failed' });
await expect(
updateRun(events, run.runId, 'run_started')
).rejects.toThrow(/terminal/i);
});
it('should reject run_completed on failed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_failed', { error: 'Failed' });
await expect(
updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([2]),
})
).rejects.toThrow(/terminal/i);
});
it('should reject run_cancelled on failed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_failed', { error: 'Failed' });
await expect(
events.create(run.runId, { eventType: 'run_cancelled' })
).rejects.toThrow(/terminal/i);
});
});
describe('cancelled run', () => {
it('should reject run_started on cancelled run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
await expect(
updateRun(events, run.runId, 'run_started')
).rejects.toThrow(/terminal/i);
});
it('should reject run_completed on cancelled run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
await expect(
updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([2]),
})
).rejects.toThrow(/terminal/i);
});
it('should reject run_failed on cancelled run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
await expect(
updateRun(events, run.runId, 'run_failed', {
error: 'Should not work',
})
).rejects.toThrow(/terminal/i);
});
});
});
describe('allowed operations on terminal runs', () => {
it('should allow step_completed on completed run for in-progress step', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// Create and start a step (making it in-progress)
await createStep(events, run.runId, {
stepId: 'step_in_progress',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, run.runId, 'step_in_progress', 'step_started');
// Complete the run while step is still running
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
// Should succeed - completing an in-progress step on a terminal run is allowed
const result = await updateStep(
events,
run.runId,
'step_in_progress',
'step_completed',
{ result: new Uint8Array([1]) }
);
expect(result.status).toBe('completed');
});
it('should allow step_failed on completed run for in-progress step', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// Create and start a step
await createStep(events, run.runId, {
stepId: 'step_in_progress_fail',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
run.runId,
'step_in_progress_fail',
'step_started'
);
// Complete the run
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
// Should succeed - failing an in-progress step on a terminal run is allowed
const result = await updateStep(
events,
run.runId,
'step_in_progress_fail',
'step_failed',
{ error: 'step failed' }
);
expect(result.status).toBe('failed');
});
it('should auto-delete hooks when run completes (postgres-specific behavior)', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// Create a hook
await createHook(events, run.runId, {
hookId: 'hook_auto_deleted',
token: 'test-token-dispose',
});
// Complete the run - this auto-deletes the hook
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
// The hook should no longer exist because run completion auto-deletes hooks
// This is intentional behavior to allow token reuse across runs
await expect(
events.create(run.runId, {
eventType: 'hook_disposed',
correlationId: 'hook_auto_deleted',
})
).rejects.toThrow(/not found/i);
});
});
describe('disallowed operations on terminal runs', () => {
it('should reject step_created on completed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
await expect(
createStep(events, run.runId, {
stepId: 'new_step',
stepName: 'test-step',
input: new Uint8Array(),
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_started on completed run for pending step', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// Create a step but don't start it
await createStep(events, run.runId, {
stepId: 'pending_step',
stepName: 'test-step',
input: new Uint8Array(),
});
// Complete the run
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
// Should reject - cannot start a pending step on a terminal run
await expect(
updateStep(events, run.runId, 'pending_step', 'step_started')
).rejects.toThrow(/terminal/i);
});
it('should reject hook_created on completed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
await expect(
createHook(events, run.runId, {
hookId: 'new_hook',
token: 'new-token',
})
).rejects.toThrow(/terminal/i);
});
it('should reject attr_set on completed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_completed', {
output: new Uint8Array([1]),
});
await expect(
events.create(run.runId, {
eventType: 'attr_set',
correlationId: 'attr_after_complete',
eventData: {
changes: [{ key: 'phase', value: 'too-late' }],
writer: { type: 'workflow' },
},
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_created on failed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_failed', { error: 'Failed' });
await expect(
createStep(events, run.runId, {
stepId: 'new_step_failed',
stepName: 'test-step',
input: new Uint8Array(),
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_created on cancelled run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
await expect(
createStep(events, run.runId, {
stepId: 'new_step_cancelled',
stepName: 'test-step',
input: new Uint8Array(),
})
).rejects.toThrow(/terminal/i);
});
it('should reject hook_created on failed run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await updateRun(events, run.runId, 'run_failed', { error: 'Failed' });
await expect(
createHook(events, run.runId, {
hookId: 'new_hook_failed',
token: 'new-token-failed',
})
).rejects.toThrow(/terminal/i);
});
it('should reject hook_created on cancelled run', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
await expect(
createHook(events, run.runId, {
hookId: 'new_hook_cancelled',
token: 'new-token-cancelled',
})
).rejects.toThrow(/terminal/i);
});
});
describe('idempotent operations', () => {
it('should allow run_cancelled on already cancelled run (idempotent)', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
// Should succeed - idempotent operation
const result = await events.create(run.runId, {
eventType: 'run_cancelled',
});
expect(result.run?.status).toBe('cancelled');
});
});
describe('step_retrying event handling', () => {
let testRunId: string;
beforeEach(async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
testRunId = run.runId;
});
it('should set step status to pending and record error', async () => {
await createStep(events, testRunId, {
stepId: 'step_retry_1',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, testRunId, 'step_retry_1', 'step_started');
// The `error` field is opaque SerializedData (Uint8Array) produced by
// dehydrateStepError. The storage layer persists it verbatim.
const serializedError = new Uint8Array([9, 9, 9]);
const result = await events.create(testRunId, {
eventType: 'step_retrying',
correlationId: 'step_retry_1',
eventData: {
error: serializedError,
retryAfter: new Date(Date.now() + 5000),
},
});
expect(result.step?.status).toBe('pending');
expect(result.step?.error).toEqual(serializedError);
expect(result.step?.retryAfter).toBeInstanceOf(Date);
});
it('should increment attempt when step_started is called after step_retrying', async () => {
await createStep(events, testRunId, {
stepId: 'step_retry_2',
stepName: 'test-step',
input: new Uint8Array(),
});
// First attempt
const started1 = await updateStep(
events,
testRunId,
'step_retry_2',
'step_started'
);
expect(started1.attempt).toBe(1);
// Retry
await events.create(testRunId, {
eventType: 'step_retrying',
correlationId: 'step_retry_2',
eventData: { error: 'Temporary failure' },
});
// Second attempt
const started2 = await updateStep(
events,
testRunId,
'step_retry_2',
'step_started'
);
expect(started2.attempt).toBe(2);
});
it('should reject step_retrying on completed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_retry_completed',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(
events,
testRunId,
'step_retry_completed',
'step_completed',
{
result: new Uint8Array([1]),
}
);
await expect(
events.create(testRunId, {
eventType: 'step_retrying',
correlationId: 'step_retry_completed',
eventData: { error: 'Should not work' },
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_retrying on failed step', async () => {
await createStep(events, testRunId, {
stepId: 'step_retry_failed',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, testRunId, 'step_retry_failed', 'step_failed', {
error: 'Permanent failure',
});
await expect(
events.create(testRunId, {
eventType: 'step_retrying',
correlationId: 'step_retry_failed',
eventData: { error: 'Should not work' },
})
).rejects.toThrow(/terminal/i);
});
});
describe('run cancellation with in-flight entities', () => {
it('should allow in-progress step to complete after run cancelled', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// Create and start a step
await createStep(events, run.runId, {
stepId: 'step_in_flight',
stepName: 'test-step',
input: new Uint8Array(),
});
await updateStep(events, run.runId, 'step_in_flight', 'step_started');
// Cancel the run
await events.create(run.runId, { eventType: 'run_cancelled' });
// Should succeed - completing an in-progress step is allowed
const result = await updateStep(
events,
run.runId,
'step_in_flight',
'step_completed',
{ result: new Uint8Array([1]) }
);
expect(result.status).toBe('completed');
});
it('should reject step_created after run cancelled', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
await events.create(run.runId, { eventType: 'run_cancelled' });
await expect(
createStep(events, run.runId, {
stepId: 'new_step_after_cancel',
stepName: 'test-step',
input: new Uint8Array(),
})
).rejects.toThrow(/terminal/i);
});
it('should reject step_started for pending step after run cancelled', async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
// Create a step but don't start it
await createStep(events, run.runId, {
stepId: 'pending_after_cancel',
stepName: 'test-step',
input: new Uint8Array(),
});
// Cancel the run
await events.create(run.runId, { eventType: 'run_cancelled' });
// Should reject - cannot start a pending step on a cancelled run
await expect(
updateStep(events, run.runId, 'pending_after_cancel', 'step_started')
).rejects.toThrow(/terminal/i);
});
});
describe('event ordering validation', () => {
let testRunId: string;
beforeEach(async () => {
const run = await createRun(events, {
deploymentId: 'deployment-123',
workflowName: 'test-workflow',
input: new Uint8Array(),
});
testRunId = run.runId;
});
it('should reject step_completed before step_created', async () => {
await expect(
events.create(testRunId, {
eventType: 'step_completed',
correlationId: 'nonexistent_step',
eventData: { result: new Uint8Array([1]) },
})
).rejects.toThrow(/not found/i);
});
it('should reject step_started before step_created', async () => {
await expect(
events.create(testRunId, {
eventType: 'step_started',
correlationId: 'nonexistent_step_started',
})
).rejects.toThrow(/not found/i);
});
it('should reject step_failed before step_created', async () => {
await expect(
events.create(testRunId, {
eventType: 'step_failed',
correlationId: 'nonexistent_step_failed',
eventData: { error: 'Failed' },
})
).rejects.toThrow(/not found/i);
});
it('should allow step_completed without step_started (instant completion)', async () => {
await createStep(events, testRunId, {
stepId: 'instant_complete',
stepName: 'test-step',
input: new Uint8Array(),
});
// Should succeed - instant completion without starting
const result = await updateStep(
events,
testRunId,
'instant_complete',
'step_completed',
{ result: new Uint8Array([1]) }
);
expect(result.status).toBe('completed');
});
it('should reject hook_disposed before hook_created', async () => {
await expect(
events.create(testRunId, {
eventType: 'hook_disposed',
correlationId: 'nonexistent_hook',
})
).rejects.toThrow(/not found/i);
});
it('should reject hook_received before hook_created', async () => {
await expect(
events.create(testRunId, {
eventType: 'hook_received',
correlationId: 'nonexistent_hook_received',
eventData: { payload: new Uint8Array() },
})
).rejects.toThrow(/not found/i);
});
});
describe('legacy/backwards compatibility', () => {
// Helper to create a legacy run directly in the database (bypassing events.create)
// Column mapping: id (runId), deployment_id, name (workflowName), spec_version, status, input
async function createLegacyRun(runId: string, specVersion: number | null) {
await pool.query(
`INSERT INTO workflow.workflow_runs (id, deployment_id, name, spec_version, status, input, created_at, updated_at)
VALUES ($1, 'legacy-deployment', 'legacy-workflow', $2, 'running', '[]'::jsonb, NOW(), NOW())`,
[runId, specVersion]
);
}
describe('legacy runs (specVersion < 2 or null)', () => {
it('should handle run_cancelled on legacy run with specVersion=1', async () => {
const runId = 'wrun_legacy_v1';
await createLegacyRun(runId, 1);
const result = await events.create(runId, {
eventType: 'run_cancelled',
});
// Legacy behavior: run is updated but event is not stored
expect(result.run?.status).toBe('cancelled');
expect(result.event).toBeUndefined();
});
it('should handle run_cancelled on legacy run with specVersion=null', async () => {
const runId = 'wrun_legacy_null';
await createLegacyRun(runId, null);
const result = await events.create(runId, {
eventType: 'run_cancelled',
});
// Legacy behavior: run is updated but event is not stored
expect(result.run?.status).toBe('cancelled');
expect(result.event).toBeUndefined();
});
it('should handle wait_completed on legacy run', async () => {
const runId = 'wrun_legacy_wait';
await createLegacyRun(runId, 1);
const result = await events.create(runId, {
eventType: 'wait_completed',
correlationId: 'wait_123',
eventData: { result: new Uint8Array([1]) },
} as any);
// Legacy behavior: event is stored but no entity mutation
expect(result.event).toBeDefined();
expect(result.event?.eventType).toBe('wait_completed');
expect(result.run).toBeUndefined();
});
it('should handle hook_received on legacy run', async () => {
const runId = 'wrun_legacy_hook_received';
await createLegacyRun(runId, 1);
const result = await events.create(runId, {
eventType: 'hook_received',
correlationId: 'hook_123',
eventData: { payload: new Uint8Array([1, 2, 3]) },
} as any);
// Legacy behavior: event is stored but no entity mutation
// (hooks exist via old system, not via events)
expect(result.event).toBeDefined();
expect(result.event?.eventType).toBe('hook_received');
expect(result.event?.correlationId).toBe('hook_123');
expect(result.hook).toBeUndefined();
});
it('should reject unsupported events on legacy runs', async () => {
const runId = 'wrun_legacy_unsupported';
await createLegacyRun(runId, 1);
// run_started is not supported for legacy runs
await expect(
events.create(runId, { eventType: 'run_started' })
).rejects.toThrow(/not supported for legacy runs/i);
// run_completed is not supported for legacy runs
await expect(
events.create(runId, {
eventType: 'run_completed',
eventData: { output: new Uint8Array([1]) },
})
).rejects.toThrow(/not supported for legacy runs/i);
// run_failed is not supported for legacy runs
await expect(
events.create(runId, {
eventType: 'run_failed',
eventData: { error: 'failed' },
})
).rejects.toThrow(/not supported for legacy runs/i);
});
it('should delete hooks when legacy run is cancelled', async () => {
const runId = 'wrun_legacy_hooks';
await createLegacyRun(runId, 1);
// Create a hook directly in the database for this run
await pool.query(
`INSERT INTO workflow.workflow_hooks (hook_id, run_id, token, owner_id, project_id, environment, created_at)
VALUES ('hook_legacy', $1, 'legacy-token', 'owner', 'project', 'test', NOW())`,
[runId]
);
// Verify hook exists
const hookBefore = await pool.query(
`SELECT hook_id FROM workflow.workflow_hooks WHERE hook_id = 'hook_legacy'`
);
expect(hookBefore.rows[0]).toBeDefined();
// Cancel the legacy run
await events.create(runId, { eventType: 'run_cancelled' });
// Hook should be deleted
const hookAfter = await pool.query(
`SELECT hook_id FROM workflow.workflow_hooks WHERE hook_id = 'hook_legacy'`
);
expect(hookAfter.rows[0]).toBeUndefined();
});
});
describe('newer runs (specVersion > current)', () => {
it('should reject events on runs with newer specVersion', async () => {
const runId = 'wrun_future';
// Create a run with a future spec version (higher than current)
await pool.query(
`INSERT INTO workflow.workflow_runs (id, deployment_id, name, spec_version, status, input, created_at, updated_at)
VALUES ($1, 'future-deployment', 'future-workflow', 999, 'running', '[]'::jsonb, NOW(), NOW())`,
[runId]
);
await expect(
events.create(runId, { eventType: 'run_started' })
).rejects.toThrow(/requires spec version 999/i);
});
});
describe('current version runs', () => {
it('should process events normally for current specVersion runs', async () => {
// Create run via events.create (gets current specVersion)
const run = await createRun(events, {
deploymentId: 'current-deployment',
workflowName: 'current-workflow',
input: new Uint8Array(),
});
// Should work normally
const result = await events.create(run.runId, {
eventType: 'run_started',
});
expect(result.run?.status).toBe('running');
expect(result.event?.eventType).toBe('run_started');
});
});
describe('legacy error column handling', () => {
// In the current event-sourced model, the `error` field on runs/steps
// is SerializedData (Uint8Array) produced by dehydrate*Error, stored in
// the `error_cbor` column. Legacy records written pre-serialization-
// pipeline (to the `error` text column) cannot be hydrated into the
// original thrown value and are surfaced as `undefined` on read.
it('should surface legacy errorJson field on runs as undefined', async () => {
const runId = 'wrun_legacy_error';
const inputCbor = encode(new Uint8Array());
await pool.query(
`INSERT INTO workflow.workflow_runs (id, deployment_id, name, spec_version, status, input_cbor, error, created_at, updated_at, completed_at)
VALUES ($1, 'deployment', 'workflow', 2, 'failed', $2, $3, NOW(), NOW(), NOW())`,
[runId, inputCbor, '{"message":"Legacy error","stack":"at foo()"}']
);
const run = await runs.get(runId);
expect(run.status).toBe('failed');
expect(run.error).toBeUndefined();
});
it('should surface legacy errorJson on steps as undefined', async () => {
const run = await createRun(events, {
deploymentId: 'deployment',
workflowName: 'workflow',
input: new Uint8Array(),
});
const inputCbor = encode(new Uint8Array());
await pool.query(
`INSERT INTO workflow.workflow_steps (run_id, step_id, step_name, status, input_cbor, error, attempt, created_at, updated_at, completed_at)
VALUES ($1, 'step_legacy_err', 'test-step', 'failed', $2, $3, 1, NOW(), NOW(), NOW())`,
[run.runId, inputCbor, '{"message":"Step error","stack":"at bar()"}']
);
const step = await steps.get(run.runId, 'step_legacy_err');
expect(step.status).toBe('failed');
expect(step.error).toBeUndefined();
});
});
});
});