use std::collections::{BTreeMap, VecDeque};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use async_trait::async_trait;
use pointlock_ir::{
ActionOutcome, ActionResult, AssetRef, EventCursor, FlowIR, Hash, IterState, Observation,
PathFrame, ReconcileResult, RunLogEvent, RunLogPayload, StepIR, StepState, VerdictStatus,
};
use pointlock_provider_kit::{
BoundActionCall, CancellationToken, CapabilityAttestation, EvidenceStream, FakeProvider,
ObserveRequest, Provider, ProviderError, ProviderSession, ScriptedOutcome, SessionHealth,
SessionOutcome, UiSnapshotOutcome, VerdictWrite,
};
use pointlock_runner::{ResumeOptions, RunOptions, RunOutcome, Runner};
use pointlock_store::Store;
use serde_json::{Value, json};
static DIR_COUNTER: AtomicU64 = AtomicU64::new(0);
struct TempStoreDir(PathBuf);
impl TempStoreDir {
fn new(tag: &str) -> Self {
let path = std::env::temp_dir().join(format!(
"pointlock-control-flow-test-{tag}-{}-{}",
std::process::id(),
DIR_COUNTER.fetch_add(1, Ordering::Relaxed),
));
std::fs::create_dir_all(&path).expect("create temp store dir");
TempStoreDir(path)
}
fn path(&self) -> &std::path::Path {
&self.0
}
}
impl Drop for TempStoreDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn h64(fill: char) -> String {
format!("sha256:{}", fill.to_string().repeat(64))
}
fn seal(flow: &mut FlowIR) {
fn seal_steps(steps: &mut [StepIR]) {
for step in steps.iter_mut() {
match step {
StepIR::If(s) => {
seal_steps(&mut s.then);
if let Some(otherwise) = &mut s.r#else {
seal_steps(otherwise);
}
}
StepIR::Foreach(s) => seal_steps(&mut s.body),
_ => {}
}
let effect = pointlock_ir::effect_hash(step);
let judge = pointlock_ir::judge_hash(step);
let base = match step {
StepIR::Action(s) => &mut s.base,
StepIR::Assert(s) => &mut s.base,
StepIR::Call(s) => &mut s.base,
StepIR::Human(s) => &mut s.base,
StepIR::If(s) => &mut s.base,
StepIR::Foreach(s) => &mut s.base,
StepIR::Let(s) => &mut s.base,
};
base.effect_hash = effect;
base.judge_hash = judge;
}
}
seal_steps(&mut flow.body);
flow.ir_hash = pointlock_ir::ir_hash(flow);
}
fn build_flow(
flow_id: &str,
lockfile_digest: &Hash,
params: Value,
outputs: Value,
body: Vec<Value>,
subflows: Value,
) -> FlowIR {
let mut flow: FlowIR = serde_json::from_value(json!({
"irVersion": 1,
"flowId": flow_id,
"irHash": h64('e'),
"provider": { "name": "devicerail", "version": "0.1.0" },
"requiredFeatures": [],
"lockfileDigest": lockfile_digest.as_str(),
"params": params,
"outputs": outputs,
"body": body,
"verdictPolicy": "standard",
"sourceMap": [],
"subflows": subflows
}))
.expect("fixture is a valid FlowIR");
seal(&mut flow);
flow
}
fn flow_fixture(lockfile_digest: &Hash, body: Vec<Value>) -> FlowIR {
build_flow(
"cf_demo",
lockfile_digest,
json!([]),
json!([]),
body,
json!({}),
)
}
fn expect_ok(assert_id: &str, sid: &str) -> Value {
json!({
"assertId": assert_id,
"predicate": { "type": "expr", "expr": { "fn": "eq", "args": [
{ "ref": format!("steps.{sid}.output.ok") },
{ "lit": true }
] } },
"verifyVia": [],
"onMissingInput": "unknown"
})
}
fn action_step(id: &str, assertions: Vec<Value>) -> Value {
json!({
"kind": "action",
"stepId": id,
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"effect": "mutating",
"idempotent": false,
"binding": { "attempts": [ {
"channel": "uiTree",
"actionName": "tapElement",
"args": { "element": { "lit": { "identifier": id } } },
"acceptExecutionModes": ["nativeSemantic", "webSemantic"],
"protection": "standard"
} ] },
"assertions": assertions
})
}
fn succeeded_with(output: Value) -> ScriptedOutcome {
ScriptedOutcome::Terminal(ActionOutcome::Succeeded {
result: Box::new(ActionResult {
call_id: String::new(),
started_at_ms: 0,
finished_at_ms: 0,
output,
before: None,
after: None,
evidence: Vec::new(),
execution: None,
}),
})
}
fn succeeded_with_after(output: Value, after: Observation) -> ScriptedOutcome {
ScriptedOutcome::Terminal(ActionOutcome::Succeeded {
result: Box::new(ActionResult {
call_id: String::new(),
started_at_ms: 0,
finished_at_ms: 0,
output,
before: None,
after: Some(after),
evidence: Vec::new(),
execution: None,
}),
})
}
fn tree(nodes: Value) -> Value {
json!({
"formatVersion": 1,
"observationId": "stamped-by-fake",
"context": { "contextKind": "native", "contextId": "ctx-1", "documentEpoch": "e1" },
"rootStableNodeIds": ["n1"],
"nodes": nodes
})
}
async fn open(provider: &FakeProvider) -> Box<dyn ProviderSession> {
provider
.open_session(provider.default_open_options())
.await
.expect("open_session")
}
fn event_types(store: &Store, run_id: &str) -> Vec<&'static str> {
store
.events(run_id)
.expect("events")
.iter()
.map(|event| event.payload.event_type())
.collect()
}
fn run_opts(run_id: &str) -> RunOptions {
let mut opts = RunOptions::new("fake-device-1");
opts.run_id = Some(run_id.to_owned());
opts
}
fn step_of(event: &RunLogEvent) -> Option<&str> {
event.run_path.iter().rev().find_map(|frame| match frame {
PathFrame::Step { step_id } => Some(step_id.as_str()),
PathFrame::Call {
step_id: Some(step_id),
..
} => Some(step_id.as_str()),
_ => None,
})
}
fn finished_verdict(outcome: RunOutcome) -> pointlock_ir::Verdict {
let RunOutcome::Finished {
verdict: Some(verdict),
} = outcome
else {
panic!("expected Finished with a verdict, got {outcome:?}");
};
verdict
}
struct StopAfter {
inner: Box<dyn ProviderSession>,
remaining: AtomicUsize,
stop: CancellationToken,
}
#[async_trait]
impl ProviderSession for StopAfter {
fn attestation(&self) -> &CapabilityAttestation {
self.inner.attestation()
}
async fn execute(
&self,
call: BoundActionCall,
cancel: Option<CancellationToken>,
) -> Result<ActionOutcome, ProviderError> {
let outcome = self.inner.execute(call, cancel).await;
if self.remaining.fetch_sub(1, Ordering::SeqCst) == 1 {
self.stop.cancel();
}
outcome
}
async fn observe(
&self,
req: ObserveRequest,
cancel: Option<CancellationToken>,
) -> Result<Observation, ProviderError> {
self.inner.observe(req, cancel).await
}
async fn ui_snapshot(&self, observation_id: &str) -> Result<UiSnapshotOutcome, ProviderError> {
self.inner.ui_snapshot(observation_id).await
}
async fn reconcile(
&self,
call_id: &str,
issuing: &pointlock_ir::EventCursor,
) -> Result<ReconcileResult, ProviderError> {
self.inner.reconcile(call_id, issuing).await
}
async fn fetch_evidence(&self, asset: &AssetRef) -> Result<EvidenceStream, ProviderError> {
self.inner.fetch_evidence(asset).await
}
async fn record_verdict(&self, verdict: VerdictWrite) -> Result<(), ProviderError> {
self.inner.record_verdict(verdict).await
}
async fn current_cursor(&self) -> Result<EventCursor, ProviderError> {
self.inner.current_cursor().await
}
async fn health(&self) -> Result<SessionHealth, ProviderError> {
self.inner.health().await
}
async fn end(
&self,
outcome: SessionOutcome,
reason: Option<String>,
) -> Result<(), ProviderError> {
self.inner.end(outcome, reason).await
}
}
fn nested_fixture(digest: &Hash) -> (FlowIR, BTreeMap<Hash, FlowIR>) {
let inner = build_flow(
"inner",
digest,
json!([ { "name": "label", "schema": { "type": "string" }, "required": true } ]),
json!([ { "name": "echo", "schema": { "type": "string" },
"from": { "ref": "params.label" } } ]),
vec![action_step("g1", vec![expect_ok("ga", "g1")])],
json!({}),
);
let mid = build_flow(
"mid",
digest,
json!([ { "name": "label", "schema": { "type": "string" }, "required": true } ]),
json!([ { "name": "midEcho", "schema": { "type": "string" },
"from": { "ref": "steps.callInner.output.echo" } } ]),
vec![
json!({
"kind": "call",
"stepId": "callInner",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "inner", "irHash": inner.ir_hash.as_str() },
"inputs": { "label": { "ref": "params.label" } }
}),
action_step("m1", vec![expect_ok("ma", "m1")]),
],
json!({ "inner": { "flowId": "inner", "irHash": inner.ir_hash.as_str() } }),
);
let root = build_flow(
"root",
digest,
json!([]),
json!([]),
vec![
action_step("r1", vec![expect_ok("ra", "r1")]),
json!({
"kind": "call",
"stepId": "callMid",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "mid", "irHash": mid.ir_hash.as_str() },
"inputs": { "label": { "lit": "hello" } }
}),
json!({
"kind": "action",
"stepId": "r2",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"effect": "mutating",
"idempotent": false,
"binding": { "attempts": [ {
"channel": "uiTree",
"actionName": "tapElement",
"args": { "element": { "lit": { "identifier": "r2" } } },
"acceptExecutionModes": ["nativeSemantic", "webSemantic"],
"protection": "standard"
} ] },
"assertions": [ {
"assertId": "echoBack",
"predicate": { "type": "expr", "expr": { "fn": "eq", "args": [
{ "ref": "steps.callMid.output.midEcho" },
{ "lit": "hello" }
] } },
"verifyVia": [],
"onMissingInput": "unknown"
} ]
}),
],
json!({ "mid": { "flowId": "mid", "irHash": mid.ir_hash.as_str() } }),
);
let mut registry = BTreeMap::new();
registry.insert(inner.ir_hash.clone(), inner);
registry.insert(mid.ir_hash.clone(), mid);
(root, registry)
}
#[tokio::test]
async fn nested_two_layer_call_passes_end_to_end() {
let dir = TempStoreDir::new("nested");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let (root, registry) = nested_fixture(&provider.lockfile().digest);
let mid_flow = registry
.values()
.find(|flow| flow.flow_id.as_str() == "mid")
.expect("mid registered");
let inner_flow = registry
.values()
.find(|flow| flow.flow_id.as_str() == "inner")
.expect("inner registered");
let (mid_hash8, inner_hash8) = (
mid_flow.ir_hash.hex_prefix8().to_owned(),
inner_flow.ir_hash.hex_prefix8().to_owned(),
);
let mut opts = run_opts("run-nested");
opts.subflows = registry;
let session = open(&provider).await;
let outcome = Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Pass);
assert!(!verdict.degraded);
let action_block = [
"stepEntered",
"actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
];
let mut expected = vec!["runStarted"];
expected.extend(action_block); expected.extend(["stepEntered", "callFramePushed"]); expected.extend(["stepEntered", "callFramePushed"]); expected.extend(action_block); expected.extend(["callFramePopped", "verdictRecorded", "stepExited"]); expected.extend(action_block); expected.extend(["callFramePopped", "verdictRecorded", "stepExited"]); expected.extend(action_block); expected.push("runFinished");
assert_eq!(event_types(&store, "run-nested"), expected);
let view = store.verify_checkpoint("run-nested").expect("verify");
assert_eq!(view.frames.len(), 1);
assert_eq!(view.frames[0].next_index, 3);
let g1 = view
.completed
.iter()
.find(|record| record.step_id.as_str() == "g1")
.expect("g1 record");
assert_eq!(
pointlock_ir::render_run_path(&g1.run_path),
format!(
"root@{}/callMid/call→mid@{mid_hash8}/callInner/call→inner@{inner_hash8}/g1",
root.ir_hash.hex_prefix8(),
)
);
let call_mid = view
.completed
.iter()
.find(|record| record.step_id.as_str() == "callMid")
.expect("callMid record");
assert_eq!(call_mid.resolved_inputs, json!({ "label": "hello" }));
assert_eq!(call_mid.output, Some(json!({ "midEcho": "hello" })));
assert_eq!(
call_mid.verdict.as_ref().map(|verdict| verdict.status),
Some(VerdictStatus::Pass)
);
assert_eq!(provider.handle().dispatched_call_ids().len(), 4);
}
#[tokio::test]
async fn call_inbound_gate_refuses_schema_invalid_inputs() {
let dir = TempStoreDir::new("gate-in");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::new());
let digest = provider.lockfile().digest.clone();
let callee = build_flow(
"callee",
&digest,
json!([ { "name": "label", "schema": { "type": "string" }, "required": true } ]),
json!([]),
vec![action_step("c1", vec![])],
json!({}),
);
let root = build_flow(
"root",
&digest,
json!([]),
json!([]),
vec![
json!({
"kind": "call",
"stepId": "callBad",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "callee", "irHash": callee.ir_hash.as_str() },
"inputs": { "label": { "lit": 42 } }
}),
action_step("after", vec![]),
],
json!({ "callee": { "flowId": "callee", "irHash": callee.ir_hash.as_str() } }),
);
let mut registry = BTreeMap::new();
registry.insert(callee.ir_hash.clone(), callee);
let mut opts = run_opts("run-gate-in");
opts.subflows = registry;
let session = open(&provider).await;
let outcome = Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Fail);
let types = event_types(&store, "run-gate-in");
assert!(!types.contains(&"callFramePushed"));
assert!(provider.handle().dispatched_call_ids().is_empty());
let events = store.events("run-gate-in").expect("events");
let call_verdict = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::VerdictRecorded { verdict, .. } if step_of(event) == Some("callBad") => {
Some(verdict.clone())
}
_ => None,
})
.expect("call step verdict");
assert_eq!(call_verdict.status, VerdictStatus::Fail);
assert!(
call_verdict.summary.contains("bind_arguments_invalid"),
"got {}",
call_verdict.summary
);
let view = store.verify_checkpoint("run-gate-in").expect("verify");
let after = view
.completed
.iter()
.find(|record| record.step_id.as_str() == "after")
.expect("after record");
assert!(after.attempts.is_empty() && after.verdict.is_none());
}
#[tokio::test]
async fn call_outbound_gate_fails_on_schema_invalid_outputs() {
let dir = TempStoreDir::new("gate-out");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
let digest = provider.lockfile().digest.clone();
let callee = build_flow(
"callee",
&digest,
json!([]),
json!([ { "name": "flag", "schema": { "type": "string" },
"from": { "ref": "steps.c1.output.ok" } } ]),
vec![action_step("c1", vec![expect_ok("ca", "c1")])],
json!({}),
);
let root = build_flow(
"root",
&digest,
json!([]),
json!([]),
vec![json!({
"kind": "call",
"stepId": "callBad",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "callee", "irHash": callee.ir_hash.as_str() },
"inputs": {}
})],
json!({ "callee": { "flowId": "callee", "irHash": callee.ir_hash.as_str() } }),
);
let mut registry = BTreeMap::new();
registry.insert(callee.ir_hash.clone(), callee);
let mut opts = run_opts("run-gate-out");
opts.subflows = registry;
let session = open(&provider).await;
let outcome = Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Fail);
assert_eq!(provider.handle().dispatched_call_ids().len(), 1);
let events = store.events("run-gate-out").expect("events");
let popped = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::CallFramePopped { outputs } => Some(outputs.clone()),
_ => None,
})
.expect("frame popped");
assert_eq!(popped, None);
let call_verdict = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::VerdictRecorded { verdict, .. } if step_of(event) == Some("callBad") => {
Some(verdict.clone())
}
_ => None,
})
.expect("call step verdict");
assert!(
call_verdict.summary.contains("outbound gate"),
"got {}",
call_verdict.summary
);
}
fn foreach_flow(digest: &Hash, tail_step: bool) -> FlowIR {
let mut body = vec![json!({
"kind": "foreach",
"stepId": "eachItem",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"items": { "lit": ["a", "b", "c"] },
"as": "item",
"body": [ {
"kind": "action",
"stepId": "perItem",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"effect": "mutating",
"idempotent": true,
"binding": { "attempts": [ {
"channel": "uiTree",
"actionName": "tapElement",
"args": { "value": { "ref": "iter.item" } },
"acceptExecutionModes": ["nativeSemantic", "webSemantic"],
"protection": "standard"
} ] },
"assertions": [ expect_ok("pa", "perItem") ]
} ]
})];
if tail_step {
body.push(action_step("tail", vec![]));
}
build_flow("fe_flow", digest, json!([]), json!([]), body, json!({}))
}
#[tokio::test]
async fn foreach_runs_three_iterations_with_positional_frames() {
let dir = TempStoreDir::new("foreach");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = foreach_flow(&provider.lockfile().digest, false);
let session = open(&provider).await;
let outcome = Runner::run(&flow, json!({}), session, &mut store, run_opts("run-fe"))
.await
.expect("run");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Pass);
assert!(
verdict.summary.contains("3 judged step(s)"),
"got {}",
verdict.summary
);
let events = store.events("run-fe").expect("events");
let fe_entered = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::StepEntered {
step_id,
resolved_inputs,
..
} if step_id.as_str() == "eachItem" => Some(resolved_inputs.clone()),
_ => None,
})
.expect("foreach entered");
assert_eq!(
fe_entered,
json!({ "items": ["a", "b", "c"], "as": "item" })
);
let intents: Vec<(u64, Value)> = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::ActionIntent { args_snapshot, .. } => {
let index = event.run_path.iter().find_map(|frame| match frame {
PathFrame::Iteration { index, .. } => Some(*index),
_ => None,
})?;
Some((index, args_snapshot.clone()))
}
_ => None,
})
.collect();
assert_eq!(
intents,
vec![
(0, json!({ "value": "a" })),
(1, json!({ "value": "b" })),
(2, json!({ "value": "c" })),
]
);
let view = store.verify_checkpoint("run-fe").expect("verify");
assert_eq!(view.completed.len(), 4);
assert_eq!(view.frames[0].next_index, 1);
}
#[tokio::test]
async fn foreach_halts_on_a_failing_iteration_and_folds_fail() {
let dir = TempStoreDir::new("foreach-fail");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": false })), ]));
let flow = foreach_flow(&provider.lockfile().digest, true);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-fe-fail"),
)
.await
.expect("run");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Fail);
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let action_block = [
"stepEntered",
"actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
];
let mut expected = vec!["runStarted", "stepEntered"]; expected.extend(action_block); expected.extend(action_block); expected.push("stepExited"); expected.extend(["stepEntered", "stepExited"]); expected.push("runFinished");
assert_eq!(event_types(&store, "run-fe-fail"), expected);
store.verify_checkpoint("run-fe-fail").expect("verify");
}
fn if_flow(digest: &Hash) -> FlowIR {
build_flow(
"if_flow",
digest,
json!([ { "name": "mode", "schema": { "type": "string" }, "required": true } ]),
json!([]),
vec![json!({
"kind": "if",
"stepId": "branch",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"cond": { "fn": "eq", "args": [ { "ref": "params.mode" }, { "lit": "yes" } ] },
"then": [ action_step("t1", vec![expect_ok("ta", "t1")]) ],
"else": [ action_step("e1", vec![expect_ok("ea", "e1")]) ]
})],
json!({}),
)
}
#[tokio::test]
async fn if_selects_then_and_accounts_the_else_branch_as_skipped() {
let dir = TempStoreDir::new("if-then");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
let flow = if_flow(&provider.lockfile().digest);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({ "mode": "yes" }),
session,
&mut store,
run_opts("run-if-then"),
)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(
event_types(&store, "run-if-then"),
vec![
"runStarted",
"stepEntered", "stepEntered", "stepExited",
"stepEntered", "actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"stepExited", "runFinished",
]
);
let events = store.events("run-if-then").expect("events");
let cond = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::StepEntered {
step_id,
resolved_inputs,
..
} if step_id.as_str() == "branch" => Some(resolved_inputs.clone()),
_ => None,
})
.expect("branch entered");
assert_eq!(cond, json!({ "cond": true }));
let skipped_exit = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::StepExited { state, .. } if step_of(event) == Some("e1") => Some(*state),
_ => None,
})
.expect("e1 exit");
assert_eq!(skipped_exit, StepState::Skipped);
let view = store.verify_checkpoint("run-if-then").expect("verify");
let e1 = view
.completed
.iter()
.find(|record| record.step_id.as_str() == "e1")
.expect("e1 record");
assert_eq!(e1.resolved_inputs, Value::Null);
assert!(e1.attempts.is_empty() && e1.verdict.is_none() && e1.output.is_none());
assert_eq!(provider.handle().dispatched_call_ids().len(), 1);
}
#[tokio::test]
async fn if_selects_else_and_accounts_the_then_branch_as_skipped() {
let dir = TempStoreDir::new("if-else");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
let flow = if_flow(&provider.lockfile().digest);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({ "mode": "no" }),
session,
&mut store,
run_opts("run-if-else"),
)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
let events = store.events("run-if-else").expect("events");
let skipped: Vec<&str> = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::StepExited { state, .. } if *state == StepState::Skipped => {
step_of(event)
}
_ => None,
})
.collect();
assert_eq!(skipped, vec!["t1"]);
let executed: Vec<&str> = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::ActionIntent { .. } => step_of(event),
_ => None,
})
.collect();
assert_eq!(executed, vec!["e1"]);
}
#[tokio::test]
async fn let_bindings_enter_vars_and_downstream_scope() {
let dir = TempStoreDir::new("let");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
let flow = build_flow(
"let_flow",
&provider.lockfile().digest,
json!([]),
json!([]),
vec![
json!({
"kind": "let",
"stepId": "bind",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"bindings": { "label": { "fn": "concat", "args": [
{ "lit": "run-" }, { "lit": "42" }
] } }
}),
json!({
"kind": "action",
"stepId": "uses",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"effect": "mutating",
"idempotent": false,
"binding": { "attempts": [ {
"channel": "uiTree",
"actionName": "tapElement",
"args": { "value": { "ref": "vars.label" } },
"acceptExecutionModes": ["nativeSemantic", "webSemantic"],
"protection": "standard"
} ] },
"assertions": []
}),
],
json!({}),
);
let session = open(&provider).await;
let outcome = Runner::run(&flow, json!({}), session, &mut store, run_opts("run-let"))
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Finished { verdict: None }));
let events = store.events("run-let").expect("events");
let bind_entered = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::StepEntered {
step_id,
resolved_inputs,
..
} if step_id.as_str() == "bind" => Some(resolved_inputs.clone()),
_ => None,
})
.expect("bind entered");
assert_eq!(bind_entered, json!({ "label": "run-42" }));
let args = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::ActionIntent { args_snapshot, .. } => Some(args_snapshot.clone()),
_ => None,
})
.expect("intent");
assert_eq!(args, json!({ "value": "run-42" }));
}
fn element_present(assert_id: &str, identifier: &str) -> Value {
json!({
"assertId": assert_id,
"predicate": { "type": "elementState",
"selector": { "identifier": identifier },
"state": "present" },
"verifyVia": ["uiTree"],
"onMissingInput": "unknown"
})
}
#[tokio::test]
async fn assert_step_fresh_observes_and_judges() {
let dir = TempStoreDir::new("assert-fresh");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
provider.handle().inject_ui_snapshot(Some(tree(json!([
{ "stableNodeId": "n1", "role": "banner", "identifier": "welcome" }
]))));
let flow = build_flow(
"assert_fresh",
&provider.lockfile().digest,
json!([]),
json!([]),
vec![
action_step("a1", vec![]),
json!({
"kind": "assert",
"stepId": "checkFresh",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"observe": "fresh",
"assertions": [ element_present("fa", "welcome") ]
}),
],
json!({}),
);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-assert-fresh"),
)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(
event_types(&store, "run-assert-fresh"),
vec![
"runStarted",
"stepEntered", "actionIntent",
"actionSettled",
"stepExited",
"stepEntered", "observationRecorded",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"runFinished",
]
);
let events = store.events("run-assert-fresh").expect("events");
let outcome_record = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::AssertionEvaluated { outcome } => Some(outcome.clone()),
_ => None,
})
.expect("assertion outcome");
assert_eq!(outcome_record.result, VerdictStatus::Pass);
assert_eq!(outcome_record.channel, Some(pointlock_ir::Channel::UiTree));
}
#[tokio::test]
async fn assert_step_from_step_reuses_archived_material_without_observing() {
let dir = TempStoreDir::new("assert-from");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::new());
let after = provider.handle().make_observation(Some(tree(json!([
{ "stableNodeId": "n1", "role": "label", "identifier": "status",
"text": "Connected" }
]))));
provider
.handle()
.push_script(succeeded_with_after(json!({ "ok": true }), after));
let flow = build_flow(
"assert_from",
&provider.lockfile().digest,
json!([]),
json!([]),
vec![
action_step("a1", vec![element_present("aa", "status")]),
json!({
"kind": "assert",
"stepId": "checkAgain",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"observe": { "fromStep": "a1", "which": "after" },
"assertions": [ {
"assertId": "fb",
"predicate": { "type": "elementText",
"selector": { "identifier": "status" },
"match": { "value": "connected", "mode": "contains",
"caseSensitive": false } },
"verifyVia": ["uiTree"],
"onMissingInput": "unknown"
} ]
}),
],
json!({}),
);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-assert-from"),
)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
let events = store.events("run-assert-from").expect("events");
let observation_events: Vec<&str> = events
.iter()
.filter(|event| matches!(event.payload, RunLogPayload::ObservationRecorded { .. }))
.map(|event| step_of(event).unwrap_or("?"))
.collect();
assert_eq!(observation_events, vec!["a1"]);
let from_step_outcome = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::AssertionEvaluated { outcome } if outcome.assert_id.as_str() == "fb" => {
Some(outcome.clone())
}
_ => None,
})
.next()
.expect("fromStep assertion");
assert_eq!(from_step_outcome.result, VerdictStatus::Pass);
}
fn preflight_flow(digest: &Hash, probe_identifier: &str) -> FlowIR {
let mut step = action_step("probed", vec![expect_ok("pa", "probed")]);
step["preflight"] = json!([element_present("pf", probe_identifier)]);
build_flow(
"pf_flow",
digest,
json!([]),
json!([]),
vec![step],
json!({}),
)
}
#[tokio::test]
async fn preflight_pass_probes_then_acts() {
let dir = TempStoreDir::new("pf-pass");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
provider.handle().inject_ui_snapshot(Some(tree(json!([
{ "stableNodeId": "n1", "role": "page", "identifier": "loginPage" }
]))));
let flow = preflight_flow(&provider.lockfile().digest, "loginPage");
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-pf-pass"),
)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(
event_types(&store, "run-pf-pass"),
vec![
"runStarted",
"stepEntered",
"observationRecorded", "preflightProbed",
"actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"runFinished",
]
);
}
#[tokio::test]
async fn preflight_drift_blocks_and_resume_reprobes() {
let dir = TempStoreDir::new("pf-drift");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([succeeded_with(json!({ "ok": true }))]));
provider.handle().inject_ui_snapshot(Some(tree(json!([
{ "stableNodeId": "n1", "role": "page", "identifier": "homePage" }
]))));
let flow = preflight_flow(&provider.lockfile().digest, "loginPage");
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-pf-drift"),
)
.await
.expect("run");
let RunOutcome::Blocked { reason } = outcome else {
panic!("expected Blocked, got {outcome:?}");
};
assert!(reason.to_string().contains("drifted"), "got {reason}");
assert_eq!(
event_types(&store, "run-pf-drift"),
vec![
"runStarted",
"stepEntered",
"observationRecorded",
"preflightProbed",
"runSuspended",
]
);
assert!(provider.handle().dispatched_call_ids().is_empty());
provider.handle().inject_ui_snapshot(Some(tree(json!([
{ "stableNodeId": "n1", "role": "page", "identifier": "loginPage" }
]))));
let session = open(&provider).await;
let events_before = store.events("run-pf-drift").expect("events").len();
let outcome = Runner::resume(
&flow,
"run-pf-drift",
session,
&mut store,
ResumeOptions::default(),
)
.await
.expect("resume");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
let events = store.events("run-pf-drift").expect("events");
let segment: Vec<&'static str> = events[events_before..]
.iter()
.map(|event| event.payload.event_type())
.collect();
assert_eq!(
segment,
vec![
"runResumed",
"observationRecorded", "preflightProbed",
"actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"runFinished",
]
);
let entered_count = events
.iter()
.filter(|event| matches!(event.payload, RunLogPayload::StepEntered { .. }))
.count();
assert_eq!(entered_count, 1);
}
fn two_step_callee_fixture(digest: &Hash) -> (FlowIR, BTreeMap<Hash, FlowIR>) {
let callee = build_flow(
"twostep",
digest,
json!([]),
json!([]),
vec![
action_step("c1", vec![expect_ok("c1a", "c1")]),
action_step("c2", vec![expect_ok("c2a", "c2")]),
],
json!({}),
);
let root = build_flow(
"outer",
digest,
json!([]),
json!([]),
vec![
action_step("r1", vec![expect_ok("r1a", "r1")]),
json!({
"kind": "call",
"stepId": "callOnce",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "twostep", "irHash": callee.ir_hash.as_str() },
"inputs": {}
}),
action_step("r2", vec![expect_ok("r2a", "r2")]),
],
json!({ "twostep": { "flowId": "twostep", "irHash": callee.ir_hash.as_str() } }),
);
let mut registry = BTreeMap::new();
registry.insert(callee.ir_hash.clone(), callee);
(root, registry)
}
#[tokio::test]
async fn suspend_inside_callee_resumes_at_the_exact_frame_position() {
let dir = TempStoreDir::new("frame-resume");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let (root, registry) = two_step_callee_fixture(&provider.lockfile().digest);
let stop = CancellationToken::new();
let session = Box::new(StopAfter {
inner: open(&provider).await,
remaining: AtomicUsize::new(2),
stop: stop.clone(),
});
let mut opts = run_opts("run-frame");
opts.stop = stop;
opts.subflows = registry.clone();
let outcome = Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run");
assert_eq!(outcome, RunOutcome::Suspended);
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let events_before = store.events("run-frame").expect("events").len();
let view = store.verify_checkpoint("run-frame").expect("verify");
assert_eq!(view.frames.len(), 2);
assert_eq!(view.frames[1].flow_id.as_str(), "twostep");
assert_eq!(view.frames[1].next_index, 1);
let session = open(&provider).await;
let outcome = Runner::resume_with_subflows(
&root,
®istry,
"run-frame",
session,
&mut store,
ResumeOptions::default(),
)
.await
.expect("resume");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(provider.handle().dispatched_call_ids().len(), 4);
let events = store.events("run-frame").expect("events");
let segment: Vec<&'static str> = events[events_before..]
.iter()
.map(|event| event.payload.event_type())
.collect();
assert_eq!(
segment,
vec![
"runResumed",
"preflightProbed",
"stepEntered", "actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"callFramePopped",
"verdictRecorded", "stepExited",
"stepEntered", "actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"runFinished",
]
);
let pushes = events
.iter()
.filter(|event| matches!(event.payload, RunLogPayload::CallFramePushed { .. }))
.count();
assert_eq!(pushes, 1);
let c1_entries = events
.iter()
.filter(|event| {
matches!(event.payload, RunLogPayload::StepEntered { .. })
&& step_of(event) == Some("c1")
})
.count();
assert_eq!(c1_entries, 1);
let view = store.verify_checkpoint("run-frame").expect("verify");
assert_eq!(view.frames.len(), 1);
assert_eq!(view.frames[0].next_index, 3);
}
#[tokio::test]
async fn foreach_crash_mid_iteration_resumes_at_the_same_index() {
let dir = TempStoreDir::new("fe-resume");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), ScriptedOutcome::TransportLostAfterDispatch, ]));
let flow = foreach_flow(&provider.lockfile().digest, false);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-fe-crash"),
)
.await
.expect("run");
assert_eq!(outcome, RunOutcome::Suspended);
let view = store.verify_checkpoint("run-fe-crash").expect("verify");
assert!(view.frontier.pending_intent.is_some());
assert_eq!(
view.frames[0].iter_stack,
vec![IterState {
var: "item".to_owned(),
index: 1,
key: None,
}]
);
let events_before = store.events("run-fe-crash").expect("events").len();
provider
.handle()
.push_script(succeeded_with(json!({ "ok": true }))); provider
.handle()
.push_script(succeeded_with(json!({ "ok": true }))); let session = open(&provider).await;
let outcome = Runner::resume(
&flow,
"run-fe-crash",
session,
&mut store,
ResumeOptions::default(),
)
.await
.expect("resume");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Pass);
assert!(
verdict.summary.contains("3 judged step(s)"),
"got {}",
verdict.summary
);
let events = store.events("run-fe-crash").expect("events");
let segment: Vec<&'static str> = events[events_before..]
.iter()
.map(|event| event.payload.event_type())
.collect();
assert_eq!(
segment,
vec![
"runResumed",
"preflightProbed",
"actionIntent", "actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"stepEntered", "actionIntent",
"actionSettled",
"assertionEvaluated",
"verdictRecorded",
"stepExited",
"stepExited", "runFinished",
]
);
let intents: Vec<(u64, u64)> = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::ActionIntent { .. } => {
let iteration = event.run_path.iter().find_map(|frame| match frame {
PathFrame::Iteration { index, .. } => Some(*index),
_ => None,
})?;
let attempt = event.run_path.iter().find_map(|frame| match frame {
PathFrame::Attempt { n } => Some(*n),
_ => None,
})?;
Some((iteration, attempt))
}
_ => None,
})
.collect();
assert_eq!(intents, vec![(0, 1), (1, 1), (1, 2), (2, 1)]);
store.verify_checkpoint("run-fe-crash").expect("verify");
}
fn visual_assert_flow(digest: &Hash) -> FlowIR {
flow_fixture(
digest,
vec![{
let mut step = action_step("shot", vec![]);
step["assertions"] = json!([ {
"assertId": "va",
"predicate": { "type": "visual", "prompt": "the page looks right" },
"verifyVia": ["vision"],
"onMissingInput": "unknown"
} ]);
step
}],
)
}
#[tokio::test]
async fn fetch_unsupported_degrades_to_unknown_instead_of_aborting() {
let dir = TempStoreDir::new("degrade-fetch");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::new());
let after = provider.handle().make_observation(None); provider
.handle()
.push_script(succeeded_with_after(json!({}), after));
provider
.handle()
.set_fetch_evidence_unsupported(Some("control plane has no asset byte channel".into()));
let flow = visual_assert_flow(&provider.lockfile().digest);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-degrade-fetch"),
)
.await
.expect("the run must not abort on a localization failure");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Unknown);
let events = store.events("run-degrade-fetch").expect("events");
let outcome_record = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::AssertionEvaluated { outcome } => Some(outcome.clone()),
_ => None,
})
.expect("assertion outcome");
assert_eq!(outcome_record.result, VerdictStatus::Unknown);
assert!(
outcome_record.reason.contains("evidence fetch failed"),
"got {}",
outcome_record.reason
);
let view = store
.verify_checkpoint("run-degrade-fetch")
.expect("verify");
let shot = &view.completed[0];
assert_eq!(shot.observations.len(), 1);
assert!(shot.observations[0].screenshot.is_none());
}
#[tokio::test]
async fn evidence_integrity_mismatch_degrades_to_unknown() {
let dir = TempStoreDir::new("degrade-integrity");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::new());
let after = provider.handle().make_observation(None);
let asset_id = after
.screenshot
.as_ref()
.expect("screenshot asset")
.id
.clone();
provider
.handle()
.push_script(succeeded_with_after(json!({}), after));
provider
.handle()
.insert_evidence(asset_id, b"tampered bytes".to_vec());
let flow = visual_assert_flow(&provider.lockfile().digest);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-degrade-integrity"),
)
.await
.expect("the run must not abort on an integrity failure");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Unknown);
let events = store.events("run-degrade-integrity").expect("events");
let outcome_record = events
.iter()
.find_map(|event| match &event.payload {
RunLogPayload::AssertionEvaluated { outcome } => Some(outcome.clone()),
_ => None,
})
.expect("assertion outcome");
assert!(
outcome_record.reason.contains("integrity"),
"got {}",
outcome_record.reason
);
}
#[tokio::test]
async fn ui_tree_channel_still_works_when_the_screenshot_cannot_localize() {
let dir = TempStoreDir::new("degrade-tree");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::new());
let after = provider.handle().make_observation(Some(tree(json!([
{ "stableNodeId": "n1", "role": "switch", "identifier": "wifi_toggle" }
]))));
provider
.handle()
.push_script(succeeded_with_after(json!({}), after));
provider
.handle()
.set_fetch_evidence_unsupported(Some("no byte channel".into()));
let flow = flow_fixture(
&provider.lockfile().digest,
vec![action_step(
"s1",
vec![element_present("ea", "wifi_toggle")],
)],
);
let session = open(&provider).await;
let outcome = Runner::run(
&flow,
json!({}),
session,
&mut store,
run_opts("run-degrade-tree"),
)
.await
.expect("run");
let verdict = finished_verdict(outcome);
assert_eq!(verdict.status, VerdictStatus::Pass);
let view = store.verify_checkpoint("run-degrade-tree").expect("verify");
let record = &view.completed[0];
assert_eq!(record.observations.len(), 1);
assert!(record.observations[0].screenshot.is_none());
assert!(record.observations[0].ui_snapshot.is_some());
}
fn if_flow_two_step(digest: &Hash, second_target: &str) -> FlowIR {
let mut second = action_step("t2", vec![expect_ok("tb", "t2")]);
second["binding"]["attempts"][0]["args"]["element"]["lit"]["identifier"] = json!(second_target);
build_flow(
"if_flow",
digest,
json!([ { "name": "mode", "schema": { "type": "string" }, "required": true } ]),
json!([]),
vec![json!({
"kind": "if",
"stepId": "branch",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"cond": { "fn": "eq", "args": [ { "ref": "params.mode" }, { "lit": "yes" } ] },
"then": [ action_step("t1", vec![expect_ok("ta", "t1")]), second ],
"else": [ action_step("e1", vec![expect_ok("ea", "e1")]) ]
})],
json!({}),
)
}
#[tokio::test]
async fn cross_ir_resume_classifies_inside_the_taken_branch() {
let dir = TempStoreDir::new("if-cross-ir");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = if_flow_two_step(&provider.lockfile().digest, "t2_original");
let outcome = Runner::run(
&flow,
json!({ "mode": "yes" }),
open(&provider).await,
&mut store,
run_opts("run-if-x"),
)
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Finished { .. }));
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let repaired = if_flow_two_step(&provider.lockfile().digest, "t2_repaired");
let outcome = Runner::resume(
&repaired,
"run-if-x",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["t2".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("the branch-internal repair resumes");
assert!(matches!(outcome, RunOutcome::Finished { .. }));
assert_eq!(provider.handle().dispatched_call_ids().len(), 3);
}
#[tokio::test]
async fn cross_ir_resume_ignores_a_change_in_the_untaken_branch() {
let dir = TempStoreDir::new("if-untaken");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = if_flow(&provider.lockfile().digest);
Runner::run(
&flow,
json!({ "mode": "yes" }),
open(&provider).await,
&mut store,
run_opts("run-if-u"),
)
.await
.expect("run");
assert_eq!(provider.handle().dispatched_call_ids().len(), 1);
let mut edited = if_flow(&provider.lockfile().digest);
let pointlock_ir::StepIR::If(branch) = &mut edited.body[0] else {
panic!("if step");
};
let pointlock_ir::StepIR::Action(otherwise) =
&mut branch.r#else.as_mut().expect("else branch")[0]
else {
panic!("action step");
};
otherwise.base.effect_hash = Hash::new(h64('9')).expect("hash");
seal(&mut edited);
Runner::resume(
&edited,
"run-if-u",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("an untaken-branch edit does not invalidate the taken branch");
assert_eq!(provider.handle().dispatched_call_ids().len(), 1);
}
fn foreach_flow_with_tail(digest: &Hash, tail_target: &str) -> FlowIR {
let mut tail = action_step("tail", vec![]);
tail["binding"]["attempts"][0]["args"]["element"]["lit"]["identifier"] = json!(tail_target);
foreach_flow_parts(digest, json!(["a", "b", "c"]), Some(tail))
}
fn foreach_flow_items(digest: &Hash, items: Value) -> FlowIR {
foreach_flow_parts(digest, items, None)
}
fn foreach_flow_parts(digest: &Hash, items: Value, tail: Option<Value>) -> FlowIR {
let mut body = vec![json!({
"kind": "foreach",
"stepId": "eachItem",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"items": { "lit": items },
"as": "item",
"body": [ {
"kind": "action",
"stepId": "perItem",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"effect": "mutating",
"idempotent": true,
"binding": { "attempts": [ {
"channel": "uiTree",
"actionName": "tapElement",
"args": { "value": { "ref": "iter.item" } },
"acceptExecutionModes": ["nativeSemantic", "webSemantic"],
"protection": "standard"
} ] },
"assertions": [ expect_ok("pa", "perItem") ]
} ]
})];
if let Some(tail) = tail {
body.push(tail);
}
build_flow("fe_flow", digest, json!([]), json!([]), body, json!({}))
}
#[tokio::test]
async fn cross_ir_resume_adopts_completed_foreach_rounds() {
let dir = TempStoreDir::new("fe-cross-ir");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = foreach_flow(&provider.lockfile().digest, true);
Runner::run(
&flow,
json!({}),
open(&provider).await,
&mut store,
run_opts("run-fe-x"),
)
.await
.expect("run");
assert_eq!(provider.handle().dispatched_call_ids().len(), 4);
let repaired = foreach_flow_with_tail(&provider.lockfile().digest, "tail_repaired");
Runner::resume(
&repaired,
"run-fe-x",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["tail".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("the rounds are adopted; only the tail re-runs");
assert_eq!(
provider.handle().dispatched_call_ids().len(),
5,
"exactly one more dispatch: the three rounds were adopted"
);
}
#[tokio::test]
async fn cross_ir_resume_invalidates_a_foreach_whose_head_changed() {
let dir = TempStoreDir::new("fe-head");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = foreach_flow(&provider.lockfile().digest, false);
Runner::run(
&flow,
json!({}),
open(&provider).await,
&mut store,
run_opts("run-fe-h"),
)
.await
.expect("run");
assert_eq!(provider.handle().dispatched_call_ids().len(), 3);
let repaired = foreach_flow_items(&provider.lockfile().digest, json!(["a", "b"]));
Runner::resume(
&repaired,
"run-fe-h",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("a changed head re-runs the whole foreach");
assert_eq!(provider.handle().dispatched_call_ids().len(), 5);
}
fn assert_flow(digest: &Hash, matcher: &str) -> FlowIR {
build_flow(
"assert_x",
digest,
json!([]),
json!([]),
vec![
action_step("a1", vec![]),
json!({
"kind": "assert",
"stepId": "checkFresh",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"observe": "fresh",
"assertions": [ element_present("fa", matcher) ]
}),
],
json!({}),
)
}
#[tokio::test]
async fn a_changed_assertion_on_an_assert_step_re_observes() {
let dir = TempStoreDir::new("assert-cross-ir");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = assert_flow(&provider.lockfile().digest, "welcome");
Runner::run(
&flow,
json!({}),
open(&provider).await,
&mut store,
run_opts("run-as-x"),
)
.await
.expect("run");
let repaired = assert_flow(&provider.lockfile().digest, "farewell");
Runner::resume(
&repaired,
"run-as-x",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("the assert step re-executes");
let report = store
.events("run-as-x")
.expect("events")
.iter()
.rev()
.find_map(|event| match &event.payload {
pointlock_ir::RunLogPayload::RunResumed {
alignment_report, ..
} => Some(alignment_report.clone()),
_ => None,
})
.expect("the resume recorded its alignment report");
let entry = report
.entries
.iter()
.find(|entry| entry.step_id.as_str() == "checkFresh")
.expect("the assert step is classified");
assert_eq!(entry.class, pointlock_ir::AlignmentClass::EffectDirty);
assert_ne!(
entry.reason.as_deref(),
Some("preflightChanged"),
"a changed assertion is not a preflight-only change"
);
}
fn call_pair(
digest: &Hash,
callee_target: &str,
tail_target: &str,
) -> (FlowIR, BTreeMap<Hash, FlowIR>) {
let mut inner_step = action_step("g1", vec![expect_ok("ga", "g1")]);
inner_step["binding"]["attempts"][0]["args"]["element"]["lit"]["identifier"] =
json!(callee_target);
let inner = build_flow(
"inner",
digest,
json!([]),
json!([]),
vec![inner_step],
json!({}),
);
let mut tail = action_step("after", vec![expect_ok("aa", "after")]);
tail["binding"]["attempts"][0]["args"]["element"]["lit"]["identifier"] = json!(tail_target);
let root = build_flow(
"caller",
digest,
json!([]),
json!([]),
vec![
json!({
"kind": "call",
"stepId": "callInner",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "inner", "irHash": inner.ir_hash.as_str() },
"inputs": {}
}),
tail,
],
json!({ "inner": { "flowId": "inner", "irHash": inner.ir_hash.as_str() } }),
);
let registry = BTreeMap::from([(inner.ir_hash.clone(), inner)]);
(root, registry)
}
#[tokio::test]
async fn cross_ir_resume_adopts_an_untouched_call_without_re_entering_the_callee() {
let dir = TempStoreDir::new("call-adopt");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let (root, registry) = call_pair(&provider.lockfile().digest, "g1_target", "after_v1");
let mut opts = run_opts("run-call-a");
opts.subflows = registry.clone();
Runner::run(&root, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let (repaired, registry) = call_pair(&provider.lockfile().digest, "g1_target", "after_v2");
Runner::resume_with_subflows(
&repaired,
®istry,
"run-call-a",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["after".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("the call adopts; only the tail re-runs");
assert_eq!(
provider.handle().dispatched_call_ids().len(),
3,
"the callee must not be re-entered"
);
}
#[tokio::test]
async fn cross_ir_resume_re_calls_a_callee_that_changed() {
let dir = TempStoreDir::new("call-dirty");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let (root, registry) = call_pair(&provider.lockfile().digest, "g1_v1", "after_v1");
let mut opts = run_opts("run-call-d");
opts.subflows = registry.clone();
Runner::run(&root, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let (repaired, registry) = call_pair(&provider.lockfile().digest, "g1_v2", "after_v1");
Runner::resume_with_subflows(
&repaired,
®istry,
"run-call-d",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["callInner".to_owned(), "after".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("the changed call is re-called");
assert_eq!(provider.handle().dispatched_call_ids().len(), 4);
}
fn drilldown_pair(digest: &Hash, c2_target: &str, label: &str) -> (FlowIR, BTreeMap<Hash, FlowIR>) {
let mut c1 = action_step("c1", vec![expect_ok("c1a", "c1")]);
c1["binding"]["attempts"][0]["args"]["element"] = json!({ "ref": "params.label" });
let mut c2 = action_step("c2", vec![expect_ok("c2a", "c2")]);
c2["binding"]["attempts"][0]["args"]["element"]["lit"]["identifier"] = json!(c2_target);
let callee = build_flow(
"twostep",
digest,
json!([ { "name": "label", "schema": {}, "required": true } ]),
json!([]),
vec![c1, c2],
json!({}),
);
let root = build_flow(
"outer",
digest,
json!([]),
json!([]),
vec![
action_step("r1", vec![expect_ok("r1a", "r1")]),
json!({
"kind": "call",
"stepId": "callOnce",
"effectHash": h64('0'),
"judgeHash": h64('0'),
"checkpoint": true,
"flowRef": { "flowId": "twostep", "irHash": callee.ir_hash.as_str() },
"inputs": { "label": { "lit": { "identifier": label } } }
}),
action_step("r2", vec![expect_ok("r2a", "r2")]),
],
json!({ "twostep": { "flowId": "twostep", "irHash": callee.ir_hash.as_str() } }),
);
let registry = BTreeMap::from([(callee.ir_hash.clone(), callee)]);
(root, registry)
}
fn last_alignment(store: &Store, run_id: &str) -> pointlock_ir::AlignmentReport {
store
.events(run_id)
.expect("events")
.iter()
.rev()
.find_map(|event| match &event.payload {
RunLogPayload::RunResumed {
alignment_report, ..
} => Some(alignment_report.clone()),
_ => None,
})
.expect("the resume recorded its alignment report")
}
#[tokio::test]
async fn down_drill_adopts_callee_steps_before_the_repair() {
let dir = TempStoreDir::new("drill-ok");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let digest = provider.lockfile().digest.clone();
let (root, registry) = drilldown_pair(&digest, "c2_v1", "same");
let stop = CancellationToken::new();
let session = Box::new(StopAfter {
inner: open(&provider).await,
remaining: AtomicUsize::new(2), stop: stop.clone(),
});
let mut opts = run_opts("run-drill");
opts.stop = stop;
opts.subflows = registry.clone();
assert_eq!(
Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run"),
RunOutcome::Suspended
);
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let view = store.verify_checkpoint("run-drill").expect("verify");
assert_eq!(view.frames.len(), 2);
let old_callee_pin = view.frames[1].ir_hash.clone();
assert!(
!view
.completed
.iter()
.any(|record| record.step_id.as_str() == "callOnce"),
"an unfinished call contributes no StepRecord"
);
let (repaired, new_registry) = drilldown_pair(&digest, "c2_v2", "same");
assert_ne!(repaired.ir_hash, root.ir_hash);
Runner::resume_with_subflows(
&repaired,
&new_registry,
"run-drill",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("the down-drill resumes inside the callee");
assert_eq!(
provider.handle().dispatched_call_ids().len(),
4,
"the concluded callee step must not run again"
);
let report = last_alignment(&store, "run-drill");
let call_entry = report
.entries
.iter()
.find(|entry| entry.step_id.as_str() == "callOnce")
.expect("the call step is classified");
assert_eq!(call_entry.class, pointlock_ir::AlignmentClass::EffectDirty);
assert!(
call_entry
.reason
.as_deref()
.is_some_and(|reason| reason.contains("down-drill")),
"the report must name the down-drill, got {:?}",
call_entry.reason
);
assert_eq!(
report
.resume_point
.as_deref()
.map(pointlock_ir::render_run_path),
Some(pointlock_ir::render_run_path(
&report
.entries
.iter()
.find(|entry| entry.step_id.as_str() == "c2")
.expect("c2 is classified")
.run_path
))
);
let c1_entry = report
.entries
.iter()
.find(|entry| entry.step_id.as_str() == "c1")
.expect("the callee step is classified");
assert_eq!(c1_entry.class, pointlock_ir::AlignmentClass::Reusable);
let events = store.events("run-drill").expect("events");
let rebases: Vec<bool> = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::CallFramePushed { rebase, .. } => Some(*rebase),
_ => None,
})
.collect();
assert_eq!(rebases, vec![false, true], "one open, one re-entry");
let mid = store
.rebuild_checkpoint("run-drill")
.expect("rebuild")
.frames
.len();
assert_eq!(mid, 1, "the frame stack is balanced again after the pop");
let new_pin = new_registry
.keys()
.next()
.expect("one callee in the registry")
.clone();
assert_ne!(new_pin, old_callee_pin);
let pins: Vec<_> = events
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::CallFramePushed {
frame,
rebase: true,
} => Some(frame.ir_hash.clone()),
_ => None,
})
.collect();
assert_eq!(pins, vec![new_pin]);
}
#[tokio::test]
async fn changed_arguments_tear_down_the_live_frame_and_re_call() {
let dir = TempStoreDir::new("drill-args");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let digest = provider.lockfile().digest.clone();
let (root, registry) = drilldown_pair(&digest, "c2_v1", "before");
let (repaired, new_registry) = drilldown_pair(&digest, "c2_v1", "after");
let old_callee = registry.values().next().expect("callee");
let new_callee = new_registry.values().next().expect("callee");
assert_eq!(old_callee.ir_hash, new_callee.ir_hash);
assert_eq!(
old_callee.body[0].base().effect_hash,
new_callee.body[0].base().effect_hash
);
assert_ne!(repaired.ir_hash, root.ir_hash);
let stop = CancellationToken::new();
let session = Box::new(StopAfter {
inner: open(&provider).await,
remaining: AtomicUsize::new(2),
stop: stop.clone(),
});
let mut opts = run_opts("run-args");
opts.stop = stop;
opts.subflows = registry.clone();
assert_eq!(
Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run"),
RunOutcome::Suspended
);
let error = Runner::resume_with_subflows(
&repaired,
&new_registry,
"run-args",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect_err("re-calling over an effective mutating step needs authorization");
let pointlock_runner::RunnerError::RequiresConfirmation { report } = error else {
panic!("expected the 07 §5.4 gate, got {error}");
};
assert_eq!(report.requires_confirmation.len(), 1, "{report:?}");
assert_eq!(
report.requires_confirmation[0]
.step_id
.as_ref()
.map(|id| id.as_str()),
Some("callOnce")
);
assert_eq!(
provider.handle().dispatched_call_ids().len(),
2,
"pre-execution refusal"
);
let outcome = Runner::resume_with_subflows(
&repaired,
&new_registry,
"run-args",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["callOnce".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("the authorized teardown re-call proceeds");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(provider.handle().dispatched_call_ids().len(), 5);
let view = store.verify_checkpoint("run-args").expect("verify");
assert_eq!(view.frames.len(), 1, "no stranded frame survives");
let events = store.events("run-args").expect("events");
let pushes = events
.iter()
.filter(|event| matches!(event.payload, RunLogPayload::CallFramePushed { .. }))
.count();
let pops = events
.iter()
.filter(|event| matches!(event.payload, RunLogPayload::CallFramePopped { .. }))
.count();
assert_eq!(
(pushes, pops),
(2, 2),
"old frame + torn down, new frame + popped"
);
let last_push = events
.iter()
.rev()
.find_map(|event| match &event.payload {
RunLogPayload::CallFramePushed { frame, rebase } => Some((frame.clone(), *rebase)),
_ => None,
})
.expect("the re-call pushed a frame");
assert!(
!last_push.1,
"a teardown re-call is a fresh push, not a rebase"
);
assert_eq!(
last_push.0.inputs_snapshot["label"]["identifier"],
json!("after"),
"inputs were re-evaluated under the new IR: {:?}",
last_push.0.inputs_snapshot
);
let call_records: Vec<_> = view
.completed
.iter()
.filter(|record| record.step_id.as_str() == "callOnce")
.collect();
assert_eq!(call_records.len(), 2, "aborted teardown + judged re-call");
assert!(
call_records[0].verdict.is_none(),
"the abort claims nothing"
);
assert!(call_records[1].verdict.is_some(), "the re-call concluded");
}
#[tokio::test]
async fn a_readonly_frame_tears_down_without_authorization() {
let dir = TempStoreDir::new("drill-args-ro");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let digest = provider.lockfile().digest.clone();
let readonly_pair = |label: &str| {
let (mut root, mut registry) = drilldown_pair(&digest, "c2_v1", label);
let mut callee = registry.values().next().expect("callee").clone();
for step in &mut callee.body {
if let pointlock_ir::StepIR::Action(action) = step {
action.effect = pointlock_ir::vocab::EffectClassAction::Readonly;
action.idempotent = true;
}
}
seal(&mut callee);
let pointlock_ir::StepIR::Call(call) = &mut root.body[1] else {
panic!("body[1] is the call");
};
call.flow_ref.ir_hash = callee.ir_hash.clone();
root.subflows.insert(
serde_json::from_value(json!("twostep")).expect("flow id"),
serde_json::from_value(json!({
"flowId": "twostep", "irHash": callee.ir_hash.as_str()
}))
.expect("flow ref"),
);
seal(&mut root);
registry.clear();
registry.insert(callee.ir_hash.clone(), callee);
(root, registry)
};
let (root, registry) = readonly_pair("before");
let (repaired, new_registry) = readonly_pair("after");
assert_ne!(repaired.ir_hash, root.ir_hash);
let stop = CancellationToken::new();
let session = Box::new(StopAfter {
inner: open(&provider).await,
remaining: AtomicUsize::new(2),
stop: stop.clone(),
});
let mut opts = run_opts("run-args-ro");
opts.stop = stop;
opts.subflows = registry.clone();
Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run");
let outcome = Runner::resume_with_subflows(
&repaired,
&new_registry,
"run-args-ro",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("a readonly frame needs no authorization to re-call");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(provider.handle().dispatched_call_ids().len(), 5);
let view = store.verify_checkpoint("run-args-ro").expect("verify");
assert_eq!(view.frames.len(), 1);
}
#[tokio::test]
async fn re_calling_a_failed_frame_gates_on_its_effective_callee_step() {
let dir = TempStoreDir::new("frame-gate");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": false })), ]));
let digest = provider.lockfile().digest.clone();
let (root, registry) = drilldown_pair(&digest, "c2_v1", "same");
let mut opts = run_opts("run-gate");
opts.subflows = registry.clone();
let outcome = Runner::run(&root, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Fail);
let view = store.verify_checkpoint("run-gate").expect("verify");
assert_eq!(view.frames.len(), 1, "the callee frame was popped");
let call_record = view
.completed
.iter()
.find(|record| record.step_id.as_str() == "callOnce")
.expect("the concluded call has a record");
assert_eq!(
call_record
.verdict
.as_ref()
.expect("aggregate verdict")
.status,
VerdictStatus::Fail
);
let (repaired, new_registry) = drilldown_pair(&digest, "c2_v2", "same");
let error = Runner::resume_with_subflows(
&repaired,
&new_registry,
"run-gate",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect_err("re-calling the frame replays c1 and must be confirmed");
let pointlock_runner::RunnerError::RequiresConfirmation { report } = error else {
panic!("expected the 07 §5.4 gate, got {error}");
};
let gated: Vec<String> = report
.requires_confirmation
.iter()
.map(|entry| pointlock_ir::render_run_path(&entry.run_path))
.collect();
assert!(
report
.requires_confirmation
.iter()
.any(|entry| pointlock_ir::render_run_path(&entry.run_path).contains("callOnce")),
"the call step must be gated for its frame's effective step, got {gated:?}"
);
}
#[tokio::test]
async fn a_changed_cond_gates_the_branch_it_replays() {
let dir = TempStoreDir::new("cond-gate");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })),
succeeded_with(json!({ "ok": true })),
]));
let flow = if_flow_two_step(&provider.lockfile().digest, "t2_original");
Runner::run(
&flow,
json!({ "mode": "yes" }),
open(&provider).await,
&mut store,
run_opts("run-cond"),
)
.await
.expect("run");
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
let mut edited = if_flow_two_step(&provider.lockfile().digest, "t2_original");
let pointlock_ir::StepIR::If(branch) = &mut edited.body[0] else {
panic!("body[0] is the if");
};
branch.cond = serde_json::from_value(json!({
"fn": "eq", "args": [ { "lit": "yes" }, { "ref": "params.mode" } ]
}))
.expect("cond");
seal(&mut edited);
let error = Runner::resume(
&edited,
"run-cond",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect_err("replaying the branch needs confirmation");
let pointlock_runner::RunnerError::RequiresConfirmation { report } = error else {
panic!("expected the 07 §5.4 gate, got {error}");
};
let gated: Vec<&str> = report
.requires_confirmation
.iter()
.map(|entry| entry.cause.as_str())
.collect();
assert_eq!(gated, vec!["mutatingReexec"]);
assert_eq!(
pointlock_ir::render_run_path(&report.requires_confirmation[0].run_path),
pointlock_ir::render_run_path(
&report
.entries
.iter()
.find(|entry| entry.step_id.as_str() == "branch")
.expect("the if is classified")
.run_path
),
"the gate names the container, which is what --allow-mutating-reexec takes"
);
Runner::resume(
&edited,
"run-cond",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["branch".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("authorized replay");
assert_eq!(provider.handle().dispatched_call_ids().len(), 4);
}
fn probe_count(store: &Store, run_id: &str) -> Vec<usize> {
store
.events(run_id)
.expect("events")
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::PreflightProbed { outcomes } => Some(outcomes.len()),
_ => None,
})
.collect()
}
#[tokio::test]
async fn a_resume_without_probes_records_that_it_checked_nothing() {
let dir = TempStoreDir::new("unprobed");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let digest = provider.lockfile().digest.clone();
let (root, registry) = drilldown_pair(&digest, "c2_v1", "same");
let stop = CancellationToken::new();
let session = Box::new(StopAfter {
inner: open(&provider).await,
remaining: AtomicUsize::new(2),
stop: stop.clone(),
});
let mut opts = run_opts("run-unprobed");
opts.stop = stop;
opts.subflows = registry.clone();
Runner::run(&root, json!({}), session, &mut store, opts)
.await
.expect("run");
assert_eq!(
probe_count(&store, "run-unprobed"),
Vec::<usize>::new(),
"a fresh run has no unwatched world to re-touch"
);
let before = store.events("run-unprobed").expect("events").len();
Runner::resume_with_subflows(
&root,
®istry,
"run-unprobed",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("resume");
let segment: Vec<usize> = store.events("run-unprobed").expect("events")[before..]
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::PreflightProbed { outcomes } => Some(outcomes.len()),
_ => None,
})
.collect();
assert_eq!(
segment,
vec![0],
"one unprobed mark at the re-entry step, and nothing after it"
);
}
fn trivial_preflight() -> Value {
json!([{
"assertId": "pf",
"predicate": { "type": "expr", "expr": { "fn": "eq", "args": [
{ "lit": 1 }, { "lit": 1 }
] } },
"verifyVia": [],
"onMissingInput": "unknown"
}])
}
#[tokio::test]
async fn a_gate_released_step_is_its_own_unprobed_mark() {
let dir = TempStoreDir::new("unprobed-gate");
let mut store = Store::open(dir.path()).expect("open store");
let provider = FakeProvider::new(VecDeque::from([
succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), succeeded_with(json!({ "ok": true })), ]));
let digest = provider.lockfile().digest.clone();
let body = |s2_target: &str| {
let mut s2 = action_step("s2", vec![expect_ok("a2", "s2")]);
s2["binding"]["attempts"][0]["args"]["element"]["lit"]["identifier"] = json!(s2_target);
s2["preflight"] = trivial_preflight();
vec![
action_step("s1", vec![expect_ok("a1", "s1")]),
s2,
action_step("s3", vec![expect_ok("a3", "s3")]),
]
};
let flow = flow_fixture(&digest, body("s2_v1"));
Runner::run(
&flow,
json!({}),
open(&provider).await,
&mut store,
run_opts("run-gate-probe"),
)
.await
.expect("run");
let before = store.events("run-gate-probe").expect("events").len();
let repaired = flow_fixture(&digest, body("s2_v2"));
Runner::resume(
&repaired,
"run-gate-probe",
open(&provider).await,
&mut store,
ResumeOptions {
allow_mutating_reexec: vec!["s2".to_owned(), "s3".to_owned()],
..ResumeOptions::default()
},
)
.await
.expect("authorized replay");
let segment: Vec<usize> = store.events("run-gate-probe").expect("events")[before..]
.iter()
.filter_map(|event| match &event.payload {
RunLogPayload::PreflightProbed { outcomes } => Some(outcomes.len()),
_ => None,
})
.collect();
assert_eq!(
segment,
vec![1, 0],
"s2 probed for real (1 outcome); s3 was released by name with no probes \
declared, so it is an unprobed mark of its own"
);
}
fn last_suspension_reason(store: &Store, run_id: &str) -> String {
let events = store.events(run_id).expect("events");
let last = events.last().expect("at least one event");
match &last.payload {
RunLogPayload::RunSuspended { reason, .. } => reason.clone().expect("reason present"),
other => panic!("ledger must end with runSuspended, got {other:?}"),
}
}
fn three_steps(digest: &Hash) -> FlowIR {
flow_fixture(
digest,
vec![
action_step("s1", vec![expect_ok("a1", "s1")]),
action_step("s2", vec![expect_ok("a2", "s2")]),
action_step("s3", vec![expect_ok("a3", "s3")]),
],
)
}
fn ok_provider(n: usize) -> FakeProvider {
FakeProvider::new(VecDeque::from(
std::iter::repeat_with(|| succeeded_with(json!({ "ok": true })))
.take(n)
.collect::<Vec<_>>(),
))
}
fn entered_steps(store: &Store, run_id: &str) -> Vec<String> {
store
.events(run_id)
.expect("events")
.iter()
.filter(|event| matches!(event.payload, RunLogPayload::StepEntered { .. }))
.filter_map(|event| step_of(event).map(str::to_owned))
.collect()
}
#[tokio::test]
async fn stop_at_suspends_before_the_target_enters_and_a_plain_resume_runs_it() {
let dir = TempStoreDir::new("bp-at");
let mut store = Store::open(dir.path()).expect("open store");
let provider = ok_provider(3);
let flow = three_steps(&provider.lockfile().digest);
let mut opts = run_opts("run-bp-at");
opts.stop_at = Some("s2".to_owned());
let outcome = Runner::run(&flow, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(
provider.handle().dispatched_call_ids().len(),
1,
"only s1 ran"
);
assert_eq!(
entered_steps(&store, "run-bp-at"),
vec!["s1"],
"s2 never entered"
);
let reason = last_suspension_reason(&store, "run-bp-at");
assert!(
reason.starts_with("stopped at breakpoint --stop-at ") && reason.ends_with("/s2"),
"honest reason names the matched instance: {reason}"
);
let outcome = Runner::resume(
&flow,
"run-bp-at",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("resume");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(provider.handle().dispatched_call_ids().len(), 3);
assert_eq!(entered_steps(&store, "run-bp-at"), vec!["s1", "s2", "s3"]);
}
#[tokio::test]
async fn stop_after_suspends_after_the_target_exits_before_the_next_enters() {
let dir = TempStoreDir::new("bp-after");
let mut store = Store::open(dir.path()).expect("open store");
let provider = ok_provider(3);
let flow = three_steps(&provider.lockfile().digest);
let mut opts = run_opts("run-bp-after");
opts.stop_after = Some("s2".to_owned());
let outcome = Runner::run(&flow, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(provider.handle().dispatched_call_ids().len(), 2);
assert_eq!(entered_steps(&store, "run-bp-after"), vec!["s1", "s2"]);
let events = store.events("run-bp-after").expect("events");
let s2_exited = events.iter().rposition(|event| {
matches!(event.payload, RunLogPayload::StepExited { .. }) && step_of(event) == Some("s2")
});
assert_eq!(
s2_exited,
Some(events.len() - 2),
"s2's exit is immediately followed by the suspension"
);
let reason = last_suspension_reason(&store, "run-bp-after");
assert!(
reason.starts_with("stopped at breakpoint --stop-after ") && reason.ends_with("/s2"),
"{reason}"
);
}
#[tokio::test]
async fn iteration_qualified_breakpoint_stops_only_at_that_iteration() {
let dir = TempStoreDir::new("bp-iter");
let mut store = Store::open(dir.path()).expect("open store");
let provider = ok_provider(4);
let flow = foreach_flow(&provider.lockfile().digest, true);
let mut opts = run_opts("run-bp-iter");
opts.stop_at = Some("eachItem[1]/perItem".to_owned());
let outcome = Runner::run(&flow, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(
provider.handle().dispatched_call_ids().len(),
1,
"iteration 0 ran; iteration 1 never entered"
);
let reason = last_suspension_reason(&store, "run-bp-iter");
assert!(
reason.ends_with("/eachItem[1]/perItem"),
"the matched instance is iteration 1: {reason}"
);
let outcome = Runner::resume(
&flow,
"run-bp-iter",
open(&provider).await,
&mut store,
ResumeOptions::default(),
)
.await
.expect("resume");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(provider.handle().dispatched_call_ids().len(), 4);
}
#[tokio::test]
async fn stop_at_does_not_fire_on_a_container_span_reopened_by_resume() {
let dir = TempStoreDir::new("bp-open-span");
let mut store = Store::open(dir.path()).expect("open store");
let provider = ok_provider(4);
let flow = foreach_flow(&provider.lockfile().digest, true);
let mut opts = run_opts("run-bp-open-span");
opts.stop_after = Some("eachItem[0]/perItem".to_owned());
let outcome = Runner::run(&flow, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(
entered_steps(&store, "run-bp-open-span"),
vec!["eachItem", "perItem"]
);
let outcome = Runner::resume(
&flow,
"run-bp-open-span",
open(&provider).await,
&mut store,
ResumeOptions {
stop_at: Some("eachItem".to_owned()),
stop_after: Some("eachItem[2]/perItem".to_owned()),
..ResumeOptions::default()
},
)
.await
.expect("resume");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(
provider.handle().dispatched_call_ids().len(),
3,
"iterations 1 and 2 ran: the open container was not re-suspended"
);
let reason = last_suspension_reason(&store, "run-bp-open-span");
assert!(
reason.starts_with("stopped at breakpoint --stop-after ")
&& reason.ends_with("/eachItem[2]/perItem"),
"the container breakpoint never fired; the inner one did: {reason}"
);
assert_eq!(
entered_steps(&store, "run-bp-open-span"),
vec!["eachItem", "perItem", "perItem", "perItem"],
"the container entered exactly once across both segments"
);
let outcome = Runner::resume(
&flow,
"run-bp-open-span",
open(&provider).await,
&mut store,
ResumeOptions {
stop_at: Some("tail".to_owned()),
..ResumeOptions::default()
},
)
.await
.expect("resume");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(provider.handle().dispatched_call_ids().len(), 3);
let reason = last_suspension_reason(&store, "run-bp-open-span");
assert!(
reason.starts_with("stopped at breakpoint --stop-at ") && reason.ends_with("/tail"),
"{reason}"
);
}
#[tokio::test]
async fn bare_step_id_breakpoint_hits_the_first_instance() {
let dir = TempStoreDir::new("bp-bare");
let mut store = Store::open(dir.path()).expect("open store");
let provider = ok_provider(4);
let flow = foreach_flow(&provider.lockfile().digest, true);
let mut opts = run_opts("run-bp-bare");
opts.stop_after = Some("perItem".to_owned());
let outcome = Runner::run(&flow, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert!(matches!(outcome, RunOutcome::Suspended), "got {outcome:?}");
assert_eq!(provider.handle().dispatched_call_ids().len(), 1);
let reason = last_suspension_reason(&store, "run-bp-bare");
assert!(
reason.ends_with("/eachItem[0]/perItem"),
"the shorthand resolves to iteration 0: {reason}"
);
}
#[tokio::test]
async fn an_unmatched_breakpoint_target_lets_the_run_complete() {
let dir = TempStoreDir::new("bp-none");
let mut store = Store::open(dir.path()).expect("open store");
let provider = ok_provider(3);
let flow = three_steps(&provider.lockfile().digest);
let mut opts = run_opts("run-bp-none");
opts.stop_at = Some("no_such_step".to_owned());
opts.stop_after = Some("s1/nested[0]".to_owned());
let outcome = Runner::run(&flow, json!({}), open(&provider).await, &mut store, opts)
.await
.expect("run");
assert_eq!(finished_verdict(outcome).status, VerdictStatus::Pass);
assert_eq!(provider.handle().dispatched_call_ids().len(), 3);
}