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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
5 changes: 5 additions & 0 deletions .changeset/replay-lineage-execution-context.md
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': minor
---

Record the source run on replays: `recreateRunFromExisting` now stamps `replayedFromRunId` into the new run's `executionContext` (and `start` accepts a matching option) so tooling can surface a run as a replay and link back to its origin.
40 changes: 39 additions & 1 deletion packages/core/src/runtime/runs.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -12,14 +12,27 @@ import { afterEach, describe, expect, it, vi } from 'vitest';
// Mock version module to avoid missing generated file
vi.mock('../version.js', () => ({ version: '0.0.0-test' }));

// Stub `start` so recreateRunFromExisting can be tested for the options it
// forwards without exercising the full run-creation pipeline (serialization,
// telemetry, queueing).
vi.mock('./start.js', () => ({ start: vi.fn() }));

// Keep serialization real except for argument hydration, which needs a real
// serialized payload we don't have in these unit tests.
vi.mock('../serialization.js', async (importActual) => {
const actual = await importActual<typeof import('../serialization.js')>();
return { ...actual, hydrateWorkflowArguments: vi.fn().mockResolvedValue([]) };
});

import { registerSerializationClass } from '../class-serialization.js';
import {
dehydrateRunError,
dehydrateStepReturnValue,
hydrateStepReturnValue,
} from '../serialization.js';
import { Run } from './run.js';
import { reenqueueRun, wakeUpRun } from './runs.js';
import { recreateRunFromExisting, reenqueueRun, wakeUpRun } from './runs.js';
import { start } from './start.js';
import { setWorld } from './world.js';

function createMockWorld(
Expand DownExpand Up@@ -155,6 +168,31 @@ describe('reenqueueRun', () => {
});
});

describe('recreateRunFromExisting', () => {
afterEach(() => {
vi.mocked(start).mockReset();
});

it('forwards the source run id to start as replayedFromRunId', async () => {
const world = createMockWorld({
run: { runId: 'wrun_source', deploymentId: 'deploy_source' },
});
vi.mocked(start).mockResolvedValue({ runId: 'wrun_new' } as Run<unknown>);

const newRunId = await recreateRunFromExisting(world, 'wrun_source');

expect(newRunId).toBe('wrun_new');
expect(start).toHaveBeenCalledWith(
{ workflowId: 'test-workflow' },
expect.any(Array),
expect.objectContaining({
replayedFromRunId: 'wrun_source',
deploymentId: 'deploy_source',
})
);
});
});

describe('Run.exists', () => {
afterEach(() => {
setWorld(undefined as unknown as World);
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/runs.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -105,6 +105,7 @@ export async function recreateRunFromExisting(
deploymentId,
world,
specVersion,
replayedFromRunId: runId,
namespace: options.namespace,
}
);
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -770,6 +770,89 @@ describe('start', () => {
});
});

describe('replay lineage (executionContext.replayedFromRunId)', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;

const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
});

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);

setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn().mockResolvedValue('deploy_123'),
events: { create: mockEventsCreate },
queue: mockQueue,
});
});

afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

it('records replayedFromRunId in executionContext when provided', async () => {
const sourceRunId = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV';
await start(validWorkflow, [], { replayedFromRunId: sourceRunId });

expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
executionContext: expect.objectContaining({
replayedFromRunId: sourceRunId,
}),
}),
}),
expect.anything()
);
});

it('omits replayedFromRunId from executionContext when not provided', async () => {
await start(validWorkflow, []);

const eventData = mockEventsCreate.mock.calls[0]?.[1]?.eventData;
expect(eventData.executionContext).not.toHaveProperty(
'replayedFromRunId'
);
});

it('rejects a replayedFromRunId without the wrun_ prefix', async () => {
await expect(
start(validWorkflow, [], { replayedFromRunId: 'not-a-run-id' })
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a wrun_-prefixed value whose body is not a valid ULID', async () => {
await expect(
start(validWorkflow, [], {
replayedFromRunId: `wrun_${'x'.repeat(300)}`,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it('rejects a non-string replayedFromRunId', async () => {
await expect(
start(validWorkflow, [], {
// Types forbid this, but JS callers can still pass it.
replayedFromRunId: 12345 as unknown as string,
})
).rejects.toThrow(/replayedFromRunId must be a run ID/);
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('overload type inference', () => {
// Type-only assertions that don't execute start() at runtime.
// We use expectTypeOf on the function signature's return type directly.
Expand Down
31 changes: 31 additions & 0 deletions packages/core/src/runtime/start.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,6 +6,7 @@ import {
SPEC_VERSION_SUPPORTS_ATTRIBUTES,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_COMPRESSION,
workflowRunIdSchema,
} from '@workflow/world';
import { monotonicFactory } from 'ulid';
import { normalizeAttributeChanges } from '../attribute-changes.js';
Expand DownExpand Up@@ -90,6 +91,20 @@ export interface StartOptionsBase {
*/
allowReservedAttributes?: boolean;

/**
* The ID of an existing run this run is being replayed from, if any.
*
* Recorded on the new run's `executionContext` as `replayedFromRunId` so
* tooling (e.g. the dashboard runs list) can show that a run originated as
* a replay and link back to its source. Set automatically by
* {@link recreateRunFromExisting}; there's usually no reason to pass it
* directly.
*
* Must be a run ID: `wrun_` followed by a 26-char ULID. It's a foreign key
* to the source run, so `start()` validates the exact shape and rejects
* anything else rather than persist a lineage link that points at garbage.
*/
replayedFromRunId?: string;
/**
* Queue namespace of the target deployment. Scopes the workflow queue
* topic to `__{namespace}_wkf_workflow_*` (e.g. `'eve'`) instead of the
Expand DownExpand Up@@ -335,6 +350,19 @@ export async function start<TArgs extends unknown[], TResult>(
}
: {};

// `replayedFromRunId` is a foreign key to the source run; reject anything
// that isn't a real run ID so the lineage link can't point at garbage.
if (
opts.replayedFromRunId !== undefined &&
!workflowRunIdSchema.safeParse(opts.replayedFromRunId).success
) {
throw new WorkflowRuntimeError(
`replayedFromRunId must be a run ID (wrun_<ulid>); received ${JSON.stringify(
String(opts.replayedFromRunId).slice(0, 64)
)}.`
);
}

// Resolve encryption key for the new run. The runId has already been
// generated above (client-generated ULID) and will be used for both
// key derivation and the run_created event. The World implementation
Expand DownExpand Up@@ -372,6 +400,9 @@ export async function start<TArgs extends unknown[], TResult>(
traceCarrier,
workflowCoreVersion,
features: { encryption: !!encryptionKey },
...(opts.replayedFromRunId
? { replayedFromRunId: opts.replayedFromRunId }
: {}),
};

// Call events.create (run_created) and queue in parallel.
Expand Down
2 changes: 2 additions & 0 deletions packages/world/src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,6 +125,8 @@ export {
DEFAULT_TIMESTAMP_THRESHOLD_PAST_MS,
ulidToDate,
validateUlidTimestamp,
workflowRunIdSchema,
} from './ulid.js';
export type { WorkflowRunId } from './ulid.js';
export type * from './waits.js';
export { WaitSchema, WaitStatusSchema } from './waits.js';
10 changes: 10 additions & 0 deletions packages/world/src/ulid.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,6 +3,16 @@ import { z } from 'zod';

const UlidSchema = z.string().ulid();

/**
* A workflow run ID: the `wrun_` prefix followed by a 26-char ULID (minted
* client-side in core's `start()`). Validates the exact shape — prefix plus a
* well-formed ULID — rather than a loose length bound, so callers can't smuggle
* arbitrary strings through APIs that persist a run ID verbatim.
*/
export const workflowRunIdSchema = z.templateLiteral(['wrun_', z.ulid()]);

export type WorkflowRunId = z.infer<typeof workflowRunIdSchema>;

/**
* Default threshold for ULID timestamps in the past (24 hours).
*
Expand Down
Loading