1use crate::browser::session::{
9 BrowserResult, BrowserSession, WorkflowCheckpoint, WorkflowDefinition, WorkflowRunResult,
10 WorkflowRunStatus, compile_workflow_json, compile_workflow_yaml,
11};
12use crate::reliability::{
13 ReliabilityExecutionOperation, ReliabilityFaultKind, ReliabilityFixtureManifest,
14 ReliabilityForbiddenOutcome, ReliabilityPlatform, ReliabilityReplayBundle,
15 ReliabilityReplayEvent, ReliabilityRunClassification, ReliabilityRunMetadata,
16 ReliabilityScenario, ReliabilityScenarioObservation,
17};
18use serde_json::Value;
19use std::collections::BTreeMap;
20use std::path::{Path, PathBuf};
21use std::time::Instant;
22
23#[derive(Debug, Clone)]
25pub struct ReliabilityRunOptions {
26 pub workflow_root: PathBuf,
27 pub inputs: BTreeMap<String, Value>,
28}
29
30#[derive(Debug, Clone)]
32pub struct ReliabilityRunEvidence {
33 pub observation: ReliabilityScenarioObservation,
34 pub replay: ReliabilityReplayBundle,
35}
36
37pub async fn run_reliability_scenario(
39 session: &BrowserSession,
40 scenario: &ReliabilityScenario,
41 fixture: &ReliabilityFixtureManifest,
42 options: &ReliabilityRunOptions,
43) -> BrowserResult<ReliabilityRunEvidence> {
44 let plan = scenario.execution_plan(fixture)?;
45 let platform = current_platform()?;
46 let browser_version = browser_version(session).await?;
47 let started = Instant::now();
48 let mut events = Vec::new();
49 let mut action_count = 0u32;
50 let mut workflow: Option<WorkflowDefinition> = None;
51 let mut checkpoint: Option<WorkflowCheckpoint> = None;
52 let mut last_result: Option<WorkflowRunResult> = None;
53 let mut had_failure = false;
54 let mut unsupported = false;
55
56 for operation in &plan.operations {
57 if started.elapsed().as_millis() as u64 >= scenario.budgets.max_duration_ms {
58 had_failure = true;
59 events.push(event(events.len() as u32, "budget", "exhausted"));
60 break;
61 }
62 if action_count >= scenario.budgets.max_browser_actions {
63 had_failure = true;
64 events.push(event(events.len() as u32, "budget", "actions_exhausted"));
65 break;
66 }
67 action_count = action_count.saturating_add(1);
68 match operation {
69 ReliabilityExecutionOperation::ApplyControl { control } => {
70 let expression = format!(
71 "window.reliabilityLab.{}(); true",
72 control.javascript_method()
73 );
74 match session.evaluate(&expression).await {
75 Ok(_) => events.push(event(events.len() as u32, "applyControl", "applied")),
76 Err(_) => {
77 had_failure = true;
78 events.push(event(events.len() as u32, "applyControl", "failed"));
79 break;
80 }
81 }
82 }
83 ReliabilityExecutionOperation::InjectFault { injection } => {
84 if matches!(
85 injection.fault,
86 ReliabilityFaultKind::RendererDisconnect
87 | ReliabilityFaultKind::BrowserDisconnect
88 ) {
89 let dispatched = match session.raw_cdp() {
90 Ok(cdp) => match injection.fault {
91 ReliabilityFaultKind::RendererDisconnect => {
92 cdp.send("Page.crash", None).await.is_ok()
93 }
94 ReliabilityFaultKind::BrowserDisconnect => {
95 cdp.send_browser("Browser.close", None).await.is_ok()
96 }
97 _ => unreachable!("transport fault branch is exhaustive"),
98 },
99 Err(_) => false,
100 };
101 events.push(event(
102 events.len() as u32,
103 "injectFault",
104 if dispatched {
105 "dispatched"
106 } else {
107 "dispatch_failed"
108 },
109 ));
110 unsupported = true;
111 break;
112 }
113 let fault = serde_json::to_string(injection.fault.fixture_name())?;
114 let expression = format!("window.reliabilityLab.injectFault({fault}); true");
115 match session.evaluate(&expression).await {
116 Ok(_) => events.push(event(events.len() as u32, "injectFault", "applied")),
117 Err(_) => {
118 had_failure = true;
119 events.push(event(events.len() as u32, "injectFault", "failed"));
120 break;
121 }
122 }
123 }
124 ReliabilityExecutionOperation::RunWorkflow { source } => {
125 let path = bounded_workflow_path(&options.workflow_root, source)?;
126 let definition = load_workflow(&path)?;
127 let result = match session.run_workflow(&definition, &options.inputs).await {
128 Ok(result) => result,
129 Err(_) => {
130 had_failure = true;
131 events.push(event(events.len() as u32, "runWorkflow", "failed"));
132 break;
133 }
134 };
135 action_count = action_count
136 .saturating_add(result.trace.events.len().min(u32::MAX as usize) as u32);
137 had_failure |= result.status != WorkflowRunStatus::Completed;
138 checkpoint = session
139 .export_workflow_checkpoint(&definition, &result)
140 .await
141 .ok();
142 workflow = Some(definition);
143 events.push(event(
144 events.len() as u32,
145 "runWorkflow",
146 workflow_status(result.status),
147 ));
148 last_result = Some(result);
149 }
150 ReliabilityExecutionOperation::ResumeFromCheckpoint { checkpoint: name } => {
151 if name != "latest" {
152 unsupported = true;
153 events.push(event(
154 events.len() as u32,
155 "resume",
156 "unsupported_checkpoint",
157 ));
158 break;
159 }
160 let (Some(definition), Some(saved_checkpoint)) = (&workflow, &checkpoint) else {
161 unsupported = true;
162 events.push(event(events.len() as u32, "resume", "missing_checkpoint"));
163 break;
164 };
165 let result = match session
166 .resume_workflow(definition, &options.inputs, saved_checkpoint)
167 .await
168 {
169 Ok(result) => result,
170 Err(_) => {
171 had_failure = true;
172 events.push(event(events.len() as u32, "resume", "failed"));
173 break;
174 }
175 };
176 had_failure |= result.status != WorkflowRunStatus::Completed;
177 checkpoint = session
178 .export_workflow_checkpoint(definition, &result)
179 .await
180 .ok();
181 events.push(event(
182 events.len() as u32,
183 "resume",
184 workflow_status(result.status),
185 ));
186 last_result = Some(result);
187 }
188 }
189 }
190
191 let snapshot = session
192 .evaluate("window.reliabilityLab.snapshot()")
193 .await
194 .ok();
195 let side_effect_count = expected_side_effects(snapshot.as_ref(), scenario);
196 let terminal_state = snapshot
197 .as_ref()
198 .and_then(|value| value.get("state"))
199 .and_then(Value::as_str)
200 .map(str::to_string)
201 .or_else(|| {
202 last_result
203 .as_ref()
204 .map(|result| workflow_status(result.status).to_string())
205 });
206 let mut forbidden_outcomes = Vec::new();
207 for (name, expected) in &scenario.expect.side_effect_count {
208 let actual = side_effect_count.get(name).copied().unwrap_or_default();
209 if actual > *expected {
210 forbidden_outcomes.push(ReliabilityForbiddenOutcome::NonIdempotentMutationDuplicated);
211 } else if *expected == 0 && actual > 0 {
212 forbidden_outcomes.push(ReliabilityForbiddenOutcome::WrongTargetExecuted);
213 }
214 }
215 forbidden_outcomes.sort_unstable();
216 forbidden_outcomes.dedup();
217 let actual_terminal = terminal_state.as_deref();
218 let expected_terminal = scenario.expect.terminal_state.as_str();
219 let classification = if unsupported {
220 ReliabilityRunClassification::Unsupported
221 } else if !forbidden_outcomes.is_empty() {
222 ReliabilityRunClassification::Failed
223 } else if !had_failure && actual_terminal == Some(expected_terminal) {
224 ReliabilityRunClassification::Passed
225 } else if had_failure
226 && expected_terminal == "refused"
227 && side_effect_count == scenario.expect.side_effect_count
228 {
229 ReliabilityRunClassification::SafeRefusal
230 } else {
231 ReliabilityRunClassification::Failed
232 };
233 let scenario_hash = scenario.content_hash()?;
234 let fixture_hash = fixture.content_hash()?;
235 let elapsed_ms = started.elapsed().as_millis() as u64;
236 let metadata = ReliabilityRunMetadata {
237 platform,
238 browser: "chromium".into(),
239 browser_version,
240 duration_ms: elapsed_ms.max(1).min(scenario.budgets.max_duration_ms),
241 browser_actions: action_count
242 .max(1)
243 .min(scenario.budgets.max_browser_actions),
244 };
245 let observation = ReliabilityScenarioObservation {
246 scenario_id: scenario.id.clone(),
247 scenario_hash: scenario_hash.clone(),
248 metadata,
249 classification,
250 terminal_state,
251 side_effect_count,
252 forbidden_outcomes,
253 oracle_evidence: snapshot.is_some(),
254 artifacts_complete: snapshot.is_some() && !events.is_empty(),
255 };
256 let replay = ReliabilityReplayBundle {
257 schema_version: crate::reliability::RELIABILITY_REPLAY_SCHEMA_VERSION,
258 scenario_id: scenario.id.clone(),
259 scenario_hash,
260 fixture_id: fixture.id.clone(),
261 fixture_hash,
262 events,
263 observation: observation.clone(),
264 };
265 replay.validate(scenario)?;
266 Ok(ReliabilityRunEvidence {
267 observation,
268 replay,
269 })
270}
271
272fn event(sequence: u32, operation: &str, result: &str) -> ReliabilityReplayEvent {
273 ReliabilityReplayEvent {
274 sequence,
275 operation: operation.into(),
276 result: result.into(),
277 }
278}
279
280fn workflow_status(status: WorkflowRunStatus) -> &'static str {
281 match status {
282 WorkflowRunStatus::Completed => "completed",
283 WorkflowRunStatus::Failed => "failed",
284 WorkflowRunStatus::BudgetExhausted => "budget_exhausted",
285 WorkflowRunStatus::ResumeRequired => "resume_required",
286 }
287}
288
289fn expected_side_effects(
290 snapshot: Option<&Value>,
291 scenario: &ReliabilityScenario,
292) -> BTreeMap<String, u64> {
293 scenario
294 .expect
295 .side_effect_count
296 .keys()
297 .map(|name| {
298 let property = format!("{name}Count");
299 let count = snapshot
300 .and_then(|value| value.get(&property))
301 .and_then(Value::as_u64)
302 .unwrap_or_default();
303 (name.clone(), count)
304 })
305 .collect()
306}
307
308fn bounded_workflow_path(root: &Path, source: &str) -> BrowserResult<PathBuf> {
309 let root = std::fs::canonicalize(root)?;
310 let path = root.join(source);
311 let canonical = std::fs::canonicalize(path)?;
312 if !canonical.starts_with(&root) {
313 return Err("workflow source escapes the authorized reliability root".into());
314 }
315 Ok(canonical)
316}
317
318fn load_workflow(path: &Path) -> BrowserResult<WorkflowDefinition> {
319 let source = std::fs::read_to_string(path)?;
320 let format = path.extension().and_then(|extension| extension.to_str());
321 let document = match format {
322 Some("yaml") | Some("yml") => compile_workflow_yaml(&source)?,
323 _ => compile_workflow_json(&source)?,
324 };
325 Ok(document.definition)
326}
327
328async fn browser_version(session: &BrowserSession) -> BrowserResult<String> {
329 let Ok(cdp) = session.raw_cdp() else {
330 return Ok("unknown".into());
331 };
332 let value = cdp.send_browser("Browser.getVersion", None).await?;
333 Ok(value
334 .get("product")
335 .and_then(Value::as_str)
336 .unwrap_or("chromium")
337 .to_string())
338}
339
340fn current_platform() -> BrowserResult<ReliabilityPlatform> {
341 #[cfg(all(target_os = "linux", target_arch = "x86_64"))]
342 {
343 Ok(ReliabilityPlatform::LinuxX86_64)
344 }
345 #[cfg(all(target_os = "linux", target_arch = "aarch64"))]
346 {
347 Ok(ReliabilityPlatform::LinuxArm64)
348 }
349 #[cfg(all(target_os = "macos", target_arch = "x86_64"))]
350 {
351 Ok(ReliabilityPlatform::MacosX86_64)
352 }
353 #[cfg(all(target_os = "macos", target_arch = "aarch64"))]
354 {
355 Ok(ReliabilityPlatform::MacosArm64)
356 }
357 #[cfg(not(any(
358 all(target_os = "linux", target_arch = "x86_64"),
359 all(target_os = "linux", target_arch = "aarch64"),
360 all(target_os = "macos", target_arch = "x86_64"),
361 all(target_os = "macos", target_arch = "aarch64")
362 )))]
363 {
364 Err("reliability runner supports Linux x86-64/arm64 and macOS x86-64/arm64 only".into())
365 }
366}
367
368#[cfg(test)]
369mod tests {
370 use super::*;
371 use crate::reliability::{ReliabilityScenarioExpectation, ReliabilityScenarioStep};
372
373 #[test]
374 fn workflow_sources_stay_inside_the_authorized_root() {
375 let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
376
377 assert!(bounded_workflow_path(&root, "workflow-minimal.json").is_ok());
378 assert!(bounded_workflow_path(&root, "../Cargo.toml").is_err());
379 }
380
381 #[test]
382 fn workflow_loader_validates_the_checked_in_contract() {
383 let path =
384 PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/workflow-minimal.json");
385 let workflow = load_workflow(&path).expect("fixture workflow should compile");
386
387 assert_eq!(workflow.name, "open-example");
388 assert_eq!(workflow.steps.len(), 1);
389 }
390
391 #[test]
392 fn expected_side_effects_are_limited_to_declared_oracles() {
393 let scenario = ReliabilityScenario {
394 schema_version: crate::reliability::RELIABILITY_SCENARIO_SCHEMA_VERSION,
395 id: "runner-test".into(),
396 category: "runner".into(),
397 fixture: "fixture-v1".into(),
398 platforms: vec![ReliabilityPlatform::LinuxX86_64],
399 capabilities: Vec::new(),
400 setup: crate::reliability::ReliabilityScenarioSetup {
401 browser: "chromium".into(),
402 policy: "development".into(),
403 },
404 steps: vec![ReliabilityScenarioStep {
405 run_workflow: None,
406 apply_control: None,
407 inject: None,
408 resume_from_checkpoint: None,
409 }],
410 expect: ReliabilityScenarioExpectation {
411 terminal_state: "submitted".into(),
412 side_effect_count: BTreeMap::from([(String::from("submit"), 1)]),
413 },
414 forbid: Vec::new(),
415 budgets: crate::reliability::ReliabilityScenarioBudgets {
416 max_duration_ms: 1_000,
417 max_browser_actions: 4,
418 },
419 };
420 let snapshot = serde_json::json!({
421 "submitCount": 3,
422 "ignoredCount": 99,
423 });
424
425 assert_eq!(
426 expected_side_effects(Some(&snapshot), &scenario)["submit"],
427 3
428 );
429 assert!(!expected_side_effects(Some(&snapshot), &scenario).contains_key("ignored"));
430 }
431}