mod common;
use std::collections::{BTreeMap, HashMap};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
use common::{
ConstTool, EchoTool, PassTool, ScriptedModel, agent_builder, event_kinds, fixed_clock,
fixed_random, fixed_run_id, text_response, tool_use_response,
};
use salvor_core::{Effect, Event, EventEnvelope, ReplayError, RunId, RunStatus, derive_state};
use salvor_engine::{EngineError, GraphOutcome, run_graph};
use salvor_graph::{
AgentSpec, BranchCondition, BranchSpec, FoldBody, FoldJoin, FoldSpec, GateSpec, Graph,
GraphBuilder, MapBody, MapSpec, ToolSpec,
};
use salvor_runtime::{
Agent, Budgets, RunCtx, RuntimeError, validate_against_schema, validate_extension_input,
};
use salvor_store::{EventStore, SqliteStore};
use salvor_tools::DynTool;
use serde_json::{Value, json};
use wiremock::MockServer;
const RESEARCH_HASH: &str =
"sha256:1111111111111111111111111111111111111111111111111111111111111111";
fn extension_input() -> Value {
json!({"extend": {"steps": 100}})
}
type ExecutionCounts = BTreeMap<String, usize>;
struct Executors {
agents: HashMap<String, Agent>,
tools: HashMap<String, Box<dyn DynTool>>,
counters: BTreeMap<String, Arc<AtomicUsize>>,
}
impl Executors {
fn new(
agents: HashMap<String, Agent>,
tools: Vec<(Box<dyn DynTool>, Arc<AtomicUsize>)>,
) -> Self {
let mut registry: HashMap<String, Box<dyn DynTool>> = HashMap::new();
let mut counters: BTreeMap<String, Arc<AtomicUsize>> = BTreeMap::new();
for (tool, calls) in tools {
counters.insert(tool.name().to_owned(), calls);
registry.insert(tool.name().to_owned(), tool);
}
Self {
agents,
tools: registry,
counters,
}
}
fn snapshot(&self) -> ExecutionCounts {
self.counters
.iter()
.map(|(name, calls)| (name.clone(), calls.load(Ordering::SeqCst)))
.collect()
}
fn zeroed(&self) -> ExecutionCounts {
self.counters.keys().map(|name| (name.clone(), 0)).collect()
}
fn total(&self) -> usize {
self.counters
.values()
.map(|calls| calls.load(Ordering::SeqCst))
.sum()
}
}
type BuildExecutors = Arc<dyn Fn() -> Executors + Send + Sync>;
type ResumeInput = Arc<dyn Fn(&RunStatus) -> Value + Send + Sync>;
struct Shape {
name: &'static str,
graph: Graph,
input: Value,
run_id: RunId,
expected_kinds: Vec<&'static str>,
control_write_executions: usize,
build: BuildExecutors,
resume: ResumeInput,
}
async fn drive_once(
shape: &Shape,
executors: &Executors,
store: &Arc<dyn EventStore>,
log: Vec<EventEnvelope>,
resume_input: Option<Value>,
) -> Result<GraphOutcome, EngineError> {
let mut ctx = RunCtx::with_hooks(
store.clone(),
shape.run_id,
log,
fixed_clock(),
fixed_random(),
)?;
if let Some(input) = resume_input {
ctx.set_resume_input(input);
}
run_graph(
&mut ctx,
&shape.graph,
&shape.input,
&executors.agents,
&executors.tools,
)
.await
}
async fn drive_to_completion(
shape: &Shape,
executors: &Executors,
store: &Arc<dyn EventStore>,
control: Option<&[EventEnvelope]>,
label: &str,
) {
for _ in 0..24 {
let log = store.read_log(shape.run_id).await.expect("log reads");
let state = derive_state(&log);
let resume_input = match &state.status {
RunStatus::Completed { .. } => return,
RunStatus::Failed { error } => panic!("{label}: the run failed: {error}"),
RunStatus::NeedsReconciliation => {
panic!("{label}: reconciliation must be handled before dispatch")
}
RunStatus::Suspended { input_schema, .. } => {
let input = (shape.resume)(&state.status);
validate_against_schema(&input, input_schema).unwrap_or_else(|error| {
panic!("{label}: the resume input is rejected: {error}")
});
if let Some(control) = control {
assert_recorded_resume(control, log.len(), &input, label);
}
Some(input)
}
RunStatus::BudgetExceeded { .. } => {
let input = (shape.resume)(&state.status);
validate_extension_input(&input).unwrap_or_else(|error| {
panic!("{label}: the extension input is rejected: {error}")
});
if let Some(control) = control {
assert_recorded_resume(control, log.len(), &input, label);
}
Some(input)
}
_ => None,
};
drive_once(shape, executors, store, log, resume_input)
.await
.unwrap_or_else(|error| panic!("{label}: the drive failed: {error}"));
}
panic!("{label}: the run did not complete within the action cap");
}
fn assert_recorded_resume(control: &[EventEnvelope], position: usize, input: &Value, label: &str) {
match &control[position].event {
Event::Resumed { input: recorded } => assert_eq!(
recorded, input,
"{label}: the resume input must be the recorded one (control position {position})"
),
other => {
panic!("{label}: expected Resumed at control position {position}, found {other:?}")
}
}
}
fn ends_at_dangling_write_intent(prefix: &[EventEnvelope]) -> bool {
matches!(
prefix.last().map(|envelope| &envelope.event),
Some(Event::ToolCallRequested {
effect: Effect::Write,
..
})
)
}
fn expected_executions(
control: &[EventEnvelope],
cut: usize,
zeroed: &ExecutionCounts,
) -> ExecutionCounts {
let mut counts = zeroed.clone();
for (index, envelope) in control.iter().enumerate() {
let Event::ToolCallRequested { tool, effect, .. } = &envelope.event else {
continue;
};
let runs = if index >= cut {
1
} else if index + 1 == cut && !matches!(effect, Effect::Write) {
1
} else {
0
};
let slot = counts
.get_mut(tool.as_str())
.unwrap_or_else(|| panic!("unregistered tool in the log: {tool}"));
*slot += runs;
}
counts
}
fn prefix_write_executions(control: &[EventEnvelope], cut: usize) -> usize {
control[..cut]
.iter()
.enumerate()
.filter(|(index, envelope)| {
matches!(envelope.event, Event::ToolCallCompleted { .. })
&& matches!(
control.get(index.wrapping_sub(1)).map(|e| &e.event),
Some(Event::ToolCallRequested {
effect: Effect::Write,
..
})
)
})
.count()
}
fn live_write_executions(control: &[EventEnvelope], executors: &Executors) -> usize {
let counts = executors.snapshot();
let mut write_tools: BTreeMap<&str, ()> = BTreeMap::new();
for envelope in control {
if let Event::ToolCallRequested {
tool,
effect: Effect::Write,
..
} = &envelope.event
{
write_tools.insert(tool.as_str(), ());
}
}
write_tools
.keys()
.map(|tool| counts.get(*tool).copied().unwrap_or_default())
.sum()
}
struct SweepDir(PathBuf);
impl SweepDir {
fn create(owner: &str) -> Self {
let path =
std::env::temp_dir().join(format!("salvor-graph-gate-{owner}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&path);
std::fs::create_dir_all(&path).expect("sweep dir creates");
Self(path)
}
fn path(&self) -> &Path {
&self.0
}
}
impl Drop for SweepDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
async fn record_control(shape: &Shape, dir: &Path) -> Vec<EventEnvelope> {
let store: Arc<dyn EventStore> = Arc::new(
SqliteStore::open(dir.join(format!("{}-control.db", shape.name)))
.expect("control store opens"),
);
let executors = (shape.build)();
drive_to_completion(shape, &executors, &store, None, "control").await;
let log = store
.read_log(shape.run_id)
.await
.expect("control log reads");
assert_eq!(
event_kinds(&log),
shape.expected_kinds,
"{}: the control run's recorded shape is pinned",
shape.name
);
assert_eq!(
executors.snapshot(),
expected_executions(&log, 0, &executors.zeroed()),
"{}: the control run executes every recorded call exactly once",
shape.name
);
assert_eq!(
live_write_executions(&log, &executors),
shape.control_write_executions,
"{}: the control run's write total is the declared one",
shape.name
);
log
}
async fn continue_from_boundary(
cut: usize,
shape: Arc<Shape>,
control: Arc<Vec<EventEnvelope>>,
db_path: PathBuf,
) -> bool {
let label = format!("{} cut {cut}", shape.name);
let store: Arc<dyn EventStore> =
Arc::new(SqliteStore::open(&db_path).expect("boundary store opens"));
for envelope in &control[..cut] {
store.append(envelope).await.expect("prefix event appends");
}
let executors = (shape.build)();
if ends_at_dangling_write_intent(&control[..cut]) {
assert_eq!(
derive_state(&control[..cut]).status,
RunStatus::NeedsReconciliation,
"{label}: a dangling write intent derives to reconciliation"
);
let prefix = control[..cut].to_vec();
let error = drive_once(&shape, &executors, &store, prefix, None)
.await
.expect_err("a continuation over a dangling write intent must refuse");
assert!(
matches!(
error,
EngineError::Runtime(RuntimeError::Replay(
ReplayError::NeedsReconciliation { .. }
))
),
"{label}: expected the reconciliation error, got {error:?}"
);
assert_eq!(
executors.total(),
0,
"{label}: a refused continuation executes nothing"
);
assert_eq!(
store.read_log(shape.run_id).await.expect("log reads").len(),
cut,
"{label}: a refused continuation appends nothing"
);
return true;
}
if cut == control.len() {
let outcome = drive_once(&shape, &executors, &store, control[..].to_vec(), None)
.await
.expect("full-log replay is divergence free");
assert!(
matches!(outcome, GraphOutcome::Completed { .. }),
"{label}: a full-log replay completes, got {outcome:?}"
);
} else {
drive_to_completion(&shape, &executors, &store, Some(&control), &label).await;
}
let log = store.read_log(shape.run_id).await.expect("final log reads");
assert_eq!(
serde_json::to_string(&log).unwrap(),
serde_json::to_string(&*control).unwrap(),
"{label}: the continued log must be byte-identical to the control log"
);
let expected = expected_executions(&control, cut, &executors.zeroed());
assert_eq!(
executors.snapshot(),
expected,
"{label}: live executions must match the re-execution table"
);
assert_eq!(
prefix_write_executions(&control, cut) + live_write_executions(&control, &executors),
shape.control_write_executions,
"{label}: write executions across the whole scenario must total the control count"
);
false
}
struct SweepReport {
name: &'static str,
events: usize,
boundaries: usize,
refusals: usize,
}
async fn sweep(shape: Shape, dir: &Path) -> SweepReport {
let name = shape.name;
let control = Arc::new(record_control(&shape, dir).await);
let shape = Arc::new(shape);
let mut boundaries = tokio::task::JoinSet::new();
for cut in 0..=control.len() {
let shape = shape.clone();
let control = control.clone();
let db_path = dir.join(format!("{name}-boundary-{cut:03}.db"));
boundaries.spawn(async move { continue_from_boundary(cut, shape, control, db_path).await });
}
let mut refusals = 0;
let mut completed = 0;
while let Some(result) = boundaries.join_next().await {
match result {
Ok(was_refusal) => {
refusals += usize::from(was_refusal);
completed += 1;
}
Err(error) => std::panic::resume_unwind(error.into_panic()),
}
}
assert_eq!(
completed,
control.len() + 1,
"{name}: every boundary was swept"
);
assert_eq!(
refusals, shape.control_write_executions,
"{name}: exactly the write-intent prefixes refuse; every other boundary completes"
);
SweepReport {
name,
events: control.len(),
boundaries: completed,
refusals,
}
}
fn approval_schema() -> Value {
json!({
"type": "object",
"properties": {"approved": {"type": "boolean"}},
"required": ["approved"]
})
}
fn targets_schema() -> Value {
json!({
"type": "object",
"properties": {"approved": {"type": "boolean"}, "targets": {"type": "array"}},
"required": ["approved", "targets"]
})
}
fn gate_shape(research_uri: &str) -> Shape {
let graph = GraphBuilder::new()
.agent(AgentSpec::new("research", RESEARCH_HASH))
.gate(
GateSpec::new("approve", approval_schema())
.prompt("Approve this draft for publication?"),
)
.tool(ToolSpec::new("publish", "http_post"))
.edge("research", "approve")
.edge("approve", "publish")
.build();
let uri = research_uri.to_owned();
Shape {
name: "gate",
graph,
input: json!({"topic": "otters"}),
run_id: fixed_run_id(50),
expected_kinds: vec![
"GraphRunStarted",
"NodeEntered", "NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"NodeExited", "NodeEntered", "Suspended",
"Resumed",
"NodeExited", "NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "RunCompleted",
],
control_write_executions: 1,
build: Arc::new(move || {
let mut agents: HashMap<String, Agent> = HashMap::new();
agents.insert(
RESEARCH_HASH.to_owned(),
agent_builder(&uri).build().expect("the agent builds"),
);
let (publish, publish_calls) = EchoTool::new("http_post", Effect::Write);
Executors::new(agents, vec![(Box::new(publish), publish_calls)])
}),
resume: Arc::new(|status| match status {
RunStatus::Suspended { .. } => json!({"approved": true}),
other => panic!("the gate shape parks only at its gate, got {other:?}"),
}),
}
}
fn branch_shape() -> Shape {
let graph = GraphBuilder::new()
.tool(ToolSpec::new("assess", "assess"))
.branch(
BranchSpec::new("route")
.on("score")
.case("high", BranchCondition::Expression("score >= 0.8".into()))
.case("low", BranchCondition::Expression("score < 0.8".into())),
)
.tool(ToolSpec::new("publish", "http_post"))
.tool(ToolSpec::new("reject", "notify"))
.edge("assess", "route")
.labeled_edge("route", "publish", "high")
.labeled_edge("route", "reject", "low")
.build();
Shape {
name: "branch",
graph,
input: json!({"topic": "otters"}),
run_id: fixed_run_id(51),
expected_kinds: vec![
"GraphRunStarted",
"NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "NodeEntered", "BranchTaken",
"NodeExited", "NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "NodeSkipped", "RunCompleted",
],
control_write_executions: 1,
build: Arc::new(|| {
let (assess, assess_calls) =
ConstTool::new("assess", Effect::Read, json!({"score": 0.9}));
let (publish, publish_calls) = EchoTool::new("http_post", Effect::Write);
let (notify, notify_calls) = EchoTool::new("notify", Effect::Write);
Executors::new(
HashMap::new(),
vec![
(Box::new(assess), assess_calls),
(Box::new(publish), publish_calls),
(Box::new(notify), notify_calls),
],
)
}),
resume: Arc::new(|status| panic!("the branch shape never parks, got {status:?}")),
}
}
fn map_shape() -> Shape {
let graph = GraphBuilder::new()
.map(MapSpec::new(
"fanout",
"targets",
2,
MapBody::Node("worker".into()),
))
.tool(ToolSpec::new("worker", "publish_worker"))
.tool(ToolSpec::new("record", "record_tool"))
.edge("fanout", "record")
.build();
let mut expected_kinds = vec!["GraphRunStarted", "NodeEntered", "MapFannedOut"];
for _ in 0..3 {
expected_kinds.extend([
"MapIterationStarted",
"ToolCallRequested",
"ToolCallCompleted",
"MapIterationJoined",
]);
}
expected_kinds.extend([
"NodeExited", "NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "RunCompleted",
]);
Shape {
name: "map",
graph,
input: json!({"targets": ["alpha", "beta", "gamma"]}),
run_id: fixed_run_id(52),
expected_kinds,
control_write_executions: 3,
build: Arc::new(|| {
let (worker, worker_calls) = EchoTool::new("publish_worker", Effect::Write);
let (record, record_calls) = EchoTool::new("record_tool", Effect::Idempotent);
Executors::new(
HashMap::new(),
vec![
(Box::new(worker), worker_calls),
(Box::new(record), record_calls),
],
)
}),
resume: Arc::new(|status| panic!("the map shape never parks, got {status:?}")),
}
}
fn fold_shape() -> Shape {
let graph = GraphBuilder::new()
.fold(FoldSpec::new(
"refine",
FoldBody::Node("worker".into()),
3,
"score >= 99",
FoldJoin::BestBy("score".into()),
))
.tool(ToolSpec::new("worker", "publish_worker"))
.tool(ToolSpec::new("record", "record_tool"))
.edge("refine", "record")
.build();
let mut expected_kinds = vec!["GraphRunStarted", "NodeEntered"];
for _ in 0..3 {
expected_kinds.extend([
"FoldIterationStarted",
"ToolCallRequested",
"ToolCallCompleted",
"FoldIterationJoined",
]);
}
expected_kinds.extend([
"FoldConverged",
"NodeExited", "NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "RunCompleted",
]);
Shape {
name: "fold",
graph,
input: json!({"pass": 0}),
run_id: fixed_run_id(54),
expected_kinds,
control_write_executions: 3,
build: Arc::new(|| {
let (worker, worker_calls) = PassTool::new(
"publish_worker",
Effect::Write,
vec![json!(1), json!(5), json!(2)],
);
let (record, record_calls) = EchoTool::new("record_tool", Effect::Idempotent);
Executors::new(
HashMap::new(),
vec![
(Box::new(worker), worker_calls),
(Box::new(record), record_calls),
],
)
}),
resume: Arc::new(|status| panic!("the fold shape never parks, got {status:?}")),
}
}
fn flagship_shape(research_uri: &str) -> Shape {
let graph = GraphBuilder::new()
.agent(AgentSpec::new("research", RESEARCH_HASH))
.tool(ToolSpec::new("assess", "assess"))
.branch(
BranchSpec::new("route")
.on("score")
.case("high", BranchCondition::Expression("score >= 0.8".into()))
.case("low", BranchCondition::Expression("score < 0.8".into())),
)
.gate(
GateSpec::new("approve", targets_schema())
.prompt("Approve these targets for publication?"),
)
.map(MapSpec::new(
"fanout",
"targets",
2,
MapBody::Node("worker".into()),
))
.tool(ToolSpec::new("worker", "publish_worker"))
.tool(ToolSpec::new("record", "record_tool"))
.tool(ToolSpec::new("reject", "notify"))
.edge("research", "assess")
.edge("assess", "route")
.labeled_edge("route", "approve", "high")
.labeled_edge("route", "reject", "low")
.edge("approve", "fanout")
.edge("fanout", "record")
.build();
let uri = research_uri.to_owned();
let mut expected_kinds = vec![
"GraphRunStarted",
"NodeEntered", "NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"ToolCallRequested", "ToolCallCompleted",
"NowObserved",
"BudgetExceeded", "Resumed",
"ModelCallRequested",
"ModelCallCompleted",
"NodeExited", "NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "NodeEntered", "BranchTaken",
"NodeExited", "NodeEntered", "Suspended",
"Resumed",
"NodeExited", "NodeEntered", "MapFannedOut",
];
for _ in 0..2 {
expected_kinds.extend([
"MapIterationStarted",
"ToolCallRequested",
"ToolCallCompleted",
"MapIterationJoined",
]);
}
expected_kinds.extend([
"NodeExited", "NodeEntered", "ToolCallRequested",
"ToolCallCompleted",
"NodeExited", "NodeSkipped", "RunCompleted",
]);
Shape {
name: "flagship",
graph,
input: json!({"topic": "otters"}),
run_id: fixed_run_id(53),
expected_kinds,
control_write_executions: 2,
build: Arc::new(move || {
let (lookup, lookup_calls) = EchoTool::new("lookup", Effect::Read);
let mut agents: HashMap<String, Agent> = HashMap::new();
agents.insert(
RESEARCH_HASH.to_owned(),
agent_builder(&uri)
.tool_dyn(Box::new(lookup))
.budgets(Budgets {
max_steps: Some(1),
..Budgets::default()
})
.build()
.expect("the agent builds"),
);
let (assess, assess_calls) =
ConstTool::new("assess", Effect::Read, json!({"score": 0.9}));
let (worker, worker_calls) = EchoTool::new("publish_worker", Effect::Write);
let (record, record_calls) = EchoTool::new("record_tool", Effect::Idempotent);
let (notify, notify_calls) = EchoTool::new("notify", Effect::Write);
let mut executors = Executors::new(
agents,
vec![
(Box::new(assess), assess_calls),
(Box::new(worker), worker_calls),
(Box::new(record), record_calls),
(Box::new(notify), notify_calls),
],
);
executors.counters.insert("lookup".to_owned(), lookup_calls);
executors
}),
resume: Arc::new(|status| match status {
RunStatus::Suspended { .. } => json!({"approved": true, "targets": ["alpha", "beta"]}),
RunStatus::BudgetExceeded { .. } => extension_input(),
other => panic!("the flagship shape parks only at its gate or budget, got {other:?}"),
}),
}
}
async fn gate_server() -> MockServer {
ScriptedModel::mount(vec![(1, text_response("a draft about otters", 5, 3))]).await
}
async fn flagship_server() -> MockServer {
ScriptedModel::mount(vec![
(
1,
tool_use_response("tu_01", "lookup", json!({"q": "otters"}), 101, 11),
),
(3, text_response("a draft about otters", 102, 12)),
])
.await
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
async fn graph_release_gate_kill_at_every_event_boundary_resumes_identically() {
let started = Instant::now();
let dir = SweepDir::create("sweep");
let gate_model = gate_server().await;
let flagship_model = flagship_server().await;
let shapes = vec![
gate_shape(&gate_model.uri()),
branch_shape(),
map_shape(),
fold_shape(),
flagship_shape(&flagship_model.uri()),
];
let mut reports = Vec::with_capacity(shapes.len());
for shape in shapes {
reports.push(sweep(shape, dir.path()).await);
}
let counts: Vec<(&str, usize, usize, usize)> = reports
.iter()
.map(|report| {
(
report.name,
report.events,
report.boundaries,
report.refusals,
)
})
.collect();
assert_eq!(
counts,
vec![
("gate", 15, 16, 1),
("branch", 14, 15, 1),
("map", 21, 22, 3),
("fold", 21, 22, 3),
("flagship", 41, 42, 2),
],
"each shape sweeps every boundary of its control log and refuses at each write intent"
);
let total: usize = reports.iter().map(|report| report.boundaries).sum();
eprintln!(
"graph release gate: {total} boundaries swept across {} shapes in {:.2?}",
reports.len(),
started.elapsed()
);
}
#[tokio::test]
async fn resuming_a_graph_run_must_re_supply_the_matching_document() {
let dir = SweepDir::create("re-supply");
let model = gate_server().await;
let shape = gate_shape(&model.uri());
let control = record_control(&shape, dir.path()).await;
let park = control
.iter()
.position(|envelope| matches!(envelope.event, Event::Resumed { .. }))
.expect("the control run parked at its gate");
let store: Arc<dyn EventStore> =
Arc::new(SqliteStore::open(dir.path().join("re-supply.db")).expect("boundary store opens"));
for envelope in &control[..park] {
store.append(envelope).await.expect("prefix event appends");
}
let executors = (shape.build)();
let mut altered = gate_shape(&model.uri());
altered.graph = GraphBuilder::new()
.agent(AgentSpec::new("research", RESEARCH_HASH))
.gate(GateSpec::new("approve", approval_schema()).prompt("Approve? (reworded)"))
.tool(ToolSpec::new("publish", "http_post"))
.edge("research", "approve")
.edge("approve", "publish")
.build();
let error = drive_once(
&altered,
&executors,
&store,
control[..park].to_vec(),
Some(json!({"approved": true})),
)
.await
.expect_err("a changed document must not resume an old run");
assert!(
matches!(
error,
EngineError::Runtime(RuntimeError::Replay(ReplayError::Divergence { .. }))
),
"expected a divergence at the recorded head, got {error:?}"
);
assert_eq!(
store.read_log(shape.run_id).await.expect("log reads").len(),
park,
"the refused resume appended nothing"
);
assert_eq!(executors.total(), 0, "the refused resume executed nothing");
let mut wrong_input = gate_shape(&model.uri());
wrong_input.input = json!({"topic": "not the recorded topic"});
drive_to_completion(
&wrong_input,
&executors,
&store,
Some(&control),
"re-supply",
)
.await;
let log = store.read_log(shape.run_id).await.expect("log reads");
assert_eq!(
serde_json::to_string(&log).unwrap(),
serde_json::to_string(&control).unwrap(),
"the recorded input wins over the re-supplied one, byte for byte"
);
}