diff --git a/.changeset/start-explicit-deployment-without-current.md b/.changeset/start-explicit-deployment-without-current.md new file mode 100644 index 0000000000..52f0d44aee --- /dev/null +++ b/.changeset/start-explicit-deployment-without-current.md @@ -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. diff --git a/packages/core/src/runtime/start.test.ts b/packages/core/src/runtime/start.test.ts index 61c4cd0f03..5c59733c20 100644 --- a/packages/core/src/runtime/start.test.ts +++ b/packages/core/src/runtime/start.test.ts @@ -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, @@ -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; + let mockQueue: ReturnType; + let mockGetDeploymentId: ReturnType; + + const validWorkflow = Object.assign(() => Promise.resolve('result'), { + workflowId: 'test-workflow', + }); + + const createWorld = (overrides: Record = {}) => + ({ + 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', diff --git a/packages/core/src/runtime/start.ts b/packages/core/src/runtime/start.ts index 910611acb9..19ef149507 100644 --- a/packages/core/src/runtime/start.ts +++ b/packages/core/src/runtime/start.ts @@ -175,34 +175,59 @@ export async function start( }); 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) } + ); } }