mirror of
https://github.com/vercel/workflow.git
synced 2026-09-14 19:59:43 +08:00
f2a7bdeb0a
* 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>
2867 lines
94 KiB
TypeScript
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();
|
|
});
|
|
});
|
|
});
|
|
});
|