Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .changeset/world-local-hook-cache-rebuild.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
---
'@workflow/world-local': patch
'@workflow/world': patch
---

Keep local hooks reachable after a crash or restart by rebuilding lost hook cache files from committed hook creation events, preventing active hook tokens from being reused.
157 changes: 142 additions & 15 deletions packages/world-local/src/storage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1304,24 +1304,24 @@ describe('Storage', () => {
// never be mistaken for the complete delta.
await updateRun(storage, testRunId, 'run_started');

// Open many hooks BEFORE the cursor so they are not part of the delta;
// their in-band deliveries (below) are.
const HOOK_COUNT = 25; // > paginatedFileSystemQuery default limit (20)
for (let i = 0; i < HOOK_COUNT; i++) {
await createHook(storage, testRunId, {
hookId: `corr_hook_${i}`,
token: `tok_${i}`,
});
}
await createHook(storage, testRunId, {
hookId: 'corr_delta_page_hook',
token: 'tok_delta_page_hook',
});
const sinceCursor = await currentCursor();

// A burst of in-band hook deliveries lands while the step runs, then
// the sequential step itself — together far more than one page.
for (let i = 0; i < HOOK_COUNT; i++) {
// A burst of in-band hook deliveries lands while the step runs. One
// hook is enough here; the assertion is about delta pagination, not
// creating many distinct hook tokens.
const DELTA_FILLER_EVENT_COUNT = 21; // > default page limit (20)
for (let i = 0; i < DELTA_FILLER_EVENT_COUNT; i++) {
await storage.events.create(testRunId, {
eventType: 'hook_received' as const,
correlationId: `corr_hook_${i}`,
eventData: { token: `tok_${i}`, payload: new Uint8Array([i]) },
correlationId: 'corr_delta_page_hook',
eventData: {
token: 'tok_delta_page_hook',
payload: new Uint8Array([i]),
},
});
}
await createStep(storage, testRunId, {
Expand Down Expand Up @@ -1351,7 +1351,9 @@ describe('Storage', () => {

expect(result.hasMore).toBe(true);
expect(firstPage.hasMore).toBe(true);
expect(result.events?.length).toBeLessThan(HOOK_COUNT + 3);
expect(result.events?.length).toBeLessThan(
DELTA_FILLER_EVENT_COUNT + 3
);
expect(result.events?.map((e) => e.eventId)).toEqual(
firstPage.data.map((e) => e.eventId)
);
Expand Down Expand Up @@ -3323,6 +3325,131 @@ describe('Storage', () => {
).toHaveLength(1);
});

it('rebuilds missing hook caches from a committed hook_created event', async () => {
// Regression for #2339: once hook_created is committed to the event log,
// the hook entity and token claim are cache files. If both are missing
// after a crash or upgrade, a normal hook read should rebuild them from
// the persisted event instead of treating the hook/token as gone.
const metadata = new Uint8Array([0xee]);
const hookId = 'hook_event_log_rebuild';
const token = 'event-log-rebuild-token';

const created = await storage.events.create(testRunId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: { token, metadata, isWebhook: true },
});
expect(created.event.eventType).toBe('hook_created');

const hookPath = path.join(testDir, 'hooks', `${hookId}.json`);
const tokenClaimPath = path.join(
testDir,
'hooks',
'tokens',
`${hashToken(token)}.json`
);
await fs.unlink(hookPath);
await fs.unlink(tokenClaimPath);
await fs.writeFile(
path.join(testDir, 'events', 'wrun_malformed-event.json'),
'{'
);

const conflict = await storage.events.create(testRunId, {
eventType: 'hook_created',
correlationId: 'hook_event_log_rebuild_conflict',
eventData: { token },
});
expect(conflict.event.eventType).toBe('hook_conflict');
expect((conflict.event as any).eventData.conflictingRunId).toBe(
testRunId
);

await fs.unlink(hookPath);
await fs.unlink(tokenClaimPath);

await expect(storage.hooks.get(hookId)).resolves.toMatchObject({
runId: testRunId,
hookId,
token,
metadata,
isWebhook: true,
});

const claim = JSON.parse(await fs.readFile(tokenClaimPath, 'utf8'));
expect(claim).toMatchObject({
runId: testRunId,
hookId,
eventId: created.event.eventId,
});
});

it('preserves legacy webhook default when rebuilding a hook without isWebhook', async () => {
const metadata = new Uint8Array([0xab]);
const hookId = 'hook_legacy_webhook_default';
const token = 'legacy-webhook-default-token';
const created = await storage.events.create(testRunId, {
eventType: 'hook_created',
correlationId: hookId,
eventData: { token, metadata },
});
expect(created.event.eventType).toBe('hook_created');

await fs.unlink(path.join(testDir, 'hooks', `${hookId}.json`));
await fs.unlink(
path.join(testDir, 'hooks', 'tokens', `${hashToken(token)}.json`)
);

await expect(storage.hooks.get(hookId)).resolves.toMatchObject({
hookId,
token,
metadata,
isWebhook: true,
});
});

it('does not rebuild a hook for a run already marked terminal', async () => {
const hookId = 'hook_terminal_run_cache';
const token = 'terminal-run-cache-token';
await createHook(storage, testRunId, { hookId, token });

const hookPath = path.join(testDir, 'hooks', `${hookId}.json`);
const tokenClaimPath = path.join(
testDir,
'hooks',
'tokens',
`${hashToken(token)}.json`
);
await fs.unlink(hookPath);
await fs.unlink(tokenClaimPath);

const run = await storage.runs.get(testRunId);
await writeJSON(
path.join(testDir, 'runs', `${testRunId}.json`),
{
...run,
status: 'cancelled',
completedAt: new Date(),
updatedAt: new Date(),
},
{ overwrite: true }
);

const nextRun = await createRun(storage, {
deploymentId: 'deployment-next',
workflowName: 'next-workflow',
input: new Uint8Array(),
});

const created = await storage.events.create(nextRun.runId, {
eventType: 'hook_created',
correlationId: 'hook_terminal_run_cache_next',
eventData: { token },
});
expect(created.event.eventType).toBe('hook_created');
expect(created.hook?.runId).toBe(nextRun.runId);
});

it('repairs an event-first orphan via the legacy-claim probe path', async () => {
// Same crash window as the test above, but exercised through
// the legacy-claim branch: the claim file lacks `eventId` (as
Expand Down
41 changes: 22 additions & 19 deletions packages/world-local/src/storage/events-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,9 @@ import {
EventSchema,
HookSchema,
isLegacySpecVersion,
isTerminalRunEventType,
isTerminalStepStatus,
isTerminalWorkflowRunStatus,
requiresNewerWorld,
SPEC_VERSION_CURRENT,
StepSchema,
Expand Down Expand Up @@ -57,7 +60,10 @@ import {
hookRecoveryMarkerPath,
monotonicUlid,
} from './helpers.js';
import { deleteAllHooksForRun } from './hooks-storage.js';
import {
deleteAllHooksForRun,
rebuildLiveHookByTokenFromEventLog,
} from './hooks-storage.js';
import { handleLegacyEvent } from './legacy.js';
import { withRunFileLock } from './runs-storage.js';

Expand Down Expand Up @@ -259,7 +265,7 @@ async function repairHookEntityFromPersistedEvent(
environment: 'local',
createdAt: persistedEvent.createdAt,
specVersion: persistedEvent.specVersion,
isWebhook: eventData.isWebhook ?? false,
isWebhook: eventData.isWebhook ?? true,
isSystem: eventData.isSystem ?? false,
};
await writeExclusive(
Expand Down Expand Up @@ -485,11 +491,7 @@ export function createEventsStorage(
): void {
// Terminal runs release their cached history so a long-lived dev
// server doesn't retain completed runs forever.
if (
event.eventType === 'run_completed' ||
event.eventType === 'run_failed' ||
event.eventType === 'run_cancelled'
) {
if (isTerminalRunEventType(event.eventType)) {
clearRunCache(event.runId);
return;
}
Expand Down Expand Up @@ -626,14 +628,6 @@ export function createEventsStorage(
// specVersion is always sent by the runtime, but we provide a fallback for safety
const effectiveSpecVersion = data.specVersion ?? SPEC_VERSION_CURRENT;

// Helper to check if run is in terminal state
const isRunTerminal = (status: string) =>
['completed', 'failed', 'cancelled'].includes(status);

// Helper to check if step is in terminal state
const isStepTerminal = (status: string) =>
['completed', 'failed', 'cancelled'].includes(status);

// Get current run state for validation (if not creating a new run)
// Skip run validation for step_completed and step_retrying - they only operate
// on running steps, and running steps are always allowed to modify regardless
Expand Down Expand Up @@ -801,7 +795,7 @@ export function createEventsStorage(
(data.eventData as { input?: unknown }).input !== undefined;

// Run terminal state validation
if (currentRun && isRunTerminal(currentRun.status)) {
if (currentRun && isTerminalWorkflowRunStatus(currentRun.status)) {
const runTerminalEvents = [
'run_started',
'run_completed',
Expand Down Expand Up @@ -916,14 +910,14 @@ export function createEventsStorage(
// the lazy-start path (no step yet) — there is nothing terminal to
// guard against in that case, so these checks are skipped.
if (validatedStep) {
if (isStepTerminal(validatedStep.status)) {
if (isTerminalStepStatus(validatedStep.status)) {
throw new EntityConflictError(
`Cannot modify step in terminal state "${validatedStep.status}"`
);
}

// On terminal runs: only allow completing/failing in-progress steps
if (currentRun && isRunTerminal(currentRun.status)) {
if (currentRun && isTerminalWorkflowRunStatus(currentRun.status)) {
if (validatedStep.status !== 'running') {
throw new RunExpiredError(
`Cannot modify non-running step on run in terminal state "${currentRun.status}"`
Expand Down Expand Up @@ -1429,7 +1423,7 @@ export function createEventsStorage(
StepSchema,
tag
);
if (freshStep && isStepTerminal(freshStep.status)) {
if (freshStep && isTerminalStepStatus(freshStep.status)) {
throw new EntityConflictError(
`Cannot modify step in terminal state "${freshStep.status}"`
);
Expand Down Expand Up @@ -1569,6 +1563,15 @@ export function createEventsStorage(
'tokens',
`${hashToken(hookData.token)}.json`
);
// When the claim is absent, the event log is the only durable source
// that can distinguish a first hook from a crash-lost token cache.
if (!(await readHookTokenClaim(constraintPath))) {
await rebuildLiveHookByTokenFromEventLog(
basedir,
hookData.token,
tag
);
}
// Persist `eventId` in the claim so concurrent / cross-
// process retries can converge on a single canonical
// `hook_created` event path. See the recovery comment
Expand Down
Loading
Loading