From 2fa82dd183dd7e4e23e6f1f1f42f7c341abd864b Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 24 Sep 2026 07:19:41 -0700 Subject: [PATCH 1/2] fix(kernel): stamp each start when journaled, not when its batch was elected An independent batch is elected at one instant and driven serially. A start journaled after a deterministic peer ran kept the election time, so: - its wall clock included the peer's runtime (a sleep 14 step recorded 24055ms after a 10s peer), and - its lease deadline, derived from the same stale time, was already partly spent: with the 30s default, a step started 10s late held a 20s lease, and could be judged expired while running. The dispatched lease of an llm/agent step in the same batch was shortened the same way. Stamp each start when it is appended and move its lease deadline, and its dispatch's, by the same delay. Serial execution itself is unchanged: the budget gate and stop_after tests pin it deliberately. Under a simulated clock the delay is zero, so replay is unaffected. Co-Authored-By: Claude Opus 5.5 (1M context) --- kernel/relayflowd/src/engine/drive.rs | 31 ++++++++- kernel/relayflowd/tests/parallel_driver.rs | 77 ++++++++++++++++++++++ 2 files changed, 107 insertions(+), 1 deletion(-) diff --git a/kernel/relayflowd/src/engine/drive.rs b/kernel/relayflowd/src/engine/drive.rs index 2005b9e05..0d3b16cc4 100644 --- a/kernel/relayflowd/src/engine/drive.rs +++ b/kernel/relayflowd/src/engine/drive.rs @@ -1,4 +1,8 @@ -use std::{collections::BTreeSet, thread, time::Duration}; +use std::{ + collections::{BTreeMap, BTreeSet}, + thread, + time::Duration, +}; use anyhow::{Context, Result, bail}; use relayflowd_core::{ @@ -50,6 +54,12 @@ impl Engine { let mut backpressured = false; let mut handoff_failed = false; let mut skipped_dispatches = BTreeSet::new(); + // The batch was elected at one instant but runs serially, so a + // start journaled after a deterministic peer ran is late by that + // peer's runtime. Stamp each start when it is journaled and move + // its lease deadline — and its dispatch's — by the same delay, so + // neither the attempt's time nor its lease is spent in advance. + let mut start_delays = BTreeMap::new(); for action in actions { match action { Action::Append(mut entry) => { @@ -102,6 +112,22 @@ impl Engine { continue; } } + if entry.entry_type == relayflowd_core::EntryType::StepAttemptStarted { + let now_ms = self.clock.now_ms(); + let delay = now_ms.saturating_sub(entry.at_ms); + if delay > 0 { + entry.at_ms = now_ms; + if let Some(deadline) = entry.payload["lease_deadline_ms"].as_i64() { + entry.payload["lease_deadline_ms"] = + deadline.saturating_add(delay).into(); + } + if let (Some(step_id), Some(attempt)) = + (entry.step_id.clone(), entry.attempt) + { + start_delays.insert((step_id, attempt), delay); + } + } + } let prepared = (|| -> Result<()> { self.prepare_start_entry(&state, &mut entry)?; self.assign_executor(&state, &mut entry)?; @@ -194,6 +220,9 @@ impl Engine { if skipped_dispatches.remove(&(step.id.clone(), attempt)) { continue; } + let lease_deadline_ms = lease_deadline_ms.saturating_add( + start_delays.remove(&(step.id.clone(), attempt)).unwrap_or(0), + ); let input = self.resolve_step_input(&mut journal, &step, attempt); if !matches!(input, Ok(Some(_))) { if let Some(dispatcher) = &self.dispatcher { diff --git a/kernel/relayflowd/tests/parallel_driver.rs b/kernel/relayflowd/tests/parallel_driver.rs index af9f82ba3..eb0a2ae34 100644 --- a/kernel/relayflowd/tests/parallel_driver.rs +++ b/kernel/relayflowd/tests/parallel_driver.rs @@ -478,3 +478,80 @@ fn pause_before_second_independent_step_holds_the_driver_boundary() { ); assert_eq!(fs::read_to_string(marker).unwrap(), "first\nsecond\n"); } + +fn entry_of(entries: &[JournalEntry], entry_type: EntryType, step_id: &str) -> JournalEntry { + entries + .iter() + .find(|entry| entry.entry_type == entry_type && entry.step_id.as_deref() == Some(step_id)) + .cloned() + .unwrap_or_else(|| panic!("no {entry_type:?} for {step_id}")) +} + +/// An independent batch runs serially, so a start journaled after a slow +/// deterministic peer is stamped when it is journaled, not when the batch was +/// elected: its wall clock is its own, and its full lease is still ahead of it. +#[test] +fn a_start_after_a_slow_deterministic_peer_is_stamped_when_journaled() { + let spec = RunSpec::parse(&json!({ + "name": "late-start", + "steps": [ + {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, + {"id": "quick", "type": "deterministic", "command": "true"} + ] + })) + .unwrap(); + let directory = tempdir().unwrap(); + let engine = Engine::new(directory.path()); + let outcome = engine.start(spec, "test", None).unwrap(); + assert_eq!(outcome.status, RunStatus::Completed); + let entries = engine.journal_entries(&outcome.run_id, 1, usize::MAX).unwrap(); + + let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); + let quick_start = entry_of(&entries, EntryType::StepAttemptStarted, "quick"); + let quick_done = entry_of(&entries, EntryType::StepCompleted, "quick"); + assert!( + quick_start.at_ms >= slow_done.at_ms, + "quick started at {} before slow completed at {}", + quick_start.at_ms, + slow_done.at_ms + ); + let wallclock = quick_done.payload["spend"]["wallclock_ms"].as_i64().unwrap(); + assert!(wallclock < 300, "quick's wall clock includes slow's runtime: {wallclock}ms"); + for step in ["slow", "quick"] { + let start = entry_of(&entries, EntryType::StepAttemptStarted, step); + assert_eq!( + start.payload["lease_deadline_ms"].as_i64().unwrap() - start.at_ms, + 30_000, + "{step}'s lease must run from its own start" + ); + } +} + +/// The same delay moves the dispatched lease: a worker handed a step after a +/// slow deterministic peer holds the deadline its start entry records. +#[test] +fn a_dispatch_after_a_slow_deterministic_peer_keeps_its_full_lease() { + let spec = RunSpec::parse(&json!({ + "name": "late-dispatch", + "steps": [ + {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, + {"id": "model", "type": "llm", "prompt": "p", "model": "stub"} + ] + })) + .unwrap(); + let directory = tempdir().unwrap(); + let dispatcher = Arc::new(RecordingDispatcher::llm()); + let observer = Arc::new(StartCrashObserver::default()); + let engine = Engine::with_runtime(directory.path(), dispatcher.clone(), observer.clone()); + engine.start(spec, "test", None).unwrap(); + let entries = engine + .journal_entries(&observer.run_id(), 1, usize::MAX) + .unwrap(); + + let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); + let start = entry_of(&entries, EntryType::StepAttemptStarted, "model"); + let dispatch = dispatcher.calls().into_iter().find(|call| call.step_id == "model").unwrap(); + assert!(start.at_ms >= slow_done.at_ms); + assert_eq!(dispatch.lease_deadline_ms, start.payload["lease_deadline_ms"].as_i64().unwrap()); + assert_eq!(dispatch.lease_deadline_ms - start.at_ms, 30_000); +} From 04e96886b0b41bbbd1230a9fd56d280f69392351 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 24 Sep 2026 07:30:03 -0700 Subject: [PATCH 2/2] fix(kernel): stamp the start after placement; move late-start tests out - Stamp the start after route_start, just before the append, so a slow placement cannot hand the step a lease already partly spent. - Move the late-start tests to their own file: parallel_driver.rs had grown past the 500-line limit. - Replace the real-time wall-clock bound with an ordering and a structural check, so a loaded CI host cannot fail it spuriously. Co-Authored-By: Claude Opus 5.5 (1M context) --- kernel/relayflowd/src/engine/drive.rs | 48 +++++--- kernel/relayflowd/tests/late_start.rs | 130 +++++++++++++++++++++ kernel/relayflowd/tests/parallel_driver.rs | 77 ------------ 3 files changed, 161 insertions(+), 94 deletions(-) create mode 100644 kernel/relayflowd/tests/late_start.rs diff --git a/kernel/relayflowd/src/engine/drive.rs b/kernel/relayflowd/src/engine/drive.rs index 0d3b16cc4..7eb27f0eb 100644 --- a/kernel/relayflowd/src/engine/drive.rs +++ b/kernel/relayflowd/src/engine/drive.rs @@ -112,26 +112,12 @@ impl Engine { continue; } } - if entry.entry_type == relayflowd_core::EntryType::StepAttemptStarted { - let now_ms = self.clock.now_ms(); - let delay = now_ms.saturating_sub(entry.at_ms); - if delay > 0 { - entry.at_ms = now_ms; - if let Some(deadline) = entry.payload["lease_deadline_ms"].as_i64() { - entry.payload["lease_deadline_ms"] = - deadline.saturating_add(delay).into(); - } - if let (Some(step_id), Some(attempt)) = - (entry.step_id.clone(), entry.attempt) - { - start_delays.insert((step_id, attempt), delay); - } - } - } let prepared = (|| -> Result<()> { self.prepare_start_entry(&state, &mut entry)?; self.assign_executor(&state, &mut entry)?; self.route_start(&mut journal, &state, &mut entry)?; + // After placement, which can itself be slow. + self.stamp_start(&mut entry, &mut start_delays); self.append(&mut journal, &entry)?; Ok(()) })(); @@ -221,7 +207,9 @@ impl Engine { continue; } let lease_deadline_ms = lease_deadline_ms.saturating_add( - start_delays.remove(&(step.id.clone(), attempt)).unwrap_or(0), + start_delays + .remove(&(step.id.clone(), attempt)) + .unwrap_or(0), ); let input = self.resolve_step_input(&mut journal, &step, attempt); if !matches!(input, Ok(Some(_))) { @@ -391,6 +379,32 @@ impl Engine { } } + /// Stamp a start with the moment it is journaled and move its lease + /// deadline by the delay since the batch was elected, recording the delay + /// so the batch's dispatch for the same attempt moves with it. Anything + /// but an attempt start is left alone. + fn stamp_start( + &self, + entry: &mut relayflowd_core::JournalEntry, + start_delays: &mut BTreeMap<(String, u32), i64>, + ) { + if entry.entry_type != relayflowd_core::EntryType::StepAttemptStarted { + return; + } + let now_ms = self.clock.now_ms(); + let delay = now_ms.saturating_sub(entry.at_ms); + if delay <= 0 { + return; + } + entry.at_ms = now_ms; + if let Some(deadline) = entry.payload["lease_deadline_ms"].as_i64() { + entry.payload["lease_deadline_ms"] = deadline.saturating_add(delay).into(); + } + if let (Some(step_id), Some(attempt)) = (entry.step_id.clone(), entry.attempt) { + start_delays.insert((step_id, attempt), delay); + } + } + /// Journal a dispatch that could not be honored as a declared `worker_error` /// completion of the attempt. It runs through the same completion path as /// any other rejected attempt, so retry policy and the failure taxonomy diff --git a/kernel/relayflowd/tests/late_start.rs b/kernel/relayflowd/tests/late_start.rs new file mode 100644 index 000000000..b0b8a71e2 --- /dev/null +++ b/kernel/relayflowd/tests/late_start.rs @@ -0,0 +1,130 @@ +//! A batch is elected at one instant and driven serially, so a start journaled +//! after a slow deterministic peer must be stamped when it is journaled: its +//! wall clock is its own, and its full lease — journaled and dispatched — is +//! still ahead of it. + +use std::sync::{Arc, Mutex}; + +use relayflowd::worker::{DispatchOutcome, JournalObserver, StepDispatch, StepDispatcher}; +use relayflowd::{Engine, RunStatus}; +use relayflowd_core::{EntryType, JournalEntry, RunSpec, StepType}; +use serde_json::json; +use tempfile::tempdir; + +/// The kernel's default lease for a step that declares none. +const DEFAULT_LEASE_MS: i64 = 30_000; + +#[derive(Default)] +struct Recorder { + run_id: Mutex>, + dispatches: Mutex>, +} + +impl JournalObserver for Recorder { + fn appended(&self, entry: &JournalEntry) { + if entry.entry_type == EntryType::RunSpawned { + *self.run_id.lock().unwrap() = Some(entry.run_id.clone()); + } + } +} + +impl StepDispatcher for Recorder { + fn executor(&self, _step_type: StepType) -> Option { + Some("late-start-test".to_owned()) + } + + fn available(&self, step_type: StepType) -> bool { + step_type == StepType::Llm + } + + fn dispatch(&self, dispatch: StepDispatch) -> anyhow::Result { + self.dispatches.lock().unwrap().push(dispatch); + Ok(DispatchOutcome::Dispatched) + } +} + +fn entry_of(entries: &[JournalEntry], entry_type: EntryType, step_id: &str) -> JournalEntry { + entries + .iter() + .find(|entry| entry.entry_type == entry_type && entry.step_id.as_deref() == Some(step_id)) + .cloned() + .unwrap_or_else(|| panic!("no {entry_type:?} for {step_id}")) +} + +#[test] +fn a_start_after_a_slow_deterministic_peer_is_stamped_when_journaled() { + let spec = RunSpec::parse(&json!({ + "name": "late-start", + "steps": [ + {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, + {"id": "quick", "type": "deterministic", "command": "true"} + ] + })) + .unwrap(); + let directory = tempdir().unwrap(); + let engine = Engine::new(directory.path()); + let outcome = engine.start(spec, "test", None).unwrap(); + assert_eq!(outcome.status, RunStatus::Completed); + let entries = engine + .journal_entries(&outcome.run_id, 1, usize::MAX) + .unwrap(); + + let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); + let quick_start = entry_of(&entries, EntryType::StepAttemptStarted, "quick"); + let quick_done = entry_of(&entries, EntryType::StepCompleted, "quick"); + // Ordering, not a wall-clock bound: the start is journaled after the peer + // finished, so the step's own span cannot contain the peer's runtime. + assert!( + quick_start.at_ms >= slow_done.at_ms, + "quick started at {} before slow completed at {}", + quick_start.at_ms, + slow_done.at_ms + ); + assert_eq!( + quick_done.payload["spend"]["wallclock_ms"].as_i64().unwrap(), + quick_done.at_ms - quick_start.at_ms, + ); + for step in ["slow", "quick"] { + let start = entry_of(&entries, EntryType::StepAttemptStarted, step); + assert_eq!( + start.payload["lease_deadline_ms"].as_i64().unwrap() - start.at_ms, + DEFAULT_LEASE_MS, + "{step}'s lease must run from its own start" + ); + } +} + +#[test] +fn a_dispatch_after_a_slow_deterministic_peer_keeps_its_full_lease() { + let spec = RunSpec::parse(&json!({ + "name": "late-dispatch", + "steps": [ + {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, + {"id": "model", "type": "llm", "prompt": "p", "model": "stub"} + ] + })) + .unwrap(); + let directory = tempdir().unwrap(); + let recorder = Arc::new(Recorder::default()); + let engine = Engine::with_runtime(directory.path(), recorder.clone(), recorder.clone()); + engine.start(spec, "test", None).unwrap(); + let run_id = recorder.run_id.lock().unwrap().clone().unwrap(); + let entries = engine.journal_entries(&run_id, 1, usize::MAX).unwrap(); + + let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); + let start = entry_of(&entries, EntryType::StepAttemptStarted, "model"); + let dispatch = recorder + .dispatches + .lock() + .unwrap() + .iter() + .find(|dispatch| dispatch.step_id == "model") + .cloned() + .unwrap(); + assert!(start.at_ms >= slow_done.at_ms); + assert_eq!( + dispatch.lease_deadline_ms, + start.payload["lease_deadline_ms"].as_i64().unwrap() + ); + assert_eq!(dispatch.lease_deadline_ms - start.at_ms, DEFAULT_LEASE_MS); +} diff --git a/kernel/relayflowd/tests/parallel_driver.rs b/kernel/relayflowd/tests/parallel_driver.rs index eb0a2ae34..af9f82ba3 100644 --- a/kernel/relayflowd/tests/parallel_driver.rs +++ b/kernel/relayflowd/tests/parallel_driver.rs @@ -478,80 +478,3 @@ fn pause_before_second_independent_step_holds_the_driver_boundary() { ); assert_eq!(fs::read_to_string(marker).unwrap(), "first\nsecond\n"); } - -fn entry_of(entries: &[JournalEntry], entry_type: EntryType, step_id: &str) -> JournalEntry { - entries - .iter() - .find(|entry| entry.entry_type == entry_type && entry.step_id.as_deref() == Some(step_id)) - .cloned() - .unwrap_or_else(|| panic!("no {entry_type:?} for {step_id}")) -} - -/// An independent batch runs serially, so a start journaled after a slow -/// deterministic peer is stamped when it is journaled, not when the batch was -/// elected: its wall clock is its own, and its full lease is still ahead of it. -#[test] -fn a_start_after_a_slow_deterministic_peer_is_stamped_when_journaled() { - let spec = RunSpec::parse(&json!({ - "name": "late-start", - "steps": [ - {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, - {"id": "quick", "type": "deterministic", "command": "true"} - ] - })) - .unwrap(); - let directory = tempdir().unwrap(); - let engine = Engine::new(directory.path()); - let outcome = engine.start(spec, "test", None).unwrap(); - assert_eq!(outcome.status, RunStatus::Completed); - let entries = engine.journal_entries(&outcome.run_id, 1, usize::MAX).unwrap(); - - let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); - let quick_start = entry_of(&entries, EntryType::StepAttemptStarted, "quick"); - let quick_done = entry_of(&entries, EntryType::StepCompleted, "quick"); - assert!( - quick_start.at_ms >= slow_done.at_ms, - "quick started at {} before slow completed at {}", - quick_start.at_ms, - slow_done.at_ms - ); - let wallclock = quick_done.payload["spend"]["wallclock_ms"].as_i64().unwrap(); - assert!(wallclock < 300, "quick's wall clock includes slow's runtime: {wallclock}ms"); - for step in ["slow", "quick"] { - let start = entry_of(&entries, EntryType::StepAttemptStarted, step); - assert_eq!( - start.payload["lease_deadline_ms"].as_i64().unwrap() - start.at_ms, - 30_000, - "{step}'s lease must run from its own start" - ); - } -} - -/// The same delay moves the dispatched lease: a worker handed a step after a -/// slow deterministic peer holds the deadline its start entry records. -#[test] -fn a_dispatch_after_a_slow_deterministic_peer_keeps_its_full_lease() { - let spec = RunSpec::parse(&json!({ - "name": "late-dispatch", - "steps": [ - {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, - {"id": "model", "type": "llm", "prompt": "p", "model": "stub"} - ] - })) - .unwrap(); - let directory = tempdir().unwrap(); - let dispatcher = Arc::new(RecordingDispatcher::llm()); - let observer = Arc::new(StartCrashObserver::default()); - let engine = Engine::with_runtime(directory.path(), dispatcher.clone(), observer.clone()); - engine.start(spec, "test", None).unwrap(); - let entries = engine - .journal_entries(&observer.run_id(), 1, usize::MAX) - .unwrap(); - - let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); - let start = entry_of(&entries, EntryType::StepAttemptStarted, "model"); - let dispatch = dispatcher.calls().into_iter().find(|call| call.step_id == "model").unwrap(); - assert!(start.at_ms >= slow_done.at_ms); - assert_eq!(dispatch.lease_deadline_ms, start.payload["lease_deadline_ms"].as_i64().unwrap()); - assert_eq!(dispatch.lease_deadline_ms - start.at_ms, 30_000); -}