use std::collections::BTreeMap;
use std::sync::mpsc;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use crate::cancel::CancellationToken;
use crate::diag::Diag;
use crate::engine::{ArtifactRef, EngineFactory, HttpDefaults, ScenarioCtx};
use crate::error::ExitCode;
use crate::event::{EVENT_SCHEMA_VERSION, Event, EventSink};
use crate::step::{Status, StepBatch, StepOutcome};
use crate::world::{GlobalStore, World};
pub struct ScenarioSpec {
pub file: Arc<str>,
pub name: Arc<str>,
pub line: usize,
pub file_root: Option<std::path::PathBuf>,
pub prepare: PrepareFn,
}
pub type PrepareFn = Box<dyn FnOnce(&World) -> Result<Prepared, Vec<Diag>> + Send>;
pub struct Prepared {
pub batches: Vec<StepBatch>,
pub artifact: Option<ArtifactRef>,
}
#[derive(Clone)]
pub struct RunConfig {
pub run_id: Arc<str>,
pub jobs: usize,
pub default_batch_budget: Duration,
pub secrets: Arc<BTreeMap<String, String>>,
pub http: HttpDefaults,
}
#[derive(Debug)]
pub struct RunSummary {
pub outcomes: Vec<ScenarioOutcome>,
pub passed: usize,
pub failed: usize,
pub skipped: usize,
pub cancelled: bool,
}
impl RunSummary {
pub fn exit_code(&self) -> ExitCode {
self.exit_code_excluding(&[])
}
pub fn exit_code_excluding(&self, non_gating: &[(String, String)]) -> ExitCode {
let mut worst = if self.cancelled {
ExitCode::TestFailure
} else {
ExitCode::Success
};
for outcome in &self.outcomes {
let quarantined = non_gating.iter().any(|(file, name)| {
file.as_str() == outcome.file.as_ref() && name.as_str() == outcome.name.as_ref()
});
let code = match (&outcome.fault, outcome.status) {
(Some(Fault::System(_)), _) => ExitCode::SystemError,
(Some(Fault::User(_)), _) => ExitCode::UserError,
(None, Status::Failed) if !quarantined => ExitCode::TestFailure,
_ => ExitCode::Success,
};
worst = pick_worse(worst, code);
}
worst
}
}
fn pick_worse(a: ExitCode, b: ExitCode) -> ExitCode {
let rank = |c: ExitCode| match c {
ExitCode::SystemError => 3,
ExitCode::UserError => 2,
ExitCode::TestFailure => 1,
ExitCode::Success => 0,
};
if rank(b) > rank(a) { b } else { a }
}
#[derive(Debug)]
pub struct ScenarioOutcome {
pub file: Arc<str>,
pub name: Arc<str>,
pub line: usize,
pub status: Status,
pub steps: Vec<StepOutcome>,
pub fault: Option<Fault>,
pub artifact_slug: Option<Arc<str>>,
}
#[derive(Debug)]
pub enum Fault {
User(String),
System(String),
}
enum Msg {
BatchBegin {
scenario: usize,
deadline: Instant,
},
Done {
scenario: usize,
outcome: ScenarioOutcome,
},
}
#[allow(clippy::too_many_lines)]
pub fn run(
specs: Vec<ScenarioSpec>,
engines: &Arc<Vec<Box<dyn EngineFactory>>>,
store: &Arc<Mutex<GlobalStore>>,
config: &RunConfig,
events: &EventSink,
cancel: &CancellationToken,
) -> RunSummary {
events.emit(&Event::RunStarted {
schema: EVENT_SCHEMA_VERSION,
run_id: Arc::clone(&config.run_id),
});
let (tx, rx) = mpsc::channel::<Msg>();
let identities: Vec<(Arc<str>, Arc<str>, usize)> = specs
.iter()
.map(|spec| (Arc::clone(&spec.file), Arc::clone(&spec.name), spec.line))
.collect();
let mut queue: std::collections::VecDeque<(usize, ScenarioSpec)> =
specs.into_iter().enumerate().collect();
let total = queue.len();
let mut active: BTreeMap<usize, (Instant, CancellationToken)> = BTreeMap::new();
let mut outcomes: Vec<ScenarioOutcome> = Vec::new();
let grace = Duration::from_secs(2);
while outcomes.len() < total {
while active.len() < config.jobs.max(1) {
let Some((index, spec)) = queue.pop_front() else {
break;
};
if cancel.is_cancelled() {
events.emit(&Event::ScenarioFinished {
scenario: Arc::clone(&spec.name),
file: Arc::clone(&spec.file),
status: Status::Skipped,
timestamp_ms: None,
worker: None,
});
outcomes.push(ScenarioOutcome {
file: spec.file,
name: spec.name,
line: spec.line,
status: Status::Skipped,
steps: Vec::new(),
fault: None,
artifact_slug: None,
});
continue;
}
let initial_deadline = Instant::now() + config.default_batch_budget + grace;
let child = cancel.child_token();
active.insert(index, (initial_deadline, child.clone()));
spawn_scenario(
index,
spec,
Arc::clone(engines),
Arc::clone(store),
config.clone(),
events.clone(),
child,
tx.clone(),
);
}
if active.is_empty() {
continue; }
let now = Instant::now();
let next_deadline = active
.values()
.map(|(deadline, _)| *deadline)
.min()
.unwrap_or(now + grace);
let wait = next_deadline
.saturating_duration_since(now)
.max(Duration::from_millis(20));
match rx.recv_timeout(wait) {
Ok(Msg::BatchBegin { scenario, deadline }) => {
if let Some(entry) = active.get_mut(&scenario) {
entry.0 = deadline + grace;
}
}
Ok(Msg::Done { scenario, outcome }) => {
if active.remove(&scenario).is_some() {
events.emit(&Event::ScenarioFinished {
scenario: Arc::clone(&outcome.name),
file: Arc::clone(&outcome.file),
status: outcome.status,
timestamp_ms: None,
worker: None,
});
outcomes.push(outcome);
}
}
Err(mpsc::RecvTimeoutError::Timeout) => {}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
}
sweep_expired(&mut active, &mut outcomes, events, &identities);
}
let passed = outcomes
.iter()
.filter(|o| matches!(o.status, Status::Passed | Status::Warned))
.count();
let failed = outcomes
.iter()
.filter(|o| o.status == Status::Failed)
.count();
let skipped = outcomes
.iter()
.filter(|o| o.status == Status::Skipped)
.count();
let cancelled = cancel.is_cancelled();
events.emit(&Event::RunFinished {
passed,
failed,
skipped,
cancelled,
});
RunSummary {
outcomes,
passed,
failed,
skipped,
cancelled,
}
}
fn sweep_expired(
active: &mut BTreeMap<usize, (Instant, CancellationToken)>,
outcomes: &mut Vec<ScenarioOutcome>,
events: &EventSink,
identities: &[(Arc<str>, Arc<str>, usize)],
) {
let now = Instant::now();
let expired: Vec<usize> = active
.iter()
.filter(|(_, (deadline, _))| *deadline <= now)
.map(|(index, _)| *index)
.collect();
for index in expired {
if let Some((_, token)) = active.remove(&index) {
token.cancel();
}
let (file, name, line) = &identities[index];
let outcome = ScenarioOutcome {
file: Arc::clone(file),
name: Arc::clone(name),
line: *line,
status: Status::Failed,
steps: Vec::new(),
fault: Some(Fault::System(
"batch budget exceeded — scenario thread abandoned (ADR-0007)".to_owned(),
)),
artifact_slug: None,
};
events.emit(&Event::ScenarioFinished {
scenario: Arc::clone(&outcome.name),
file: Arc::clone(&outcome.file),
status: Status::Failed,
timestamp_ms: None,
worker: None,
});
outcomes.push(outcome);
}
}
#[allow(clippy::too_many_arguments)]
fn spawn_scenario(
index: usize,
spec: ScenarioSpec,
engines: Arc<Vec<Box<dyn EngineFactory>>>,
store: Arc<Mutex<GlobalStore>>,
config: RunConfig,
events: EventSink,
cancel: CancellationToken,
tx: mpsc::Sender<Msg>,
) {
std::thread::spawn(move || {
let identity = (Arc::clone(&spec.file), Arc::clone(&spec.name), spec.line);
let heartbeat_tx = tx.clone();
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
run_scenario(
spec,
&engines,
&store,
&config,
&events,
&cancel,
|budget| {
let _ = heartbeat_tx.send(Msg::BatchBegin {
scenario: index,
deadline: Instant::now() + budget,
});
},
)
}));
let outcome = result.unwrap_or_else(|panic| {
let message = panic
.downcast_ref::<&str>()
.map(ToString::to_string)
.or_else(|| panic.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "opaque panic payload".to_owned());
ScenarioOutcome {
file: identity.0,
name: identity.1,
line: identity.2,
status: Status::Failed,
steps: Vec::new(),
fault: Some(Fault::System(format!(
"scenario thread panicked: {message}"
))),
artifact_slug: None,
}
});
let _ = tx.send(Msg::Done {
scenario: index,
outcome,
});
});
}
#[allow(clippy::too_many_lines)]
fn run_scenario(
spec: ScenarioSpec,
engines: &[Box<dyn EngineFactory>],
store: &Mutex<GlobalStore>,
config: &RunConfig,
events: &EventSink,
cancel: &CancellationToken,
heartbeat: impl Fn(Duration),
) -> ScenarioOutcome {
let (file, name, line) = (Arc::clone(&spec.file), Arc::clone(&spec.name), spec.line);
let outcome = move |status, steps, fault, artifact_slug| ScenarioOutcome {
file: Arc::clone(&file),
name: Arc::clone(&name),
line,
status,
steps,
fault,
artifact_slug,
};
events.emit(&Event::ScenarioStarted {
scenario: Arc::clone(&spec.name),
file: Arc::clone(&spec.file),
timestamp_ms: None,
worker: None,
});
let snapshot = match store.lock() {
Ok(guard) => guard.clone(),
Err(_) => {
return outcome(
Status::Failed,
Vec::new(),
Some(Fault::System("global store lock poisoned".to_owned())),
None,
);
}
};
let mut world = World::new(snapshot);
let prepared = match (spec.prepare)(&world) {
Ok(prepared) => prepared,
Err(diags) => {
let detail = diags
.iter()
.map(|d| d.message.clone())
.collect::<Vec<_>>()
.join("; ");
return outcome(Status::Failed, Vec::new(), Some(Fault::User(detail)), None);
}
};
let mut sessions: Vec<(String, Box<dyn crate::engine::EngineSession>)> = Vec::new();
let mut steps: Vec<StepOutcome> = Vec::new();
let mut fault: Option<Fault> = None;
let mut failed = false;
let mut interrupted = false;
let mut processed = 0usize;
for batch in &prepared.batches {
if failed {
break;
}
if cancel.is_cancelled() {
interrupted = true;
break;
}
let engine_id = batch.engine.as_str().to_owned();
if !sessions.iter().any(|(id, _)| *id == engine_id) {
let Some(factory) = engines.iter().find(|f| f.id() == engine_id) else {
fault = Some(Fault::System(format!(
"no engine registered for `{engine_id}`"
)));
failed = true;
break;
};
let ctx = ScenarioCtx {
run_id: Arc::clone(&config.run_id),
scenario: Arc::clone(&spec.name),
artifact: prepared.artifact.clone(),
secrets: Arc::clone(&config.secrets),
http: config.http,
file_root: spec.file_root.clone(),
};
match factory.open(&ctx) {
Ok(session) => sessions.push((engine_id.clone(), session)),
Err(err) => {
fault = Some(Fault::System(format!(
"cannot open engine `{engine_id}`: {err}"
)));
failed = true;
break;
}
}
}
let Some((_, session)) = sessions.iter_mut().find(|(id, _)| *id == engine_id) else {
break; };
let budget = session
.batch_budget(batch)
.unwrap_or(config.default_batch_budget);
heartbeat(budget);
events.emit(&Event::BatchStarted {
scenario: Arc::clone(&spec.name),
engine: Arc::from(engine_id.as_str()),
steps: batch.steps.len(),
});
let result = session.run_batch(batch, &mut world, events, cancel);
let all_optional = batch.steps.iter().all(|s| s.optional);
for mut step_outcome in result.steps {
if step_outcome.status == Status::Failed && all_optional {
step_outcome.status = Status::Warned;
}
steps.push(step_outcome);
}
if let Some(err) = result.error {
if all_optional {
processed += 1;
continue;
}
match err.class {
crate::error::EngineErrorClass::AssertFailed => {}
crate::error::EngineErrorClass::UserInput => {
fault = Some(Fault::User(err.message.clone()));
}
crate::error::EngineErrorClass::Infra | crate::error::EngineErrorClass::Setup => {
fault = Some(Fault::System(err.message.clone()));
}
}
failed = true;
}
processed += 1;
}
let unreached_reason = if interrupted {
"not run (run cancelled)"
} else {
"not run (an earlier step failed)"
};
for batch in prepared.batches.iter().skip(processed) {
for step in &batch.steps {
events.emit(&Event::StepFinished {
scenario: Arc::clone(&spec.name),
engine: Arc::from(batch.engine.as_str()),
step: step.step.clone(),
status: Status::Skipped,
attempts: 0,
duration_ms: 0,
captures: Vec::new(),
detail: Some(unreached_reason.to_owned()),
attempt_details: Vec::new(),
});
steps.push(StepOutcome {
step: step.step.clone(),
status: Status::Skipped,
attempts: 0,
duration: std::time::Duration::ZERO,
detail: Some(unreached_reason.to_owned()),
attempt_details: Vec::new(),
reproduce_hint: None,
});
}
}
for (engine_id, session) in sessions.iter_mut().rev() {
if let Err(err) = session.finish()
&& fault.is_none()
{
fault = Some(Fault::System(format!(
"engine `{engine_id}` teardown failed: {}",
err.message
)));
}
}
{
let mut guard = store
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for (key, value) in world.promotions() {
guard.insert(key, value.clone());
}
}
let status = if failed || steps.iter().any(|s| s.status == Status::Failed) {
Status::Failed
} else if interrupted {
Status::Skipped
} else {
Status::Passed
};
let artifact_slug = prepared
.artifact
.as_ref()
.map(|artifact| Arc::clone(&artifact.slug));
outcome(status, steps, fault, artifact_slug)
}