use salvor_core::{Effect, Event, EventEnvelope, ForkOrigin, RunId, SequenceNumber};
use serde_json::Value;
use thiserror::Error;
use time::OffsetDateTime;
#[derive(Clone, Debug)]
pub struct WriteHazard {
pub seq: u64,
pub tool: String,
pub input: Value,
pub idempotency_key: Option<String>,
pub recorded_at: OffsetDateTime,
}
#[derive(Debug, Error)]
pub enum ForkError {
#[error("run is not a graph run: its log does not open with GraphRunStarted")]
NotAGraphRun,
#[error(
"the origin never entered node `{node}`; a fork point must be a node boundary the run reached"
)]
NodeNeverEntered {
node: String,
},
}
#[derive(Debug)]
pub struct ForkPlan {
origin_run: RunId,
graph_hash: String,
from_node: String,
boundary: SequenceNumber,
prefix: Vec<EventEnvelope>,
hazards: Vec<WriteHazard>,
}
impl ForkPlan {
#[must_use]
pub fn origin_run(&self) -> RunId {
self.origin_run
}
#[must_use]
pub fn graph_hash(&self) -> &str {
&self.graph_hash
}
#[must_use]
pub fn from_node(&self) -> &str {
&self.from_node
}
#[must_use]
pub fn through_seq(&self) -> SequenceNumber {
SequenceNumber::new(self.boundary.get() - 1)
}
#[must_use]
pub fn prefix_len(&self) -> usize {
self.prefix.len()
}
#[must_use]
pub fn hazards(&self) -> &[WriteHazard] {
&self.hazards
}
#[must_use]
pub fn hazard_seqs(&self) -> Vec<u64> {
let mut seqs: Vec<u64> = self.hazards.iter().map(|h| h.seq).collect();
seqs.sort_unstable();
seqs
}
#[must_use]
pub fn build_child_prefix(
&self,
child: RunId,
acknowledged_writes: Vec<u64>,
) -> Vec<EventEnvelope> {
let origin = ForkOrigin {
run_id: self.origin_run,
through_seq: self.through_seq(),
from_node: self.from_node.clone(),
graph_hash: self.graph_hash.clone(),
acknowledged_writes,
};
self.prefix
.iter()
.enumerate()
.map(|(index, envelope)| {
let mut event = envelope.event.clone();
if index == 0
&& let Event::GraphRunStarted { forked_from, .. } = &mut event
{
*forked_from = Some(origin.clone());
}
EventEnvelope {
run_id: child,
seq: envelope.seq,
schema_version: envelope.schema_version,
recorded_at: envelope.recorded_at,
event,
}
})
.collect()
}
}
pub fn plan_fork(origin_log: &[EventEnvelope], from_node: &str) -> Result<ForkPlan, ForkError> {
let graph_hash = match origin_log.first().map(|envelope| &envelope.event) {
Some(Event::GraphRunStarted { graph_hash, .. }) => graph_hash.clone(),
_ => return Err(ForkError::NotAGraphRun),
};
let boundary = origin_log
.iter()
.find_map(|envelope| match &envelope.event {
Event::NodeEntered { node } if node == from_node => Some(envelope.seq),
_ => None,
})
.ok_or_else(|| ForkError::NodeNeverEntered {
node: from_node.to_owned(),
})?;
let prefix: Vec<EventEnvelope> = origin_log
.iter()
.filter(|envelope| envelope.seq < boundary)
.cloned()
.collect();
let hazards: Vec<WriteHazard> = origin_log
.iter()
.filter(|envelope| envelope.seq >= boundary)
.filter_map(|envelope| match &envelope.event {
Event::ToolCallRequested {
tool,
input,
effect: Effect::Write,
idempotency_key,
..
} => Some(WriteHazard {
seq: envelope.seq.get(),
tool: tool.clone(),
input: input.clone(),
idempotency_key: idempotency_key.clone(),
recorded_at: envelope.recorded_at,
}),
_ => None,
})
.collect();
Ok(ForkPlan {
origin_run: origin_log[0].run_id,
graph_hash,
from_node: from_node.to_owned(),
boundary,
prefix,
hazards,
})
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use time::macros::datetime;
use uuid::Uuid;
fn run(tag: u8) -> RunId {
let mut bytes = [0u8; 16];
bytes[15] = tag;
bytes[6] = 0x40;
bytes[8] = 0x80;
RunId::from_uuid(Uuid::from_bytes(bytes))
}
fn envelope(run_id: RunId, seq: u64, event: Event) -> EventEnvelope {
EventEnvelope::new(
run_id,
SequenceNumber::new(seq),
datetime!(2026-07-14 12:00:00 UTC),
event,
)
}
fn origin_log() -> Vec<EventEnvelope> {
let r = run(1);
vec![
envelope(
r,
0,
Event::GraphRunStarted {
graph_hash: "sha256:graph".into(),
input: json!({"topic": "otters"}),
labels: None,
forked_from: None,
},
),
envelope(
r,
1,
Event::NodeEntered {
node: "research".into(),
},
),
envelope(
r,
2,
Event::NodeExited {
node: "research".into(),
},
),
envelope(
r,
3,
Event::NodeEntered {
node: "publish".into(),
},
),
envelope(
r,
4,
Event::ToolCallRequested {
seq: SequenceNumber::new(4),
tool: "http_post".into(),
input: json!({"body": "draft"}),
effect: Effect::Write,
idempotency_key: None,
},
),
envelope(
r,
5,
Event::ToolCallCompleted {
seq: SequenceNumber::new(4),
output: json!({"ok": true}),
},
),
envelope(
r,
6,
Event::NodeExited {
node: "publish".into(),
},
),
envelope(
r,
7,
Event::RunCompleted {
output: json!(null),
},
),
]
}
#[test]
fn a_non_graph_run_is_refused() {
let log = vec![envelope(
run(1),
0,
Event::RunStarted {
agent_def_hash: "sha256:a".into(),
input: json!(null),
labels: None,
},
)];
assert!(matches!(
plan_fork(&log, "anything"),
Err(ForkError::NotAGraphRun)
));
}
#[test]
fn a_node_the_run_never_entered_is_refused() {
let error = plan_fork(&origin_log(), "ghost").expect_err("ghost was never entered");
match error {
ForkError::NodeNeverEntered { node } => assert_eq!(node, "ghost"),
other => panic!("expected NodeNeverEntered, got {other:?}"),
}
}
#[test]
fn forking_before_the_write_finds_no_hazard() {
let plan = plan_fork(&origin_log(), "research").expect("research was entered");
assert_eq!(plan.through_seq().get(), 0, "prefix is just the head");
assert_eq!(plan.prefix_len(), 1);
assert_eq!(
plan.hazard_seqs(),
vec![4],
"the downstream write is still in the re-walked segment"
);
}
#[test]
fn forking_at_the_write_node_lists_the_write() {
let plan = plan_fork(&origin_log(), "publish").expect("publish was entered");
assert_eq!(
plan.through_seq().get(),
2,
"prefix is through research-exited"
);
assert_eq!(plan.prefix_len(), 3, "head + research entered + exited");
let hazards = plan.hazards();
assert_eq!(hazards.len(), 1);
assert_eq!(hazards[0].seq, 4);
assert_eq!(hazards[0].tool, "http_post");
assert_eq!(hazards[0].input, json!({"body": "draft"}));
}
#[test]
fn the_child_prefix_is_byte_identical_modulo_run_id_and_forked_from() {
let origin = origin_log();
let plan = plan_fork(&origin, "publish").expect("publish was entered");
let child = run(2);
let child_prefix = plan.build_child_prefix(child, plan.hazard_seqs());
for (index, env) in child_prefix.iter().enumerate() {
assert_eq!(env.run_id, child);
assert_eq!(env.seq, origin[index].seq);
assert_eq!(env.recorded_at, origin[index].recorded_at);
}
match &child_prefix[0].event {
Event::GraphRunStarted { forked_from, .. } => {
let origin_link = forked_from.as_ref().expect("head carries the fork origin");
assert_eq!(origin_link.run_id, run(1));
assert_eq!(origin_link.through_seq.get(), 2);
assert_eq!(origin_link.from_node, "publish");
assert_eq!(origin_link.graph_hash, "sha256:graph");
assert_eq!(origin_link.acknowledged_writes, vec![4]);
}
other => panic!("head is not GraphRunStarted: {other:?}"),
}
let blank = |envelopes: &[EventEnvelope]| -> String {
let mut copy: Vec<EventEnvelope> = envelopes.to_vec();
for env in &mut copy {
env.run_id = child;
}
if let Event::GraphRunStarted { forked_from, .. } = &mut copy[0].event {
*forked_from = None;
}
serde_json::to_string(©).unwrap()
};
let parent_prefix: Vec<EventEnvelope> = origin[..3].to_vec();
let mut child_blanked = child_prefix.clone();
if let Event::GraphRunStarted { forked_from, .. } = &mut child_blanked[0].event {
*forked_from = None;
}
assert_eq!(
blank(&parent_prefix),
serde_json::to_string(&child_blanked).unwrap(),
"child prefix is byte-identical to the parent's modulo run id and forked_from"
);
}
}