Skip to content
Open
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/temporal-node-activity-heartbeat-timeout.md
Original file line number Diff line number Diff line change
@@ -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.
4 changes: 4 additions & 0 deletions packages/temporal/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
50 changes: 47 additions & 3 deletions packages/temporal/activity-profiles.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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.

Expand Down
3 changes: 3 additions & 0 deletions packages/temporal/src/workflow/activity-profiles.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
13 changes: 13 additions & 0 deletions packages/temporal/src/workflow/node-activity-options.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions packages/temporal/src/workflow/node-activity-options.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
29 changes: 26 additions & 3 deletions packages/temporal/src/workflow/profile-validation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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', () => {
Expand Down Expand Up @@ -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 };

Expand Down
10 changes: 8 additions & 2 deletions packages/temporal/src/workflow/profile-validation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,19 +42,25 @@ 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<string> = new Set(['startToCloseTimeout', 'retry']);
const PROFILE_KEYS: ReadonlySet<string> = new Set(['startToCloseTimeout', 'retry', 'heartbeatTimeout']);
const RETRY_KEYS: ReadonlySet<string> = new Set(['maximumAttempts']);

function assertNoUnknownKeys(path: string, value: object, allowed: ReadonlySet<string>): void {
const unknown = Object.keys(value).filter((key) => !allowed.has(key));
if (unknown.length === 0) return;

throw new TypeError(
`${path} has unknown ${unknown.length === 1 ? 'key' : 'keys'} ${unknown.map((key) => 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.`,
);
}

Expand Down
5 changes: 5 additions & 0 deletions packages/temporal/test/fixtures/heartbeat-profiles.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
import type { NodeActivityProfiles } from '../../src/index';

export const heartbeatProfiles: NodeActivityProfiles = {
'test/shell': { startToCloseTimeout: '30s', retry: { maximumAttempts: 1 }, heartbeatTimeout: '2s' },
};
69 changes: 69 additions & 0 deletions packages/temporal/test/fixtures/shell-executor.ts
Original file line number Diff line number Diff line change
@@ -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<AbortController>();
const outputs: string[] = [];
const inputs: Record<string, Record<string, unknown>> = {};

const executors: NodeExecutorRegistry<ShellNode> = {
'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<Error>((resolve) => child.once('error', resolve));
const closed = new Promise<void>((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();
},
};
}
4 changes: 4 additions & 0 deletions packages/temporal/test/fixtures/workflows-with-heartbeat.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
import { createRunWorkflow } from '../../src/workflow/index';
import { heartbeatProfiles } from './heartbeat-profiles';

export const runWorkflow = createRunWorkflow({ nodeActivityProfiles: heartbeatProfiles });
Loading