use crate::browser::session::{
BrowserResult, BrowserSession, WorkflowCheckpoint, WorkflowDefinition, WorkflowRunResult,
WorkflowRunStatus, compile_workflow_json, compile_workflow_yaml,
};
use crate::reliability::{
ReliabilityExecutionOperation, ReliabilityFaultKind, ReliabilityFixtureManifest,
ReliabilityForbiddenOutcome, ReliabilityPlatform, ReliabilityReplayBundle,
ReliabilityReplayEvent, ReliabilityRunClassification, ReliabilityRunMetadata,
ReliabilityScenario, ReliabilityScenarioObservation,
};
use serde_json::Value;
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::time::Instant;
#[derive(Debug, Clone)]
pub struct ReliabilityRunOptions {
pub workflow_root: PathBuf,
pub inputs: BTreeMap<String, Value>,
}
#[derive(Debug, Clone)]
pub struct ReliabilityRunEvidence {
pub observation: ReliabilityScenarioObservation,
pub replay: ReliabilityReplayBundle,
}
pub async fn run_reliability_scenario(
session: &BrowserSession,
scenario: &ReliabilityScenario,
fixture: &ReliabilityFixtureManifest,
options: &ReliabilityRunOptions,
) -> BrowserResult<ReliabilityRunEvidence> {
let plan = scenario.execution_plan(fixture)?;
let platform = current_platform()?;
let browser_version = browser_version(session).await?;
let started = Instant::now();
let mut events = Vec::new();
let mut action_count = 0u32;
let mut workflow: Option<WorkflowDefinition> = None;
let mut checkpoint: Option<WorkflowCheckpoint> = None;
let mut last_result: Option<WorkflowRunResult> = None;
let mut had_failure = false;
let mut unsupported = false;
for operation in &plan.operations {
if started.elapsed().as_millis() as u64 >= scenario.budgets.max_duration_ms {
had_failure = true;
events.push(event(events.len() as u32, "budget", "exhausted"));
break;
}
if action_count >= scenario.budgets.max_browser_actions {
had_failure = true;
events.push(event(events.len() as u32, "budget", "actions_exhausted"));
break;
}
action_count = action_count.saturating_add(1);
match operation {
ReliabilityExecutionOperation::ApplyControl { control } => {
let expression = format!(
"window.reliabilityLab.{}(); true",
control.javascript_method()
);
match session.evaluate(&expression).await {
Ok(_) => events.push(event(events.len() as u32, "applyControl", "applied")),
Err(_) => {
had_failure = true;
events.push(event(events.len() as u32, "applyControl", "failed"));
break;
}
}
}
ReliabilityExecutionOperation::InjectFault { injection } => {
if matches!(
injection.fault,
ReliabilityFaultKind::RendererDisconnect
| ReliabilityFaultKind::BrowserDisconnect
) {
let dispatched = match session.raw_cdp() {
Ok(cdp) => match injection.fault {
ReliabilityFaultKind::RendererDisconnect => {
cdp.send("Page.crash", None).await.is_ok()
}
ReliabilityFaultKind::BrowserDisconnect => {
cdp.send_browser("Browser.close", None).await.is_ok()
}
_ => unreachable!("transport fault branch is exhaustive"),
},
Err(_) => false,
};
events.push(event(
events.len() as u32,
"injectFault",
if dispatched {
"dispatched"
} else {
"dispatch_failed"
},
));
unsupported = true;
break;
}
let fault = serde_json::to_string(injection.fault.fixture_name())?;
let expression = format!("window.reliabilityLab.injectFault({fault}); true");
match session.evaluate(&expression).await {
Ok(_) => events.push(event(events.len() as u32, "injectFault", "applied")),
Err(_) => {
had_failure = true;
events.push(event(events.len() as u32, "injectFault", "failed"));
break;
}
}
}
ReliabilityExecutionOperation::RunWorkflow { source } => {
let path = bounded_workflow_path(&options.workflow_root, source)?;
let definition = load_workflow(&path)?;
let result = match session.run_workflow(&definition, &options.inputs).await {
Ok(result) => result,
Err(_) => {
had_failure = true;
events.push(event(events.len() as u32, "runWorkflow", "failed"));
break;
}
};
action_count = action_count
.saturating_add(result.trace.events.len().min(u32::MAX as usize) as u32);
had_failure |= result.status != WorkflowRunStatus::Completed;
checkpoint = session
.export_workflow_checkpoint(&definition, &result)
.await
.ok();
workflow = Some(definition);
events.push(event(
events.len() as u32,
"runWorkflow",
workflow_status(result.status),
));
last_result = Some(result);
}
ReliabilityExecutionOperation::ResumeFromCheckpoint { checkpoint: name } => {
if name != "latest" {
unsupported = true;
events.push(event(
events.len() as u32,
"resume",
"unsupported_checkpoint",
));
break;
}
let (Some(definition), Some(saved_checkpoint)) = (&workflow, &checkpoint) else {
unsupported = true;
events.push(event(events.len() as u32, "resume", "missing_checkpoint"));
break;
};
let result = match session
.resume_workflow(definition, &options.inputs, saved_checkpoint)
.await
{
Ok(result) => result,
Err(_) => {
had_failure = true;
events.push(event(events.len() as u32, "resume", "failed"));
break;
}
};
had_failure |= result.status != WorkflowRunStatus::Completed;
checkpoint = session
.export_workflow_checkpoint(definition, &result)
.await
.ok();
events.push(event(
events.len() as u32,
"resume",
workflow_status(result.status),
));
last_result = Some(result);
}
}
}
let snapshot = session
.evaluate("window.reliabilityLab.snapshot()")
.await
.ok();
let side_effect_count = expected_side_effects(snapshot.as_ref(), scenario);
let terminal_state = snapshot
.as_ref()
.and_then(|value| value.get("state"))
.and_then(Value::as_str)
.map(str::to_string)
.or_else(|| {
last_result
.as_ref()
.map(|result| workflow_status(result.status).to_string())
});
let mut forbidden_outcomes = Vec::new();
for (name, expected) in &scenario.expect.side_effect_count {
let actual = side_effect_count.get(name).copied().unwrap_or_default();
if actual > *expected {
forbidden_outcomes.push(ReliabilityForbiddenOutcome::NonIdempotentMutationDuplicated);
} else if *expected == 0 && actual > 0 {
forbidden_outcomes.push(ReliabilityForbiddenOutcome::WrongTargetExecuted);
}
}
forbidden_outcomes.sort_unstable();
forbidden_outcomes.dedup();
let actual_terminal = terminal_state.as_deref();
let expected_terminal = scenario.expect.terminal_state.as_str();
let classification = if unsupported {
ReliabilityRunClassification::Unsupported
} else if !forbidden_outcomes.is_empty() {
ReliabilityRunClassification::Failed
} else if !had_failure && actual_terminal == Some(expected_terminal) {
ReliabilityRunClassification::Passed
} else if had_failure
&& actual_terminal == Some(expected_terminal)
&& side_effect_count == scenario.expect.side_effect_count
{
ReliabilityRunClassification::SafeRefusal
} else {
ReliabilityRunClassification::Failed
};
let scenario_hash = scenario.content_hash()?;
let fixture_hash = fixture.content_hash()?;
let elapsed_ms = started.elapsed().as_millis() as u64;
let metadata = ReliabilityRunMetadata {
platform,
browser: "chromium".into(),
browser_version,
duration_ms: elapsed_ms.max(1).min(scenario.budgets.max_duration_ms),
browser_actions: action_count
.max(1)
.min(scenario.budgets.max_browser_actions),
};
let observation = ReliabilityScenarioObservation {
scenario_id: scenario.id.clone(),
scenario_hash: scenario_hash.clone(),
metadata,
classification,
terminal_state,
side_effect_count,
forbidden_outcomes,
oracle_evidence: snapshot.is_some(),
artifacts_complete: snapshot.is_some() && !events.is_empty(),
};
let replay = ReliabilityReplayBundle {
schema_version: crate::reliability::RELIABILITY_REPLAY_SCHEMA_VERSION,
scenario_id: scenario.id.clone(),
scenario_hash,
fixture_id: fixture.id.clone(),
fixture_hash,
events,
observation: observation.clone(),
};
replay.validate(scenario)?;
Ok(ReliabilityRunEvidence {
observation,
replay,
})
}
fn event(sequence: u32, operation: &str, result: &str) -> ReliabilityReplayEvent {
ReliabilityReplayEvent {
sequence,
operation: operation.into(),
result: result.into(),
}
}
fn workflow_status(status: WorkflowRunStatus) -> &'static str {
match status {
WorkflowRunStatus::Completed => "completed",
WorkflowRunStatus::Failed => "failed",
WorkflowRunStatus::BudgetExhausted => "budget_exhausted",
WorkflowRunStatus::ResumeRequired => "resume_required",
}
}
fn expected_side_effects(
snapshot: Option<&Value>,
scenario: &ReliabilityScenario,
) -> BTreeMap<String, u64> {
scenario
.expect
.side_effect_count
.keys()
.map(|name| {
let property = format!("{name}Count");
let count = snapshot
.and_then(|value| value.get(&property))
.and_then(Value::as_u64)
.unwrap_or_default();
(name.clone(), count)
})
.collect()
}
fn bounded_workflow_path(root: &Path, source: &str) -> BrowserResult<PathBuf> {
let root = std::fs::canonicalize(root)?;
let path = root.join(source);
let canonical = std::fs::canonicalize(path)?;
if !canonical.starts_with(&root) {
return Err("workflow source escapes the authorized reliability root".into());
}
Ok(canonical)
}
fn load_workflow(path: &Path) -> BrowserResult<WorkflowDefinition> {
let source = std::fs::read_to_string(path)?;
let format = path.extension().and_then(|extension| extension.to_str());
let document = match format {
Some("yaml") | Some("yml") => compile_workflow_yaml(&source)?,
_ => compile_workflow_json(&source)?,
};
Ok(document.definition)
}
async fn browser_version(session: &BrowserSession) -> BrowserResult<String> {
let Ok(cdp) = session.raw_cdp() else {
return Ok("unknown".into());
};
let value = cdp.send_browser("Browser.getVersion", None).await?;
Ok(value
.get("product")
.and_then(Value::as_str)
.unwrap_or("chromium")
.to_string())
}
fn current_platform() -> BrowserResult<ReliabilityPlatform> {
#[cfg(all(target_os = "linux", target_arch = "x86_64"))]
{
Ok(ReliabilityPlatform::LinuxX86_64)
}
#[cfg(all(target_os = "linux", target_arch = "aarch64"))]
{
Ok(ReliabilityPlatform::LinuxArm64)
}
#[cfg(all(target_os = "macos", target_arch = "x86_64"))]
{
Ok(ReliabilityPlatform::MacosX86_64)
}
#[cfg(all(target_os = "macos", target_arch = "aarch64"))]
{
Ok(ReliabilityPlatform::MacosArm64)
}
#[cfg(not(any(
all(target_os = "linux", target_arch = "x86_64"),
all(target_os = "linux", target_arch = "aarch64"),
all(target_os = "macos", target_arch = "x86_64"),
all(target_os = "macos", target_arch = "aarch64")
)))]
{
Err("reliability runner supports Linux x86-64/arm64 and macOS x86-64/arm64 only".into())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::reliability::{ReliabilityScenarioExpectation, ReliabilityScenarioStep};
#[test]
fn workflow_sources_stay_inside_the_authorized_root() {
let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
assert!(bounded_workflow_path(&root, "workflow-minimal.json").is_ok());
assert!(bounded_workflow_path(&root, "../Cargo.toml").is_err());
}
#[test]
fn workflow_loader_validates_the_checked_in_contract() {
let path =
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/workflow-minimal.json");
let workflow = load_workflow(&path).expect("fixture workflow should compile");
assert_eq!(workflow.name, "open-example");
assert_eq!(workflow.steps.len(), 1);
}
#[test]
fn expected_side_effects_are_limited_to_declared_oracles() {
let scenario = ReliabilityScenario {
schema_version: crate::reliability::RELIABILITY_SCENARIO_SCHEMA_VERSION,
id: "runner-test".into(),
category: "runner".into(),
fixture: "fixture-v1".into(),
platforms: vec![ReliabilityPlatform::LinuxX86_64],
capabilities: Vec::new(),
setup: crate::reliability::ReliabilityScenarioSetup {
browser: "chromium".into(),
policy: "development".into(),
},
steps: vec![ReliabilityScenarioStep {
run_workflow: None,
apply_control: None,
inject: None,
resume_from_checkpoint: None,
}],
expect: ReliabilityScenarioExpectation {
terminal_state: "submitted".into(),
side_effect_count: BTreeMap::from([(String::from("submit"), 1)]),
},
forbid: Vec::new(),
budgets: crate::reliability::ReliabilityScenarioBudgets {
max_duration_ms: 1_000,
max_browser_actions: 4,
},
};
let snapshot = serde_json::json!({
"submitCount": 3,
"ignoredCount": 99,
});
assert_eq!(
expected_side_effects(Some(&snapshot), &scenario)["submit"],
3
);
assert!(!expected_side_effects(Some(&snapshot), &scenario).contains_key("ignored"));
}
}