Skip to content

Commit b9627be

Browse files
feat(surface,sdk,kernel): deterministic-step timeout option (f.run with timeout) (#343)
1 parent 5ab55e0 commit b9627be

16 files changed

Lines changed: 290 additions & 30 deletions

File tree

‎docs/SURFACE.md‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -426,6 +426,20 @@ example selects `stdout_tail`. JSON-emitting agent/LLM outputs use their declare
426426
value shape directly. Use matching SDK and kernel builds for input bindings;
427427
older kernels refuse the new field.
428428

429+
### Command timeouts
430+
431+
`f.run(command, { timeout?: string | number })` gives each command its own
432+
lease. The default is **30 seconds**. Use `await f.run(command, { timeout: '5m' })`
433+
or `{ timeout: 300000 }` for longer work. Strings accept `ms`, `s`, and `m`;
434+
the resolved value must be a positive whole number of milliseconds.
435+
436+
The hard ceiling is **15 minutes** (900000 ms), inclusive. A larger timeout
437+
is refused during step compilation, before dispatch, with `lease_exceeded`;
438+
malformed durations are refused with `timeout_invalid`. On reaching its timeout,
439+
the kernel kills the command's process group and journals `completionReason: timeout`;
440+
`f.run` refuses with code `lease_exceeded`. The override applies only to that
441+
invocation; calls without options retain the default.
442+
429443
### The authored operation lifecycle
430444

431445
An authored TypeScript body reaches `done()` only if every step it created was

‎kernel/relayflowd-core/src/machine.rs‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -280,7 +280,13 @@ fn start_actions(state: &RunState, step: &StepSpec, attempt: u32, now_ms: i64) -
280280
"unassigned"
281281
};
282282
let lease_id = deterministic_ulid(&state.run_id, &step.id, attempt, now_ms, "lease");
283-
let lease_deadline_ms = now_ms.saturating_add(LEASE_DURATION_MS);
283+
let lease_duration_ms = match &step.kind {
284+
StepKind::Deterministic {
285+
lease_ms: Some(ms), ..
286+
} => i64::try_from(*ms).unwrap_or(i64::MAX),
287+
_ => LEASE_DURATION_MS,
288+
};
289+
let lease_deadline_ms = now_ms.saturating_add(lease_duration_ms);
284290
let started = JournalEntry::new(
285291
EntryType::StepAttemptStarted,
286292
state.run_id.clone(),

‎kernel/relayflowd-core/src/machine/tests.rs‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -596,3 +596,43 @@ fn every_reason_label_matches_its_serialized_form() {
596596
);
597597
}
598598
}
599+
600+
#[test]
601+
fn deterministic_lease_override_and_default_are_journaled() {
602+
for (lease, duration) in [(None, 30_000), (Some(300_000), 300_000), (Some(10), 10)] {
603+
let mut value =
604+
json!({"steps": [{"id": "cmd", "type": "deterministic", "command": "true"}]});
605+
if let Some(ms) = lease {
606+
value["steps"][0]["lease_ms"] = json!(ms);
607+
}
608+
let spec = crate::RunSpec::parse(&value).unwrap();
609+
spec.validate().unwrap();
610+
let state = RunState::fold("run", spec, &[]).unwrap();
611+
let clock = SimClock::new(1000);
612+
let Action::Append(started) = &next_actions(&state, clock.now_ms())[0] else {
613+
panic!("expected journaled lease");
614+
};
615+
let payload: AttemptStartedPayload =
616+
serde_json::from_value(started.payload.clone()).unwrap();
617+
assert_eq!(payload.lease_deadline_ms, clock.now_ms() + duration);
618+
}
619+
}
620+
621+
#[test]
622+
fn deterministic_lease_rejects_invalid_and_foreign_fields() {
623+
for ms in [0, u64::MAX] {
624+
let spec = crate::RunSpec::parse(&json!({"steps": [{
625+
"id": "cmd", "type": "deterministic", "command": "true", "lease_ms": ms
626+
}]}))
627+
.unwrap();
628+
assert!(spec.validate().is_err());
629+
}
630+
for kind in ["llm", "agent"] {
631+
assert!(
632+
crate::RunSpec::parse(&json!({"steps": [{
633+
"id": "worker", "type": kind, "lease_ms": 1000
634+
}]}))
635+
.is_err()
636+
);
637+
}
638+
}

‎kernel/relayflowd-core/src/spec.rs‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,16 @@ impl RunSpec {
146146
if step.max_iterations == 0 {
147147
return Err(SpecError::ZeroIterations(step.id.clone()));
148148
}
149+
if let StepKind::Deterministic {
150+
lease_ms: Some(ms), ..
151+
} = &step.kind
152+
&& (*ms == 0 || *ms > i64::MAX as u64)
153+
{
154+
return Err(SpecError::Malformed(format!(
155+
"step {}: lease_ms must be positive and fit in i64",
156+
step.id
157+
)));
158+
}
149159
let cli = match &step.kind {
150160
StepKind::Llm { cli, .. } | StepKind::Agent { cli, .. } => cli,
151161
StepKind::Deterministic { .. } => &None,
@@ -297,7 +307,7 @@ const STEP_COMMON_FIELDS: &[&str] = &[
297307
"memory",
298308
"requirements",
299309
];
300-
const STEP_DETERMINISTIC_FIELDS: &[&str] = &["command", "timeout_ms"];
310+
const STEP_DETERMINISTIC_FIELDS: &[&str] = &["command", "timeout_ms", "lease_ms"];
301311
const STEP_LLM_FIELDS: &[&str] = &["prompt", "model", "cli"];
302312
const STEP_AGENT_FIELDS: &[&str] = &[
303313
"instruction",
@@ -393,6 +403,8 @@ pub enum StepKind {
393403
command: CommandSpec,
394404
#[serde(default, skip_serializing_if = "Option::is_none")]
395405
timeout_ms: Option<u64>,
406+
#[serde(default, skip_serializing_if = "Option::is_none")]
407+
lease_ms: Option<u64>,
396408
},
397409
Llm {
398410
prompt: String,

‎kernel/relayflowd/src/exec_det.rs‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ pub(crate) fn execute_placed_with_input(
4242
let StepKind::Deterministic {
4343
command,
4444
timeout_ms,
45+
lease_ms,
4546
} = &step.kind
4647
else {
4748
return worker_error("deterministic executor received a non-deterministic step");
@@ -93,7 +94,11 @@ pub(crate) fn execute_placed_with_input(
9394
let stderr = child.stderr.take().expect("piped stderr");
9495
let stdout_reader = thread::spawn(move || read_all(stdout));
9596
let stderr_reader = thread::spawn(move || read_all(stderr));
96-
let timeout = Duration::from_millis(timeout_ms.unwrap_or(30_000));
97+
let timeout = Duration::from_millis(match (lease_ms, timeout_ms) {
98+
(Some(lease), Some(command)) => (*lease).min(*command),
99+
(Some(lease), None) => *lease,
100+
(None, command) => command.unwrap_or(30_000),
101+
});
97102
let (status, timed_out) = match child.wait_timeout(timeout) {
98103
Ok(Some(status)) => (Some(status), false),
99104
Ok(None) => {
@@ -196,6 +201,23 @@ mod tests {
196201
assert_eq!(result.failure_reason, Some(CompletionReason::Timeout));
197202
}
198203

204+
#[test]
205+
fn lease_override_bounds_execution_and_preserves_command_timeout() {
206+
for (lease, command_timeout) in [(5, None), (1000, Some(5))] {
207+
let mut value = json!({
208+
"id": "slow", "type": "deterministic", "command": "sleep 1", "lease_ms": lease
209+
});
210+
if let Some(ms) = command_timeout {
211+
value["timeout_ms"] = json!(ms);
212+
}
213+
let step: StepSpec = serde_json::from_value(value).unwrap();
214+
let started = std::time::Instant::now();
215+
let result = execute(&step);
216+
assert_eq!(result.failure_reason, Some(CompletionReason::Timeout));
217+
assert!(started.elapsed() < Duration::from_millis(900));
218+
}
219+
}
220+
199221
#[test]
200222
fn timeout_kills_the_whole_process_group() {
201223
// The backgrounded sleep inherits the stdout/stderr pipes. If a

‎packages/sdk/src/authored-flow-error.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ export type AuthoredFlowExecutionErrorCode =
1919
| 'operation_after_completion'
2020
| 'operation_callback_failed'
2121
| 'step_failed'
22+
| 'lease_exceeded'
2223
| 'unsupported_completion'
2324
| 'unsupported_gate'
2425
| 'unsupported_header'

‎packages/sdk/src/authored-flow-executor.ts‎

Lines changed: 7 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,8 @@ import { checkMcpHeader, McpPreflightError } from './cli/check-typescript.js';
1010
import { buildMcpProxy, runMcpEffect } from './authored-mcp.js';
1111
import { AuthoredBudget } from './authored-budget.js';
1212
import { assertMemoryReachable, authoredMemory, scriptMemoryScope } from './authored-memory.js';
13-
import { authoredWorkerRunner } from './authored-worker-step.js';
14-
import { readSuccessfulOutput, isSurfaceRunCompletionReason } from './authored-step-output.js';
13+
import { authoredDeterministicRunner, authoredWorkerRunner } from './authored-worker-step.js';
14+
import { isSurfaceRunCompletionReason } from './authored-step-output.js';
1515
import {
1616
type LlmOptions,
1717
type CloudHelper,
@@ -23,7 +23,7 @@ import {
2323
import type { FlowHandle } from '@relayflows/surface/runtime';
2424
import { join } from 'node:path';
2525
import { observeStep, type ProgressEvent } from './progress.js';
26-
import { compileSpec, toKernelSpec } from './compile.js';
26+
import { parseStepTimeout } from './compile.js';
2727
import { getAuthoredFlowDefinition } from './authored-flow.js';
2828
import type { GetFlowDefinition } from './authored-flow-loader.js';
2929
import type { RunLifecycleOptions } from './cli/run.js';
@@ -42,7 +42,6 @@ import type {
4242
CompletionReason as ProtocolCompletionReason,
4343
RunCompletionReason as ProtocolRunCompletionReason,
4444
} from './protocol.js';
45-
import { SPEC_SCHEMA_VERSION } from './spec.js';
4645

4746
type Assert<T extends true> = T;
4847
type Equal<A, B> = [A] extends [B]
@@ -169,19 +168,7 @@ export async function executeAuthoredFlow<Input = undefined>(
169168
let nextStep = 1;
170169
let requestedCompletion: SurfaceRunCompletionReason | undefined;
171170

172-
const lowerDeterministic = async (
173-
id: string,
174-
command: string,
175-
terminal = false,
176-
): Promise<string> => {
177-
const spec = toKernelSpec(compileSpec({
178-
version: SPEC_SCHEMA_VERSION,
179-
name: `${definition.name}/${id}`,
180-
steps: [{ id, type: 'deterministic', command }],
181-
}));
182-
if (terminal) return readSuccessfulOutput(journal, await journal.runStart(spec), id, journalSteps);
183-
return budget.execute(journal, spec, outcome => readSuccessfulOutput(journal, outcome, id, journalSteps));
184-
};
171+
const lowerDeterministic = authoredDeterministicRunner(definition.name, journal, journalSteps, budget);
185172

186173
const worker = authoredWorkerRunner(definition, journal, flowPath, journalSteps, waitOptions, localAgentStream, budget, definition.header.budget);
187174

@@ -248,14 +235,15 @@ export async function executeAuthoredFlow<Input = undefined>(
248235
() => assertOperationAllowed('memory', definition.name, requestedCompletion),
249236
definition.header.memory?.script !== false,
250237
),
251-
run(command) {
238+
run(command, runOptions) {
252239
assertOperationAllowed('run', definition.name, requestedCompletion);
240+
const leaseMs = runOptions?.timeout === undefined ? undefined : parseStepTimeout(runOptions.timeout);
253241
const id = `run-${nextStep++}`;
254242
return trackStep(authoredSteps, new AuthoredFlowOperation(
255243
id,
256244
'run',
257245
() => assertOperationAllowed('run', definition.name, requestedCompletion),
258-
() => observeStep(id, 'deterministic', () => lowerDeterministic(id, command), options.onProgress),
246+
() => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs), options.onProgress),
259247
lifecycle,
260248
));
261249
},

‎packages/sdk/src/authored-step-output.ts‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,14 @@ export async function readSuccessfulOutput(
5050
stepId: string,
5151
journalSteps: AuthoredFlowJournalStep[],
5252
): Promise<string> {
53-
const output = await readCompletedStepOutput(journal, outcome.run_id, stepId, journalSteps);
53+
const output = await readCompletedStepOutput(journal, outcome.run_id, stepId, journalSteps).catch(error => {
54+
if (error instanceof AuthoredFlowExecutionError
55+
&& (error.completionReason === 'timeout' || error.completionReason === 'lease_expired')) {
56+
throw new AuthoredFlowExecutionError('lease_exceeded',
57+
`f.run step "${stepId}" exceeded its command timeout.`, error.completionReason, error.runId);
58+
}
59+
throw error;
60+
});
5461
if (!isRecord(output) || typeof output['stdout_tail'] !== 'string') {
5562
throw protocolViolation(outcome.run_id, `step "${stepId}" has no string stdout_tail`);
5663
}

‎packages/sdk/src/authored-worker-step.ts‎

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,14 @@
11
import type { AuthoredBudget } from './authored-budget.js';
22
import { parseBudget } from './budget.js';
33
import type { AgentOptions, AgentResult, LlmOptions } from '@relayflows/surface';
4-
import { toKernelSpec } from './compile.js';
4+
import { compileSpec, toKernelSpec } from './compile.js';
55
import { checkAuthoredFlow } from './cli/check.js';
66
import { classifyOutcome, type RunLifecycleOptions } from './cli/run.js';
77
import type { PreflightDiagnostic } from './preflight.js';
88
import { AuthoredFlowExecutionError } from './authored-flow-error.js';
99
import type { JournalClient } from './journal-client.js';
1010
import { SPEC_SCHEMA_VERSION, type FlowSpec, type StepSpec } from './spec.js';
11-
import { isSurfaceCompletionReason, readCompletedStepOutput } from './authored-step-output.js';
11+
import { isSurfaceCompletionReason, readCompletedStepOutput, readSuccessfulOutput } from './authored-step-output.js';
1212
import type { AuthoredFlowJournalStep } from './authored-flow-executor.js';
1313
import { snapshotJsonValue } from './json-value.js';
1414

@@ -142,3 +142,18 @@ export function authoredWorkerRunner(
142142
},
143143
};
144144
}
145+
146+
/** Deterministic commands execute inline under their per-invocation lease. */
147+
export function authoredDeterministicRunner(
148+
name: string, journal: JournalClient, journalSteps: AuthoredFlowJournalStep[], budget: AuthoredBudget,
149+
) {
150+
return async (id: string, command: string, terminal = false, leaseMs?: number): Promise<string> => {
151+
const spec = toKernelSpec(compileSpec({
152+
version: SPEC_SCHEMA_VERSION,
153+
name: `${name}/${id}`,
154+
steps: [{ id, type: 'deterministic', command, ...(leaseMs === undefined ? {} : { lease_ms: leaseMs }) }],
155+
}));
156+
if (terminal) return readSuccessfulOutput(journal, await journal.runStart(spec), id, journalSteps);
157+
return budget.execute(journal, spec, outcome => readSuccessfulOutput(journal, outcome, id, journalSteps));
158+
};
159+
}

‎packages/sdk/src/compile.ts‎

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,21 @@ export class CompileError extends Error {
5353
}
5454
}
5555

56+
/** Parse the f.run timeout before submitting any command to the kernel. */
57+
export function parseStepTimeout(timeout: unknown): number {
58+
const units: Record<string, number> = { ms: 1, s: 1000, m: 60_000 };
59+
const match = typeof timeout === 'string' ? /^(\d+(?:\.\d+)?)(ms|s|m)$/.exec(timeout) : null;
60+
const milliseconds = typeof timeout === 'number' ? timeout
61+
: match === null ? NaN : Number(match[1]) * units[match[2]!]!;
62+
if (!Number.isSafeInteger(milliseconds) || milliseconds <= 0) {
63+
throw new CompileError(['f.run timeout must be a positive whole number of milliseconds or a duration such as "10s" or "5m".'], 'timeout_invalid');
64+
}
65+
if (milliseconds > 15 * 60_000) {
66+
throw new CompileError(['f.run timeout exceeds the maximum of 15 minutes declared in SURFACE.md.'], 'lease_exceeded');
67+
}
68+
return milliseconds;
69+
}
70+
5671
/**
5772
* Compile a YAML string into a validated authoring `FlowSpec`.
5873
* Throws `CompileError` on a YAML parse error or any validation failure.
@@ -156,6 +171,7 @@ function compileStep(step: StepSpec): StepSpec {
156171
// #138: `timeoutMs` is deterministic-only — worker-backed verbs own
157172
// their dispatch timeout. It must be spread HERE and nowhere in `base`.
158173
...(s.timeoutMs !== undefined ? { timeoutMs: s.timeoutMs } : {}),
174+
...(s.lease_ms !== undefined ? { lease_ms: parseStepTimeout(s.lease_ms) } : {}),
159175
};
160176
}
161177
case 'llm': {
@@ -372,14 +388,14 @@ function kernelTriggerToAuthoring(value: unknown, at: string): unknown {
372388
function kernelStepToAuthoring(value: unknown, at: string): unknown {
373389
const unionKeys = [
374390
'id', 'type', 'depends_on', 'max_iterations', 'retry', 'verification', 'memory', 'requirements', 'input',
375-
'command', 'timeout_ms', 'prompt', 'model', 'cli', 'instruction',
391+
'command', 'timeout_ms', 'lease_ms', 'prompt', 'model', 'cli', 'instruction',
376392
'recovery_mode', 'surfaces', 'permissions',
377393
] as const;
378394
const step = requireKernelObject(value, unionKeys, at);
379395
const type = step['type'];
380396
const commonKeys = ['id', 'type', 'depends_on', 'max_iterations', 'retry', 'verification', 'memory', 'requirements', 'input'] as const;
381397
const typeKeys = type === 'deterministic'
382-
? ['command', 'timeout_ms'] as const
398+
? ['command', 'timeout_ms', 'lease_ms'] as const
383399
: type === 'llm'
384400
? ['prompt', 'model', 'cli'] as const
385401
: type === 'agent'
@@ -405,6 +421,7 @@ function kernelStepToAuthoring(value: unknown, at: string): unknown {
405421
...common,
406422
command: step['command'],
407423
...(step['timeout_ms'] !== undefined ? { timeoutMs: step['timeout_ms'] } : {}),
424+
...(step['lease_ms'] !== undefined ? { lease_ms: step['lease_ms'] } : {}),
408425
};
409426
}
410427
if (type === 'llm') {
@@ -554,6 +571,7 @@ function toKernelStep(step: StepSpec): KernelStepSpec {
554571
type: 'deterministic',
555572
command: step.command,
556573
...(step.timeoutMs !== undefined ? { timeout_ms: step.timeoutMs } : {}),
574+
...(step.lease_ms !== undefined ? { lease_ms: parseStepTimeout(step.lease_ms) } : {}),
557575
};
558576
case 'llm': {
559577
return {

0 commit comments

Comments
 (0)