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/start-explicit-deployment-without-current.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

`start()` with an explicit `deploymentId` (and so `recreateRunFromExisting`, i.e. Replay Run) no longer fails in a process that is not itself a deployment; it takes the cross-deployment path instead.
117 changes: 117 additions & 0 deletions packages/core/src/runtime/start.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import {
SPEC_VERSION_LEGACY,
SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT,
SPEC_VERSION_SUPPORTS_EVENT_SOURCING,
type World,
} from '@workflow/world';
import {
afterEach,
Expand Down Expand Up @@ -478,6 +479,122 @@ describe('start', () => {
});
});

// A process that is not itself a deployment (e.g. `workflow web` replaying
// a production run via recreateRunFromExisting, or a recovery script) has
// no current deployment: the Vercel world's getDeploymentId() throws when
// VERCEL_DEPLOYMENT_ID is unset.
describe('without a current deployment', () => {
let mockEventsCreate: ReturnType<typeof vi.fn>;
let mockQueue: ReturnType<typeof vi.fn>;
let mockGetDeploymentId: ReturnType<typeof vi.fn>;

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

const createWorld = (overrides: Record<string, unknown> = {}) =>
({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: mockGetDeploymentId,
events: { create: mockEventsCreate },
queue: mockQueue,
...overrides,
}) as unknown as World;

beforeEach(() => {
mockEventsCreate = vi.fn().mockImplementation((runId) => {
return Promise.resolve({
run: { runId: runId ?? 'wrun_test123', status: 'pending' },
});
});
mockQueue = vi.fn().mockResolvedValue(undefined);
mockGetDeploymentId = vi
.fn()
.mockRejectedValue(
new Error('Starting a workflow run requires VERCEL_DEPLOYMENT_ID')
);
});

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

it('starts a run targeted at an explicit deploymentId', async () => {
const run = await start(validWorkflow, [], {
deploymentId: 'dpl_target_456',
world: createWorld(),
});

expect(run.runId).toMatch(/^wrun_/);
expect(mockGetDeploymentId).toHaveBeenCalledTimes(1);
expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({
eventType: 'run_created',
eventData: expect.objectContaining({
deploymentId: 'dpl_target_456',
}),
}),
expect.anything()
);
expect(mockQueue).toHaveBeenCalledWith(
expect.any(String),
expect.any(Object),
expect.objectContaining({ deploymentId: 'dpl_target_456' })
);
});

it('treats the explicit deploymentId as a cross-deployment target', async () => {
// With a probe channel, the start must probe the target's capabilities
// rather than assume it runs this SDK (the same-deployment shortcut).
// Failing the probe's enqueue ends it at once instead of polling out
// its timeout; the start then falls back to the universally readable
// formats.
const mockReadFromStream = vi.fn();
mockQueue.mockImplementation((queueName: string) =>
queueName.endsWith('health_check')
? Promise.reject(new Error('probe unavailable'))
: Promise.resolve(undefined)
);

await start(validWorkflow, [], {
deploymentId: 'dpl_target_456',
world: createWorld({ readFromStream: mockReadFromStream }),
});

expect(mockQueue).toHaveBeenCalledWith(
expect.stringContaining('health_check'),
expect.anything(),
expect.objectContaining({ deploymentId: 'dpl_target_456' })
);
expect(mockEventsCreate).toHaveBeenCalledWith(
expect.stringMatching(/^wrun_/),
expect.objectContaining({ eventType: 'run_created' }),
expect.anything()
);
});

it('still fails when no deploymentId is given', async () => {
await expect(
start(validWorkflow, [], { world: createWorld() })
).rejects.toThrow('requires VERCEL_DEPLOYMENT_ID');
expect(mockEventsCreate).not.toHaveBeenCalled();
});

it("still fails for deploymentId: 'latest'", async () => {
const mockResolveLatest = vi.fn().mockResolvedValue('dpl_latest');

await expect(
start(validWorkflow, [], {
deploymentId: 'latest',
world: createWorld({ resolveLatestDeploymentId: mockResolveLatest }),
})
).rejects.toThrow('requires VERCEL_DEPLOYMENT_ID');
expect(mockResolveLatest).not.toHaveBeenCalled();
expect(mockEventsCreate).not.toHaveBeenCalled();
});
});

describe('resilient start (run_created failure)', () => {
const validWorkflow = Object.assign(() => Promise.resolve('result'), {
workflowId: 'test-workflow',
Expand Down
79 changes: 52 additions & 27 deletions packages/core/src/runtime/start.ts
Original file line number Diff line number Diff line change
Expand Up @@ -175,34 +175,59 @@ export async function start<TArgs extends unknown[], TResult>(
});

const world = opts?.world ?? getWorld();
const currentDeploymentId = await world.getDeploymentId();
let deploymentId = opts.deploymentId ?? currentDeploymentId;

// When 'latest' is requested, resolve the actual latest deployment ID
// for the current deployment's environment (same production target or
// same git branch for preview deployments).
//
// Resolving 'latest' only means something in worlds with atomic,
// immutable deployments (e.g. Vercel), which implement
// resolveLatestDeploymentId(). Worlds without that concept (local dev,
// self-hosted Postgres) have nothing to resolve between, so rather than
// fail a run that works fine on Vercel, we warn and fall back to the
// current deployment — making 'latest' an effective no-op there.
if (deploymentId === 'latest') {
if (world.resolveLatestDeploymentId) {
deploymentId = await world.resolveLatestDeploymentId();
} else {
// Warn once per process — see hasWarnedLatestNoOp above.
if (!hasWarnedLatestNoOp) {
hasWarnedLatestNoOp = true;
runtimeLogger.warn(
"deploymentId: 'latest' has no effect in this world and was ignored. " +
'It is only supported by worlds with atomic deployments, such as Vercel. ' +
'The run will target the current deployment.',
{ currentDeploymentId }
);
// `undefined` when this process is not itself a deployment and the
// caller named a concrete target; see below.
let currentDeploymentId: string | undefined;
let deploymentId: string;
if (opts.deploymentId === undefined || opts.deploymentId === 'latest') {
// Defaulting the target and resolving 'latest' both need the current
// deployment, so a world that cannot report one fails the start here.
const current = await world.getDeploymentId();
currentDeploymentId = current;
deploymentId = opts.deploymentId ?? current;

// When 'latest' is requested, resolve the actual latest deployment ID
// for the current deployment's environment (same production target or
// same git branch for preview deployments).
//
// Resolving 'latest' only means something in worlds with atomic,
// immutable deployments (e.g. Vercel), which implement
// resolveLatestDeploymentId(). Worlds without that concept (local dev,
// self-hosted Postgres) have nothing to resolve between, so rather than
// fail a run that works fine on Vercel, we warn and fall back to the
// current deployment — making 'latest' an effective no-op there.
if (deploymentId === 'latest') {
if (world.resolveLatestDeploymentId) {
deploymentId = await world.resolveLatestDeploymentId();
} else {
// Warn once per process — see hasWarnedLatestNoOp above.
if (!hasWarnedLatestNoOp) {
hasWarnedLatestNoOp = true;
runtimeLogger.warn(
"deploymentId: 'latest' has no effect in this world and was ignored. " +
'It is only supported by worlds with atomic deployments, such as Vercel. ' +
'The run will target the current deployment.',
{ currentDeploymentId }
);
}
deploymentId = current;
}
deploymentId = currentDeploymentId;
}
} else {
// With a concrete target the current deployment only decides whether
// the start is same-deployment. A process that is not itself a
// deployment (e.g. `workflow web` replaying a production run, or a
// recovery script) must still be able to start one, so an unavailable
// current deployment means "not the target": take the
// cross-deployment probe path below instead of failing.
deploymentId = opts.deploymentId;
try {
currentDeploymentId = await world.getDeploymentId();
} catch (err) {
runtimeLogger.debug(
'Current deployment is unavailable; starting as a cross-deployment run',
{ deploymentId, error: String(err) }
);
}
}

Expand Down