From b1c3fddf2e0461a4218c3f8453ed295a98ee2f8c Mon Sep 17 00:00:00 2001 From: Tom Brandenburg Date: Fri, 2 Oct 2026 23:06:28 +0200 Subject: [PATCH] feat(temporal): support node activity heartbeat timeouts --- ...emporal-node-activity-heartbeat-timeout.md | 5 + packages/temporal/README.md | 4 + packages/temporal/activity-profiles.md | 50 ++++++- .../src/workflow/activity-profiles.ts | 3 + .../workflow/node-activity-options.test.ts | 13 ++ .../src/workflow/node-activity-options.ts | 1 + .../src/workflow/profile-validation.test.ts | 29 +++- .../src/workflow/profile-validation.ts | 10 +- .../test/fixtures/heartbeat-profiles.ts | 5 + .../temporal/test/fixtures/shell-executor.ts | 69 +++++++++ .../test/fixtures/workflows-with-heartbeat.ts | 4 + packages/temporal/test/heartbeat.test.ts | 132 ++++++++++++++++++ 12 files changed, 317 insertions(+), 8 deletions(-) create mode 100644 .changeset/temporal-node-activity-heartbeat-timeout.md create mode 100644 packages/temporal/test/fixtures/heartbeat-profiles.ts create mode 100644 packages/temporal/test/fixtures/shell-executor.ts create mode 100644 packages/temporal/test/fixtures/workflows-with-heartbeat.ts create mode 100644 packages/temporal/test/heartbeat.test.ts diff --git a/.changeset/temporal-node-activity-heartbeat-timeout.md b/.changeset/temporal-node-activity-heartbeat-timeout.md new file mode 100644 index 000000000..642da4db5 --- /dev/null +++ b/.changeset/temporal-node-activity-heartbeat-timeout.md @@ -0,0 +1,5 @@ +--- +'@workflowbuilder/temporal': minor +--- + +Add optional `heartbeatTimeout` to node activity profiles so heartbeating executors can detect worker loss sooner than their execution timeout. Document native Temporal heartbeat and cancellation handling for long-running executors. diff --git a/packages/temporal/README.md b/packages/temporal/README.md index eb2e45d25..10a636e50 100644 --- a/packages/temporal/README.md +++ b/packages/temporal/README.md @@ -101,6 +101,10 @@ Keep the export named `runWorkflow`: that is the name the client starts. Entries **Use one constant on both sides.** The plugin validates the map at `Worker.create`, so a bad profile fails the deploy. `createRunWorkflow` alone validates only on the first workflow activation, and a map handed to the plugin but not to `createRunWorkflow` gives you a green deploy with every node on the default profile. +### Long-running activities and cancellation + +Profiles accept optional `heartbeatTimeout`, using the same duration grammar as `startToCloseTimeout`. A long-running executor must send heartbeats itself and connect Temporal's cancellation signal to its running operation. See [heartbeat and cancellation examples](/packages/temporal/activity-profiles.md#heartbeats-and-cancellation), including a shell-backed step. Leaving the field unset preserves the existing defaults. + ### Failures and retries A node activity gets the attempts its profile allows. An executor can settle the question itself: diff --git a/packages/temporal/activity-profiles.md b/packages/temporal/activity-profiles.md index 4c237cda5..af151e632 100644 --- a/packages/temporal/activity-profiles.md +++ b/packages/temporal/activity-profiles.md @@ -29,9 +29,53 @@ Keep the export named `runWorkflow`: that is the name the client starts, and a t Entries are whole profiles rather than partials on purpose. A partial would let you set a timeout and silently drop the retry cap, and what Temporal falls back to is unlimited retries with backoff, which on a permanently failing model call is an unbounded bill. A node type with no entry resolves to `DEFAULT_NODE_ACTIVITY_PROFILE` and nothing else. +## Heartbeats and cancellation + +For a long-running step, add optional `heartbeatTimeout` to its whole profile and share the map between the plugin and `createRunWorkflow` as described below: + +```ts +export const nodeActivityProfiles: NodeActivityProfiles = { + 'my-product/shell': { + startToCloseTimeout: '45m', + retry: { maximumAttempts: 1 }, + heartbeatTimeout: '5s', + }, +}; +``` + +The executor runs inside a Temporal activity, so the native activity context is available even in nested async calls. No extra executor argument is needed: + +```ts +import { Context } from '@temporalio/activity'; +import { execFile } from 'node:child_process'; +import { promisify } from 'node:util'; + +const execute = promisify(execFile); + +async function executeShell() { + const activity = Context.current(); + activity.heartbeat(); + const timer = setInterval(() => activity.heartbeat(), 1000); + try { + const { stdout } = await execute('sh', ['-c', 'exec ./long-running-step'], { + signal: activity.cancellationSignal, + }); + return { output: stdout }; + } finally { + clearInterval(timer); + } +} +``` + +Only enable the timeout for executors that heartbeat. When heartbeats stop, Temporal times out that attempt and applies its retry policy; leaving the field unset means no heartbeat monitoring. Timer-driven heartbeats detect worker loss, not a hung command while the timer still runs: keep `startToCloseTimeout` and, where needed, an operation-specific idle timeout. + +Cancellation is cooperative. A running activity must heartbeat to receive workflow cancellation, then use `cancellationSignal` or the rejecting `cancelled` promise to stop its operation. The signal does not kill subprocess descendants automatically; the example uses shell `exec` to replace the wrapper, and commands that create children need their own process-tree cleanup. Do not catch an abort and report a successful result. + +The SDK throttles heartbeat RPCs: normally to 80% of `heartbeatTimeout`, capped by the worker's `maxHeartbeatThrottleInterval`; without a timeout it uses `defaultHeartbeatThrottleInterval` (30 seconds by default). Calling `heartbeat()` every second therefore does not imply one-second cancellation delivery. See [Temporal's activity context](https://typescript.temporal.io/api/classes/activity.Context) and [worker options](https://typescript.temporal.io/api/interfaces/worker.WorkerOptions#defaultheartbeatthrottleinterval). + ## Duration grammar -A `startToCloseTimeout` is a number followed by `ms`, `s`, `m`, `h` or `d`. Decimals are fine (`'1.5h'`). It has to fit a protobuf `Duration`, so anything under one nanosecond or over `'3652500d'` is out. Zero, negative values and exponent notation are rejected even though TypeScript's template literal type admits them: `'0s'` type-checks, and Temporal treats a zero timeout as unset and refuses to schedule the activity. +`startToCloseTimeout` and optional `heartbeatTimeout` are a number followed by `ms`, `s`, `m`, `h` or `d`. Decimals are fine (`'1.5h'`). Each has to fit a protobuf `Duration`, so anything under one nanosecond or over `'3652500d'` is out. Zero, negative values and exponent notation are rejected even though TypeScript's template literal type admits them: `'0s'` type-checks, but Temporal treats a zero timeout as unset. This grammar is narrower than Temporal's own, which parses durations with the `ms` package and also takes `'30 minutes'` or `'1 week'`. One documented form is deliberate. If you think in the wider grammar, convert before the value reaches this map. @@ -65,9 +109,9 @@ Profile **keys** cannot be validated inside the workflow: it runs in Temporal's ## What the profile check covers -It rejects a map whose entry is missing or `undefined`, whose `startToCloseTimeout` falls outside what a protobuf `Duration` carries, whose `retry.maximumAttempts` is not a positive integer that fits Temporal's `int32` field, or that carries any key beyond those two. +It rejects a map whose entry is missing or `undefined`, whose `startToCloseTimeout` or optional `heartbeatTimeout` falls outside the documented duration grammar and protobuf range, whose `retry.maximumAttempts` is not a positive integer that fits Temporal's `int32` field, or that carries any other key. -Unknown keys throw rather than being quietly dropped. Only those two fields are forwarded to `proxyActivities`, so a third would do nothing, and a map built from configuration gets no excess-property check from TypeScript to catch it at the keyboard. +Unknown keys throw rather than being quietly dropped. Only validated fields are forwarded to `proxyActivities`, and a map built from configuration gets no excess-property check from TypeScript to catch an unsupported field at the keyboard. Both bounds guard the same failure, where a value becomes its own opposite on the wire. Under one nanosecond a duration rounds to zero, which the server reads as unset and refuses, leaving the workflow task in a retry loop with nothing written to your database. A retry cap of `4294967296` arrives as `0`, which Temporal reads as unlimited. diff --git a/packages/temporal/src/workflow/activity-profiles.ts b/packages/temporal/src/workflow/activity-profiles.ts index 55b7a65c0..8fc68123c 100644 --- a/packages/temporal/src/workflow/activity-profiles.ts +++ b/packages/temporal/src/workflow/activity-profiles.ts @@ -12,6 +12,8 @@ export type DurationString = `${number}${'ms' | 's' | 'm' | 'h' | 'd'}`; export type ActivityProfile = { startToCloseTimeout: DurationString; retry: { maximumAttempts: number }; + // Requires heartbeats from the executor; omitted means no heartbeat monitoring. + heartbeatTimeout?: DurationString; }; // Only for the two frozen singletons below. Annotating them `ActivityProfile` would @@ -20,6 +22,7 @@ export type ActivityProfile = { type ReadonlyActivityProfile = { readonly startToCloseTimeout: DurationString; readonly retry: { readonly maximumAttempts: number }; + readonly heartbeatTimeout?: DurationString; }; // Node activities may call LLMs (minutes) — generous timeout, fewer retries to limit diff --git a/packages/temporal/src/workflow/node-activity-options.test.ts b/packages/temporal/src/workflow/node-activity-options.test.ts index a474120b3..189554839 100644 --- a/packages/temporal/src/workflow/node-activity-options.test.ts +++ b/packages/temporal/src/workflow/node-activity-options.test.ts @@ -47,6 +47,19 @@ describe('the shared default profiles', () => { }); describe('resolveNodeActivityOptions', () => { + it('forwards heartbeatTimeout only for a profile that supplies it', () => { + const profiles: NodeActivityProfiles = { + 'test/step': { startToCloseTimeout: '45m', retry: { maximumAttempts: 1 }, heartbeatTimeout: '5s' }, + }; + + expect(resolveNodeActivityOptions(node(), profiles)).toEqual({ + startToCloseTimeout: '45m', + retry: { maximumAttempts: 1 }, + heartbeatTimeout: '5s', + }); + expect(resolveNodeActivityOptions(node(), {})).not.toHaveProperty('heartbeatTimeout'); + }); + describe('timeouts and retries', () => { // Temporal's own default is unlimited retries with backoff, so a resolution bug // does not fail loudly, it spends. diff --git a/packages/temporal/src/workflow/node-activity-options.ts b/packages/temporal/src/workflow/node-activity-options.ts index 585ba078e..66c6cb1bb 100644 --- a/packages/temporal/src/workflow/node-activity-options.ts +++ b/packages/temporal/src/workflow/node-activity-options.ts @@ -41,6 +41,7 @@ export function resolveFromValidatedProfiles(node: BaseNode, profiles: NodeActiv const options: NodeActivityOptions = { startToCloseTimeout: profile.startToCloseTimeout, retry: { maximumAttempts: profile.retry.maximumAttempts }, + ...(profile.heartbeatTimeout === undefined ? {} : { heartbeatTimeout: profile.heartbeatTimeout }), }; // Not trusted from the caller: any consumer can build the workflow input. Temporal diff --git a/packages/temporal/src/workflow/profile-validation.test.ts b/packages/temporal/src/workflow/profile-validation.test.ts index 28e2fdd33..0d479a5cb 100644 --- a/packages/temporal/src/workflow/profile-validation.test.ts +++ b/packages/temporal/src/workflow/profile-validation.test.ts @@ -101,7 +101,7 @@ describe('assertNodeActivityProfiles', () => { }); describe('unknown keys', () => { - // The resolver forwards only the two validated fields, so anything else would be + // The resolver forwards only validated fields, so anything else would be // dropped in silence. A map built from configuration gets no excess-property check // from TypeScript either, which is how one arrives here in the first place. it('rejects a key the profile does not carry, naming it', () => { @@ -130,16 +130,39 @@ describe('assertNodeActivityProfiles', () => { startToCloseTimeout: '10m', retry: { maximumAttempts: 2 }, taskQueue: 'x', - heartbeatTimeout: '1m', + scheduleToCloseTimeout: '1m', }, }; expect(() => assertNodeActivityProfiles(extra as unknown as NodeActivityProfiles)).toThrow( - /has unknown keys "taskQueue", "heartbeatTimeout"/, + /has unknown keys "taskQueue", "scheduleToCloseTimeout"/, ); }); }); + describe('heartbeatTimeout', () => { + it('accepts a decimal duration and preserves it in the frozen snapshot', () => { + const live = { + 'test/step': { startToCloseTimeout: '10m', retry: { maximumAttempts: 1 }, heartbeatTimeout: '1.5s' }, + }; + const snapshot = freezeNodeActivityProfiles(live as NodeActivityProfiles); + live['test/step'].heartbeatTimeout = '0s'; + + expect(resolveNodeActivityOptions(node, snapshot).heartbeatTimeout).toBe('1.5s'); + }); + + it.each(['0s', '-5s', '1e3s', '5 seconds', '0.0000001ms', '3652501d', 5000, null])( + 'rejects invalid heartbeat timeout %s with its configuration path', + (heartbeatTimeout) => { + const map = { 'test/step': { startToCloseTimeout: '10m', retry: { maximumAttempts: 1 }, heartbeatTimeout } }; + + expect(() => assertNodeActivityProfiles(map as unknown as NodeActivityProfiles)).toThrow( + /nodeActivityProfiles\["test\/step"\]\.heartbeatTimeout must be a number followed by ms/, + ); + }, + ); + }); + it('rejects an entry whose value is undefined rather than defaulting it', () => { const undefinedEntry = { 'test/step': undefined }; diff --git a/packages/temporal/src/workflow/profile-validation.ts b/packages/temporal/src/workflow/profile-validation.ts index d006e635e..3c64156d1 100644 --- a/packages/temporal/src/workflow/profile-validation.ts +++ b/packages/temporal/src/workflow/profile-validation.ts @@ -42,11 +42,17 @@ function assertActivityProfile(nodeType: string, profile: ActivityProfile | unde ); } + if (profile.heartbeatTimeout !== undefined && !isDurationString(profile.heartbeatTimeout)) { + throw new TypeError( + `${path}.heartbeatTimeout must be a number followed by ms, s, m, h or d, such as '5s', and must fit a protobuf Duration: no shorter than one nanosecond ('0.000001ms') and no longer than '3652500d'. Got ${JSON.stringify(profile.heartbeatTimeout)}.`, + ); + } + assertNoUnknownKeys(path, profile, PROFILE_KEYS); assertNoUnknownKeys(`${path}.retry`, profile.retry, RETRY_KEYS); } -const PROFILE_KEYS: ReadonlySet = new Set(['startToCloseTimeout', 'retry']); +const PROFILE_KEYS: ReadonlySet = new Set(['startToCloseTimeout', 'retry', 'heartbeatTimeout']); const RETRY_KEYS: ReadonlySet = new Set(['maximumAttempts']); function assertNoUnknownKeys(path: string, value: object, allowed: ReadonlySet): void { @@ -54,7 +60,7 @@ function assertNoUnknownKeys(path: string, value: object, allowed: ReadonlySet JSON.stringify(key)).join(', ')}. A profile carries startToCloseTimeout and retry.maximumAttempts, nothing else.`, + `${path} has unknown ${unknown.length === 1 ? 'key' : 'keys'} ${unknown.map((key) => JSON.stringify(key)).join(', ')}. A profile carries startToCloseTimeout, retry.maximumAttempts and optional heartbeatTimeout, nothing else.`, ); } diff --git a/packages/temporal/test/fixtures/heartbeat-profiles.ts b/packages/temporal/test/fixtures/heartbeat-profiles.ts new file mode 100644 index 000000000..f7e86b271 --- /dev/null +++ b/packages/temporal/test/fixtures/heartbeat-profiles.ts @@ -0,0 +1,5 @@ +import type { NodeActivityProfiles } from '../../src/index'; + +export const heartbeatProfiles: NodeActivityProfiles = { + 'test/shell': { startToCloseTimeout: '30s', retry: { maximumAttempts: 1 }, heartbeatTimeout: '2s' }, +}; diff --git a/packages/temporal/test/fixtures/shell-executor.ts b/packages/temporal/test/fixtures/shell-executor.ts new file mode 100644 index 000000000..5d95ce53b --- /dev/null +++ b/packages/temporal/test/fixtures/shell-executor.ts @@ -0,0 +1,69 @@ +import { Context } from '@temporalio/activity'; +import { spawn } from 'node:child_process'; +import { once } from 'node:events'; + +import type { BaseNode, NodeExecutorRegistry } from '../../src/index'; + +export type ShellNode = BaseNode & { type: 'test/shell'; config: { seconds: number; heartbeat: boolean } }; + +export function createShellExecutor() { + const processes: { pid: number; closed: boolean; cancelled: boolean }[] = []; + const controllers = new Set(); + const outputs: string[] = []; + const inputs: Record> = {}; + + const executors: NodeExecutorRegistry = { + 'test/shell': async (node, context) => { + inputs[node.id] = { ...context.nodeOutputs }; + const activity = Context.current(); + const controller = new AbortController(); + controllers.add(controller); + const cancel = () => controller.abort(); + activity.cancellationSignal.addEventListener('abort', cancel, { once: true }); + if (activity.cancellationSignal.aborted) cancel(); + const child = spawn('sh', ['-c', `exec sleep ${node.config.seconds}`], { signal: controller.signal }); + const process = { pid: child.pid!, closed: false, cancelled: false }; + processes.push(process); + // An abort emits error before close; observe both so cleanup waits for the real process exit. + const failure = new Promise((resolve) => child.once('error', resolve)); + const closed = new Promise((resolve) => + child.once('close', () => { + process.closed = true; + resolve(); + }), + ); + activity.heartbeat({ nodeId: node.id }); + const timer = node.config.heartbeat ? setInterval(() => activity.heartbeat({ nodeId: node.id }), 100) : undefined; + try { + const exit = once(child, 'exit').then(([code]) => { + if (code !== 0) throw new Error(`shell exited with code ${String(code)}`); + }); + await Promise.race([ + exit, + failure.then((error) => { + throw error; + }), + ]); + outputs.push(node.id); + return { output: node.id }; + } finally { + clearInterval(timer); + process.cancelled = activity.cancellationSignal.aborted; + activity.cancellationSignal.removeEventListener('abort', cancel); + controller.abort(); + await closed; + controllers.delete(controller); + } + }, + }; + + return { + executors, + processes, + outputs, + inputs, + stop: () => { + for (const controller of controllers) controller.abort(); + }, + }; +} diff --git a/packages/temporal/test/fixtures/workflows-with-heartbeat.ts b/packages/temporal/test/fixtures/workflows-with-heartbeat.ts new file mode 100644 index 000000000..263199aa2 --- /dev/null +++ b/packages/temporal/test/fixtures/workflows-with-heartbeat.ts @@ -0,0 +1,4 @@ +import { createRunWorkflow } from '../../src/workflow/index'; +import { heartbeatProfiles } from './heartbeat-profiles'; + +export const runWorkflow = createRunWorkflow({ nodeActivityProfiles: heartbeatProfiles }); diff --git a/packages/temporal/test/heartbeat.test.ts b/packages/temporal/test/heartbeat.test.ts new file mode 100644 index 000000000..b3dfc6a58 --- /dev/null +++ b/packages/temporal/test/heartbeat.test.ts @@ -0,0 +1,132 @@ +import { CancelledFailure, WorkflowFailedError } from '@temporalio/client'; +import { TestWorkflowEnvironment } from '@temporalio/testing'; +import { Worker, bundleWorkflowCode } from '@temporalio/worker'; +import { randomUUID } from 'node:crypto'; +import { fileURLToPath } from 'node:url'; +import { afterAll, beforeAll, describe, expect, it } from 'vitest'; + +import { RUN_WORKFLOW_NAME, WorkflowBuilderPlugin, type WorkflowExecutionInput } from '../src/index'; +import { heartbeatProfiles } from './fixtures/heartbeat-profiles'; +import { waitUntil } from './fixtures/helpers'; +import { createRecordingStore } from './fixtures/recording-store'; +import { type ShellNode, createShellExecutor } from './fixtures/shell-executor'; + +describe('shell activity heartbeats on a real Temporal server', () => { + let env: TestWorkflowEnvironment; + let workflowBundle: { code: string }; + + beforeAll(async () => { + [workflowBundle, env] = await Promise.all([ + bundleWorkflowCode({ + workflowsPath: fileURLToPath(new URL('fixtures/workflows-with-heartbeat.ts', import.meta.url)), + }), + TestWorkflowEnvironment.createLocal(), + ]); + }, 300_000); + + afterAll(async () => { + await env?.teardown(); + }); + + async function setup(config: ShellNode['config']) { + const id = randomUUID(); + const store = createRecordingStore(); + const shell = createShellExecutor(); + const plugin = new WorkflowBuilderPlugin({ + store, + executors: shell.executors, + taskQueue: id, + nodeActivityProfiles: heartbeatProfiles, + }); + const worker = await Worker.create({ + connection: env.nativeConnection, + namespace: env.namespace, + taskQueue: id, + workflowBundle, + plugins: [plugin], + }); + const input: WorkflowExecutionInput = { + workflowId: id, + executionId: id, + triggerPayload: {}, + variables: {}, + global: {}, + definition: { + workflowId: id, + nodes: [ + { id: 'shell', type: 'test/shell', role: 'start', config }, + { id: 'after', type: 'test/shell', config: { seconds: 0, heartbeat: true } }, + ], + edges: [{ id: 'edge', sourceNodeId: 'shell', targetNodeId: 'after' }], + }, + }; + const handle = await env.client.workflow.start(RUN_WORKFLOW_NAME, { taskQueue: id, workflowId: id, args: [input] }); + return { worker, handle, store, shell }; + } + + it('heartbeats keep a shell step alive beyond the heartbeat timeout and deliver its output downstream', async () => { + const { worker, handle, store, shell } = await setup({ seconds: 4, heartbeat: true }); + await worker.runUntil(async () => { + try { + await handle.result(); + } finally { + shell.stop(); + } + }); + + expect(shell.outputs).toEqual(['shell', 'after']); + expect(shell.inputs.after.shell).toBe('shell'); + expect(shell.processes.every((process) => process.closed)).toBe(true); + expect(store.statuses.at(-1)?.status).toBe('completed'); + const history = await handle.fetchHistory(); + const scheduled = history.events?.find( + (event) => event.activityTaskScheduledEventAttributes?.activityType?.name === 'executeNode', + )?.activityTaskScheduledEventAttributes; + expect(scheduled?.heartbeatTimeout).toMatchObject({ seconds: expect.objectContaining({ low: 2 }) }); + }, 60_000); + + it('workflow cancellation aborts the running shell process before worker shutdown and skips downstream', async () => { + const { worker, handle, store, shell } = await setup({ seconds: 20, heartbeat: true }); + await worker.runUntil(async () => { + try { + await waitUntil(() => shell.processes.length === 1, 'the shell process'); + const started = Date.now(); + await handle.cancel(); + const failure: unknown = await handle.result().catch((error: unknown) => error); + expect(failure).toBeInstanceOf(WorkflowFailedError); + expect((failure as WorkflowFailedError).cause).toBeInstanceOf(CancelledFailure); + await waitUntil(() => shell.processes[0].closed, 'shell cancellation before worker shutdown', 10_000); + expect(shell.processes[0].cancelled).toBe(true); + expect(() => process.kill(shell.processes[0].pid, 0)).toThrow(); + expect(Date.now() - started).toBeLessThan(10_000); + } finally { + shell.stop(); + } + }); + + expect(shell.outputs).toEqual([]); + expect(shell.processes).toHaveLength(1); + expect(store.statuses.at(-1)?.status).toBe('cancelled'); + }, 60_000); + + it('stopped heartbeats fail with HEARTBEAT well before start-to-close and do not retry a side-effecting step', async () => { + const { worker, handle, store, shell } = await setup({ seconds: 20, heartbeat: false }); + await worker.runUntil(async () => { + try { + const started = Date.now(); + const failure: unknown = await handle.result().catch((error: unknown) => error); + expect(failure).toBeInstanceOf(WorkflowFailedError); + const history = await handle.fetchHistory(); + const timeout = history.events?.find((event) => event.activityTaskTimedOutEventAttributes); + expect(timeout?.activityTaskTimedOutEventAttributes?.failure?.timeoutFailureInfo?.timeoutType).toBe(4); + expect(store.statuses.at(-1)?.status).toBe('failed'); + expect(shell.processes).toHaveLength(1); + expect(shell.outputs).toEqual([]); + expect(Date.now() - started).toBeLessThan(15_000); + } finally { + shell.stop(); + } + }); + expect(shell.processes[0].closed).toBe(true); + }, 60_000); +});