use crate::feature::{ExpandedStep, LoadedFeature, expand_outlines};
use crate::options::Options;
use crate::plugin::abi::{DispatchRequest, OptionsJson, Status};
use crate::polling::{AttemptError, Polling};
use crate::report::render_file;
use crate::report::{FileResult, ScenarioResult, StepResult, StepStatus};
use crate::steps::{Args, OptionsSource, Registry, StepId, StepKind, StepTarget, dispatch};
use crate::unique::Generator;
use crate::vars::{VarStack, interpolate};
use crate::world::World;
use anyhow::Result;
use std::collections::BTreeMap;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
fn prepare(
step: &ExpandedStep,
caps: Vec<String>,
vars: &VarStack,
generator: &Generator,
) -> Result<Args, String> {
let caps = caps
.iter()
.map(|c| interpolate(c, vars, generator))
.collect::<Result<Vec<_>, _>>()?;
let docstring = step
.docstring
.as_ref()
.map(|d| interpolate(d, vars, generator))
.transpose()?;
let table = step
.table
.as_ref()
.map(|rows| {
rows.iter()
.map(|r| {
r.iter()
.map(|c| interpolate(c, vars, generator))
.collect::<Result<Vec<_>, _>>()
})
.collect::<Result<Vec<_>, _>>()
})
.transpose()?;
Ok(Args {
caps,
docstring,
table,
})
}
fn execute_step<'a>(
world: &'a mut World,
reg: &'a Registry,
step: &'a ExpandedStep,
source: &'a std::path::Path,
generator: &'a Generator,
depth: usize,
) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + 'a>> {
Box::pin(async move {
let Some((target, caps)) = reg.find(&step.text)? else {
return Err("unknown step".into());
};
match target {
StepTarget::Builtin { id, kind } => {
if matches!(id, StepId::Include | StepId::IncludeScenario) {
return run_include(world, reg, step, id, caps, source, generator, depth).await;
}
let args = prepare(step, caps, &world.vars, generator)?;
match kind {
StepKind::Action => dispatch(world, id, &args, 0)
.await
.map_err(AttemptError::into_message),
StepKind::Assertion(source_kind) => {
let Some(layer) = world.take_options() else {
return dispatch(world, id, &args, 0)
.await
.map_err(AttemptError::into_message);
};
let base = match source_kind {
OptionsSource::Global => world.options.clone(),
OptionsSource::Http => world.http.options_for_last_response()?.clone(),
OptionsSource::Db => world.db.options()?.clone(),
};
let effective = base.apply(&layer)?;
let mut polling = Polling::new(&step.text, &effective.polling);
let mut attempt = 0;
loop {
match dispatch(world, id, &args, attempt).await {
Ok(()) => return Ok(()),
Err(AttemptError::Fatal(error)) => return Err(error),
Err(AttemptError::NotYet(error)) => {
polling.after_not_yet(&error).await?;
attempt += 1;
}
}
}
}
}
}
StepTarget::Macro(index) => {
if step.docstring.is_some() {
return Err("a macro call does not support a docstring".into());
}
if step.table.is_some() {
return Err("a macro call does not support a table".into());
}
if depth >= 16 {
return Err("macro nesting exceeds 16".into());
}
let definition = reg.macro_def(index);
let args = prepare(step, caps, &world.vars, generator)?;
world.vars.push_frame();
for (name, value) in definition.params.iter().zip(args.caps) {
world.vars.set(name, value);
}
if world.debug {
eprintln!("macro {:?}", step.text);
}
let macro_source = definition.source.clone();
for body_step in &definition.body {
let expanded = ExpandedStep {
keyword: String::new(),
text: body_step.text.clone(),
line: step.line,
docstring: body_step.docstring.clone(),
table: None,
};
if let Err(error) =
execute_step(world, reg, &expanded, ¯o_source, generator, depth + 1)
.await
{
world.vars.pop_frame(&[])?;
return Err(format!(" {}\n{error}", body_step.text));
}
}
let exported = world.vars.pop_frame(&definition.exports)?;
if world.debug {
for (k, v) in &exported {
eprintln!(" export {k} = {}", debug_display(v));
}
}
Ok(())
}
StepTarget::Plugin {
lib,
step: step_index,
assertion,
} => {
let args = prepare(step, caps, &world.vars, generator)?;
let Some(plugins) = world.plugins.plugins().cloned() else {
return Err("this step is served by a plugin, but no plugin is loaded".into());
};
let group = plugins.group_of_step(lib, step_index).to_string();
let instance = world.plugins.current(&group)?.to_string();
let base = plugins.options_for(&group, &instance)?.clone();
let layer = if assertion {
world.take_options()
} else {
None
};
let effective = match &layer {
Some(layer) => base.apply(layer)?,
None => base,
};
let mut polling = layer
.as_ref()
.map(|_| Polling::new(&step.text, &effective.polling));
let per_worker = plugins.is_per_worker(lib);
let workspace_dir = world.workspace_dir()?.display().to_string();
loop {
let request = serde_json::to_string(&DispatchRequest {
args: &args.caps,
docstring: args.docstring.as_ref(),
table: args.table.as_ref(),
artifacts_dir: plugins.next_artifacts_dir(),
workspace_dir: workspace_dir.clone(),
debug: world.debug,
options: OptionsJson::from(&effective),
})
.map_err(|error| format!("failed to encode the plugin request: {error}"))?;
let plugins_for_call = plugins.clone();
let (call_group, call_instance) = (group.clone(), instance.clone());
let cached = world
.plugins
.handle_for(&group, &instance)
.filter(|(l, _)| *l == lib)
.map(|(_, h)| h);
let (handle, result) = tokio::task::spawn_blocking(move || {
plugins_for_call.call_step(
&call_group,
&call_instance,
lib,
step_index,
&request,
cached,
)
})
.await
.map_err(|error| format!("the plugin dispatch task failed: {error}"))??;
world
.plugins
.record(&group, &instance, lib, handle, per_worker);
match result.status {
Status::Passed => {
for (name, value) in result.vars {
world.vars.set(&name, value);
}
return Ok(());
}
Status::Fatal => return Err(result.render_failure()),
Status::NotYet if !assertion => {
return Err(format!(
"a plugin action answered not_yet, which only an assertion may do: {}",
result.render_failure()
));
}
Status::NotYet => match &mut polling {
None => return Err(result.render_failure()),
Some(polling) => {
polling.after_not_yet(&result.render_failure()).await?
}
},
}
}
}
}
})
}
#[allow(clippy::too_many_arguments)]
fn run_include<'a>(
world: &'a mut World,
reg: &'a Registry,
caller_step: &'a ExpandedStep,
id: StepId,
caps: Vec<String>,
source: &'a std::path::Path,
generator: &'a Generator,
depth: usize,
) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + 'a>> {
Box::pin(async move {
if caller_step.docstring.is_some() {
return Err("I include does not support a docstring".into());
}
if depth >= 16 {
return Err("include nesting exceeds 16".into());
}
let path_literal = &caps[0];
let (scenario_name, prefix_raw) = if id == StepId::IncludeScenario {
(Some(caps[1].as_str()), caps[2].as_str())
} else {
(None, caps[1].as_str())
};
let prefix = (!prefix_raw.is_empty()).then_some(prefix_raw);
let base_dir = source
.parent()
.map(std::path::Path::to_path_buf)
.unwrap_or_else(|| std::path::PathBuf::from("."));
crate::include::check_literal(path_literal)?;
let resolved = crate::include::resolve(path_literal, &base_dir)?;
let included = crate::feature::load(&resolved).map_err(|e| e.to_string())?;
let scenario = crate::include::select_scenario(&included, scenario_name)?;
let exports = crate::feature::exports_of(&included, scenario)?;
let expanded = crate::include::build_scenario(scenario, caller_step.table.as_deref())?;
let with_values: Vec<(String, String)> = match &caller_step.table {
Some(table) => {
let [header, row] = table.as_slice() else {
unreachable!("build_scenario already validated exactly one data row")
};
let mut out = Vec::with_capacity(header.len());
for (name, raw) in header.iter().zip(row.iter()) {
let value = interpolate(raw, &world.vars, generator)?;
out.push((name.clone(), value));
}
out
}
None => Vec::new(),
};
if world.debug {
eprintln!(
"include {} › {:?}",
crate::feature::display_path(&resolved),
expanded.name
);
for (k, v) in &with_values {
eprintln!(" with {k} = {}", debug_display(v));
}
}
let included_base = resolved
.parent()
.map(std::path::Path::to_path_buf)
.unwrap_or_else(|| std::path::PathBuf::from("."));
let background: Vec<ExpandedStep> = included
.feature
.background
.as_ref()
.map(|bg| bg.steps.iter().map(crate::feature::to_step).collect())
.unwrap_or_default();
let fresh = VarStack::new();
let saved = std::mem::replace(&mut world.vars, fresh);
for (name, value) in &with_values {
world.vars.set(name, value.clone());
}
let mut run_result: Result<(), String> = Ok(());
for step in background.iter().chain(expanded.steps.iter()) {
if let Err(e) =
execute_step(world, reg, step, &included_base, generator, depth + 1).await
{
run_result = Err(format!(
"{}:{} → {}:{}\n {}\n{e}",
crate::feature::display_path(source),
caller_step.line,
crate::feature::display_path(&resolved),
step.line,
step.text
));
break;
}
}
let apply_prefix = |name: String| match prefix {
Some(p) => format!("{p}_{name}"),
None => name,
};
let export_result = if run_result.is_ok() {
let mut resolved_exports = Vec::new();
let mut missing = None;
for name in &exports {
if let Some(glob_prefix) = name.strip_suffix('*') {
for (k, v) in world.vars.all_vars() {
if k.starts_with(glob_prefix) {
resolved_exports.push((apply_prefix(k), v));
}
}
} else if let Some(v) = world.vars.get(name) {
resolved_exports.push((apply_prefix(name.clone()), v.to_string()));
} else {
missing = Some(format!(
"I include {path_literal:?}: declared export {name:?}, but the variable is not set"
));
break;
}
}
match missing {
Some(e) => Err(e),
None => Ok(resolved_exports),
}
} else {
Ok(Vec::new())
};
if world.debug
&& let Ok(exported) = &export_result
{
for (k, v) in exported {
eprintln!(" export {k} = {}", debug_display(v));
}
}
world.vars = saved;
run_result?;
let exported = export_result?;
for (name, value) in exported {
world.vars.set(&name, value);
}
Ok(())
})
}
fn debug_display(value: &str) -> &str {
if value == crate::vars::NULL_SENTINEL {
"<<null>>"
} else {
value
}
}
pub struct RunContext {
pub reg: Registry,
pub apis: Arc<crate::http::Apis>,
pub generator: Arc<Generator>,
pub filter: crate::feature::TagFilter,
pub db: Option<Arc<crate::db::Db>>,
pub default_db: String,
pub srp: Option<Arc<crate::srp::SrpParams>>,
pub plugins: Option<Arc<crate::plugin::Plugins>>,
pub options: Options,
fail_fast: bool,
stop: std::sync::atomic::AtomicBool,
}
impl RunContext {
#[allow(clippy::too_many_arguments)]
pub fn new(
reg: Registry,
apis: Arc<crate::http::Apis>,
generator: Arc<Generator>,
filter: crate::feature::TagFilter,
db: Option<Arc<crate::db::Db>>,
default_db: String,
srp: Option<Arc<crate::srp::SrpParams>>,
plugins: Option<Arc<crate::plugin::Plugins>>,
options: Options,
fail_fast: bool,
) -> Self {
Self {
reg,
apis,
generator,
filter,
db,
default_db,
srp,
plugins,
options,
fail_fast,
stop: std::sync::atomic::AtomicBool::new(false),
}
}
pub fn request_stop(&self) {
if self.fail_fast {
self.stop.store(true, Ordering::Relaxed);
}
}
pub fn stopped(&self) -> bool {
self.stop.load(Ordering::Relaxed)
}
}
pub async fn run_file(lf: Arc<LoadedFeature>, ctx: Arc<RunContext>) -> FileResult {
let mut world = World::new(
ctx.apis.clone(),
ctx.generator.clone(),
crate::db::DbHandle::new(ctx.db.clone(), ctx.default_db.clone()),
ctx.srp.clone(),
ctx.plugins.clone(),
ctx.options.clone(),
);
let mut scenarios = Vec::new();
let background: Vec<ExpandedStep> = lf
.feature
.background
.as_ref()
.map(|bg| bg.steps.iter().map(crate::feature::to_step).collect())
.unwrap_or_default();
let generator = world.generator.clone();
for sc in &lf.feature.scenarios {
if !ctx.filter.matches(&sc.tags) {
continue;
}
for ex in expand_outlines(sc) {
world.reset_scenario();
let mut failure = match ctx.plugins.clone() {
None => None,
Some(plugins) => {
let handles = world.plugins.to_reset();
if handles.is_empty() {
None
} else {
match tokio::task::spawn_blocking(move || plugins.reset_instances(&handles))
.await
{
Ok(result) => result
.err()
.map(|error| format!(" scenario reset\n{error}")),
Err(error) => Some(format!(
" scenario reset\nthe plugin reset task failed: {error}"
)),
}
}
}
};
let started = Instant::now();
let mut steps = Vec::with_capacity(background.len() + ex.steps.len());
for step in background.iter().chain(ex.steps.iter()) {
let step_started = Instant::now();
let status = if failure.is_some() {
StepStatus::Skipped
} else if let Err(e) =
execute_step(&mut world, &ctx.reg, step, &lf.path, &generator, 0).await
{
let mut msg = format!(" {}\n{e}", step.text);
if let Some(ex) = world.http.last() {
msg.push_str(&format!("\n\n{ex}"));
}
failure = Some(msg);
StepStatus::Failed
} else {
StepStatus::Passed
};
steps.push(StepResult {
keyword: step.keyword.clone(),
text: step.text.clone(),
line: step.line,
status,
duration: step_started.elapsed(),
});
}
scenarios.push(ScenarioResult {
name: ex.name,
line: ex.line,
failure,
steps,
duration: started.elapsed(),
});
}
}
if let Some(plugins) = ctx.plugins.clone() {
let handles = world.plugins.to_drop();
if !handles.is_empty() {
let _ = tokio::task::spawn_blocking(move || plugins.drop_instances(&handles)).await;
}
}
FileResult {
path: lf.path.clone(),
name: lf.feature.name.clone(),
scenarios,
}
}
fn worker_count(concurrency: usize, units: usize) -> usize {
concurrency.clamp(1, units.max(1))
}
fn panicked_file(lf: &LoadedFeature, error: &tokio::task::JoinError) -> FileResult {
FileResult {
path: lf.path.clone(),
name: lf.feature.name.clone(),
scenarios: vec![ScenarioResult {
name: "file run aborted".to_string(),
line: 0,
failure: Some(format!("panic while running the file: {error}")),
steps: Vec::new(),
duration: std::time::Duration::ZERO,
}],
}
}
#[derive(Debug)]
pub struct Chain {
pub priority: i64,
pub files: Vec<Arc<LoadedFeature>>,
}
pub fn build_chains(files: Vec<Arc<LoadedFeature>>) -> Result<Vec<Chain>, String> {
let mut named: BTreeMap<String, Vec<(i64, Arc<LoadedFeature>)>> = BTreeMap::new();
let mut chains: Vec<Chain> = Vec::new();
for lf in files {
let priority = crate::feature::priority_of(&lf)?;
match crate::feature::serial_of(&lf)? {
Some(name) => named.entry(name).or_default().push((priority, lf)),
None => chains.push(Chain {
priority,
files: vec![lf],
}),
}
}
for (_, mut members) in named {
members.sort_by_key(|(p, _)| std::cmp::Reverse(*p));
let priority = members.iter().map(|(p, _)| *p).max().unwrap_or(0);
chains.push(Chain {
priority,
files: members.into_iter().map(|(_, lf)| lf).collect(),
});
}
chains.sort_by(|a, b| {
b.priority.cmp(&a.priority).then_with(|| {
a.files
.first()
.map(|f| &f.path)
.cmp(&b.files.first().map(|f| &f.path))
})
});
Ok(chains)
}
pub async fn run_all(
chains: Vec<Chain>,
ctx: Arc<RunContext>,
concurrency: usize,
) -> Vec<FileResult> {
let workers = worker_count(concurrency, chains.len());
let total: usize = chains.iter().map(|c| c.files.len()).sum();
let chains = Arc::new(chains);
let cursor = Arc::new(AtomicUsize::new(0));
let results = Arc::new(Mutex::new(Vec::with_capacity(total)));
let mut handles = Vec::with_capacity(workers);
for _ in 0..workers {
let (chains, cursor, results, ctx) =
(chains.clone(), cursor.clone(), results.clone(), ctx.clone());
handles.push(tokio::spawn(async move {
loop {
if ctx.stopped() {
break;
}
let index = cursor.fetch_add(1, Ordering::Relaxed);
let Some(chain) = chains.get(index) else {
break;
};
for lf in &chain.files {
if ctx.stopped() {
break;
}
let task = tokio::spawn(run_file(lf.clone(), ctx.clone()));
let result = match task.await {
Ok(r) => r,
Err(error) => panicked_file(lf, &error),
};
if result.failed() > 0 {
ctx.request_stop();
}
let mut collected = results.lock().expect("results mutex");
print!("{}", render_file(&result));
collected.push(result);
}
}
}));
}
for handle in handles {
let _ = handle.await;
}
std::mem::take(&mut *results.lock().expect("results mutex"))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::feature::parse_str;
use crate::macros::MacroCatalog;
use crate::unique::UniqueKind;
use std::path::PathBuf;
fn apis() -> Arc<crate::http::Apis> {
apis_at("http://example.test")
}
fn apis_at(base: &str) -> Arc<crate::http::Apis> {
let mut by_name = std::collections::HashMap::new();
by_name.insert(
"default".to_string(),
crate::http::ApiResource::new(base, 1, Vec::new(), Options::default()).unwrap(),
);
Arc::new(crate::http::Apis::new(by_name, Some("default".to_string())).unwrap())
}
fn context_with(
registry: Registry,
apis: Arc<crate::http::Apis>,
generator: Arc<Generator>,
plugins: Option<Arc<crate::plugin::Plugins>>,
) -> Arc<RunContext> {
Arc::new(RunContext::new(
registry,
apis,
generator,
crate::feature::TagFilter::new(&[]),
None,
String::new(),
None,
plugins,
Options::default(),
false,
))
}
fn context(registry: Registry) -> Arc<RunContext> {
context_with(registry, apis(), Arc::new(Generator::new()), None)
}
fn echo_plugins() -> crate::plugin::Plugins {
let mut plugins = crate::plugin::Plugins::load(
vec![crate::plugin::tests::entry()],
&[crate::plugin::tests::instance("a", Some("p-"))],
&["echo".to_string()],
1,
&crate::options::Options::default(),
)
.expect("the fixture plugin loads");
plugins.add_defaults(
[("echo".to_string(), "a".to_string())]
.into_iter()
.collect(),
);
plugins
}
fn echo_registry(plugins: &crate::plugin::Plugins, force_action: Option<usize>) -> Registry {
let steps: Vec<(usize, usize, String, bool)> = plugins
.steps()
.into_iter()
.map(|(lib, index, pattern, assertion)| {
(
lib,
index,
pattern,
assertion && force_action != Some(index),
)
})
.collect();
Registry::with_macros_and_plugins(MacroCatalog::default(), &steps, &["echo".to_string()])
.expect("the plugin registry builds")
}
fn echo_world(plugins: crate::plugin::Plugins) -> (World, Arc<crate::plugin::Plugins>) {
let plugins = Arc::new(plugins);
let world = World::new(
apis(),
Arc::new(Generator::new()),
crate::db::DbHandle::new(None, String::new()),
None,
Some(plugins.clone()),
Options::default(),
);
(world, plugins)
}
async fn run_step(world: &mut World, reg: &Registry, text: &str) -> Result<(), String> {
let step = ExpandedStep {
keyword: String::new(),
text: text.to_string(),
line: 1,
docstring: None,
table: None,
};
let generator = world.generator.clone();
let source = std::path::Path::new("test.feature");
execute_step(world, reg, &step, source, &generator, 0).await
}
#[tokio::test]
async fn a_plugin_step_records_its_handle_on_the_world() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, None);
let (mut world, plugins) = echo_world(plugins);
run_step(&mut world, ®, r#"I echo "x" as "first""#)
.await
.expect("first dispatch");
let recorded = world
.plugins
.handle_for("echo", "a")
.expect("the first dispatch recorded its handle");
run_step(&mut world, ®, r#"I echo "y" as "second""#)
.await
.expect("second dispatch");
assert_eq!(
world.plugins.handle_for("echo", "a"),
Some(recorded),
"the second dispatch reused the recorded handle"
);
assert_eq!(
plugins.registered_count(),
1,
"and created no second instance"
);
}
#[tokio::test]
async fn a_plugin_step_gets_interpolated_arguments_and_publishes_its_vars() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, None);
let (mut world, _) = echo_world(plugins);
world.vars.set("who", "world".to_string());
run_step(&mut world, ®, r#"I echo "<<who>>" as "greeting""#)
.await
.expect("the action passes");
assert_eq!(world.vars.get("greeting"), Some("p-world"));
}
#[tokio::test]
async fn an_action_answering_not_yet_is_a_contract_error() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, Some(1));
let (mut world, _) = echo_world(plugins);
let error = run_step(&mut world, ®, "the echo counter should reach 3")
.await
.expect_err("an action may not answer not_yet");
assert!(error.contains("only an assertion may do"), "{error}");
assert!(error.contains("counter is 1 of 3"), "{error}");
}
#[tokio::test]
async fn a_not_yet_without_an_armed_modifier_fails_and_discards_its_vars() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, None);
let (mut world, _) = echo_world(plugins);
let error = run_step(&mut world, ®, "the echo counter should reach 3")
.await
.expect_err("one attempt, and it did not pass");
assert!(error.contains("counter is 1 of 3"), "{error}");
assert_eq!(world.vars.get("echo_attempts"), None);
}
#[tokio::test]
async fn an_armed_modifier_polls_the_plugin_until_it_passes() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, None);
let (mut world, _) = echo_world(plugins);
world.arm_options(
serde_yaml_ng::from_str("polling:\n timeout_secs: 5\n interval_ms: 1\n")
.expect("the layer parses"),
);
run_step(&mut world, ®, "the echo counter should reach 3")
.await
.expect("the third attempt passes");
assert_eq!(
world.vars.get("echo_attempts"),
Some("3"),
"the retry loop reached the third attempt"
);
}
#[tokio::test]
async fn a_plugin_action_leaves_an_armed_modifier_for_the_next_assertion() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, None);
let (mut world, _) = echo_world(plugins);
world.arm_options(
serde_yaml_ng::from_str("polling:\n timeout_secs: 5\n interval_ms: 1\n")
.expect("the layer parses"),
);
run_step(&mut world, ®, r#"I echo "x" as "ignored""#)
.await
.expect("the action passes");
run_step(&mut world, ®, "the echo counter should reach 3")
.await
.expect("the modifier survived the action");
}
#[tokio::test]
async fn a_fatal_reply_carries_its_message_and_its_diagnostics() {
let plugins = echo_plugins();
let reg = echo_registry(&plugins, None);
let (mut world, _) = echo_world(plugins);
let error = run_step(&mut world, ®, "the echo should fail")
.await
.expect_err("the step is asked to fail");
assert!(error.contains("the echo step was asked to fail"), "{error}");
assert!(error.contains("echo state"), "{error}");
assert!(error.contains("prefix=p-"), "{error}");
}
fn macro_catalog(name: &str, source: &str) -> MacroCatalog {
let dir = std::env::temp_dir().join(format!("bddkit-runner-{}-{name}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("macros.yaml");
std::fs::write(&path, source).unwrap();
MacroCatalog::load(&[path]).unwrap()
}
fn registry(name: &str, source: &str) -> Registry {
Registry::with_macros(macro_catalog(name, source)).unwrap()
}
async fn run(feature: &str, registry: Registry) -> FileResult {
run_in(feature, context(registry)).await
}
async fn run_in(feature: &str, context: Arc<RunContext>) -> FileResult {
let loaded = LoadedFeature {
path: PathBuf::from("macro.feature"),
feature: parse_str(feature).unwrap(),
};
run_file(Arc::new(loaded), context).await
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_file_run_can_be_spawned_onto_another_thread() {
let loaded = LoadedFeature {
path: PathBuf::from("spawnable.feature"),
feature: parse_str(
"Feature: f\n Scenario: s\n Given set variable \"a\" to \"1\"\n",
)
.expect("gherkin parses"),
};
let ctx = context(Registry::new().expect("built-in patterns compile"));
let result = tokio::spawn(run_file(Arc::new(loaded), ctx))
.await
.expect("the task must not panic");
assert_eq!(result.failed(), 0, "{:?}", result.scenarios[0].failure);
}
#[tokio::test]
async fn macro_interpolates_parameter_and_exports_result() {
let registry = registry(
"export",
r#"
- step: 'I remember "{value}"'
exports: [result]
do:
- set variable "private" to "<<value>>"
- set variable "result" to "<<private>>"
"#,
);
let result = run(
r#"
Feature: macro
Scenario: export
Given set variable "outer" to "saved"
When I remember "<<outer>>"
Then variable "result" should be equal to "saved"
"#,
registry,
)
.await;
assert!(
result.scenarios[0].failure.is_none(),
"{:?}",
result.scenarios[0].failure
);
}
#[tokio::test]
async fn macro_drops_private_variables() {
let registry = registry(
"private",
r#"
- step: I create private state
do:
- set variable "private" to "hidden"
"#,
);
let result = run(
r#"
Feature: macro
Scenario: private
When I create private state
Then variable "private" should be equal to "hidden"
"#,
registry,
)
.await;
let failure = result.scenarios[0].failure.as_deref().unwrap();
assert!(
failure.contains("private") && failure.contains("is not set"),
"{failure}"
);
}
#[tokio::test]
async fn macro_exports_matching_glob() {
let registry = registry(
"glob",
r#"
- step: I create a row
exports: [last_insert_id_*]
do:
- set variable "last_insert_id_users" to "42"
"#,
);
let result = run(
r#"
Feature: macro
Scenario: glob
When I create a row
Then variable "last_insert_id_users" should be equal to "42"
"#,
registry,
)
.await;
assert!(
result.scenarios[0].failure.is_none(),
"{:?}",
result.scenarios[0].failure
);
}
#[tokio::test]
async fn nested_macro_exports_through_each_frame() {
let registry = registry(
"nested-export",
r#"
- step: 'I make inner "{value}"'
exports: [inner]
do:
- set variable "inner" to "<<value>>"
- step: 'I make outer "{value}"'
exports: [result]
do:
- I make inner "<<value>>"
- set variable "result" to "<<inner>>"
"#,
);
let result = run(
r#"
Feature: macro
Scenario: nested
When I make outer "done"
Then variable "result" should be equal to "done"
"#,
registry,
)
.await;
assert!(
result.scenarios[0].failure.is_none(),
"{:?}",
result.scenarios[0].failure
);
}
#[tokio::test]
async fn missing_declared_export_fails_scenario() {
let registry = registry(
"missing-export",
r#"
- step: I forget the result
exports: [result]
do:
- set variable "private" to "x"
"#,
);
let result = run(
"Feature: macro\n Scenario: missing\n When I forget the result\n",
registry,
)
.await;
let failure = result.scenarios[0].failure.as_deref().unwrap();
assert!(
failure.contains("result") && failure.contains("is not set"),
"{failure}"
);
}
#[tokio::test]
async fn macro_call_rejects_docstring_argument() {
let registry = registry(
"call-docstring",
"- step: I do business\n do: [the response code is 200]\n",
);
let result = run(
r#"
Feature: macro
Scenario: docstring
When I do business
"""
unsupported
"""
"#,
registry,
)
.await;
let failure = result.scenarios[0].failure.as_deref().unwrap();
assert!(failure.contains("docstring"), "{failure}");
}
#[tokio::test]
async fn macro_call_rejects_table_argument() {
let registry = registry(
"call-table",
"- step: I do business\n do: [the response code is 200]\n",
);
let result = run(
r#"
Feature: macro
Scenario: table
When I do business
| value |
| x |
"#,
registry,
)
.await;
let failure = result.scenarios[0].failure.as_deref().unwrap();
assert!(failure.contains("table"), "{failure}");
}
mod eventual {
use super::*;
use tokio::time::Instant;
async fn failure(feature: &str, registry: Registry) -> String {
run(feature, registry).await.scenarios[0]
.failure
.clone()
.expect("scenario should fail")
}
async fn run_with_echo(feature: &str, macros: MacroCatalog) -> FileResult {
let plugins = echo_plugins();
let registry =
Registry::with_macros_and_plugins(macros, &plugins.steps(), &["echo".to_string()])
.expect("the plugin registry builds");
let context = context_with(
registry,
apis(),
Arc::new(Generator::new()),
Some(Arc::new(plugins)),
);
run_in(feature, context).await
}
async fn echo_failure(feature: &str, macros: MacroCatalog) -> String {
run_with_echo(feature, macros).await.scenarios[0]
.failure
.clone()
.expect("scenario should fail")
}
#[tokio::test(start_paused = true)]
async fn a_counter_that_never_arrives_times_out_with_the_last_observation() {
let failure = echo_failure(
r#"
Feature: eventual assertion
Scenario: timeout
Given I expect the next assertion to pass within "1" seconds, checking every "100" milliseconds
Then the echo counter should reach 1000
"#,
MacroCatalog::default(),
)
.await;
for expected in [
"the echo counter should reach 1000",
"1s",
"100ms",
"11 attempts",
"counter is 11 of 1000",
] {
assert!(
failure.contains(expected),
"missing {expected:?} in {failure}"
);
}
}
#[tokio::test(start_paused = true)]
async fn a_variable_mismatch_under_a_modifier_fails_at_once() {
let started = Instant::now();
let failure = failure(
r#"
Feature: eventual assertion
Scenario: final
Given set variable "state" to "pending"
And I expect the next assertion to pass within "1" seconds, checking every "100" milliseconds
Then variable "state" should be equal to "ready"
"#,
Registry::new().unwrap(),
)
.await;
assert!(!failure.contains("did not pass"), "{failure}");
assert!(failure.contains("expected: ready"), "{failure}");
assert!(failure.contains("actual: pending"), "{failure}");
assert_eq!(Instant::now(), started);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_retrying_assertion_prepares_unique_capture_once() {
let app =
axum::Router::new().route("/state", axum::routing::get(|| async { "pending" }));
let (base, _server) = crate::http::tests::spawn_app(app).await;
let generator = Arc::new(Generator::new());
let context = context_with(
Registry::new().unwrap(),
apis_at(&base),
generator.clone(),
None,
);
let before = generator.next(UniqueKind::Number).parse::<u64>().unwrap();
let result = run_in(
r#"
Feature: eventual assertion
Scenario: stable preparation
When I request "/state"
And I expect the next assertion to pass within "1" seconds, checking every "100" milliseconds
Then the response body contains "<<unique(number)>>"
"#,
context,
)
.await;
let failure = result.scenarios[0].failure.as_deref().unwrap();
let attempts: u64 = regex::Regex::new(r"(\d+) attempts")
.unwrap()
.captures(failure)
.and_then(|caps| caps[1].parse().ok())
.unwrap_or_else(|| panic!("no attempt count in {failure}"));
assert!(attempts > 1, "the assertion was never retried: {failure}");
let after = generator.next(UniqueKind::Number).parse::<u64>().unwrap();
assert_eq!(after - before, 2, "unique() must be evaluated only once");
}
#[tokio::test(start_paused = true)]
async fn a_missing_variable_is_fatal_without_waiting() {
let started = Instant::now();
let failure = failure(
r#"
Feature: eventual assertion
Scenario: fatal
Given I expect the next assertion to pass within "1" seconds, checking every "100" milliseconds
Then variable "missing" should be equal to "ready"
"#,
Registry::new().unwrap(),
)
.await;
assert!(
failure.contains("missing") && failure.contains("is not set"),
"{failure}"
);
assert_eq!(Instant::now(), started);
}
#[tokio::test(start_paused = true)]
async fn a_later_modifier_silently_replaces_the_first() {
let failure = echo_failure(
r#"
Feature: eventual assertion
Scenario: replacement
Given I expect the next assertion to pass within "5" seconds
And I expect the next assertion to pass within "1" seconds
Then the echo counter should reach 1000
"#,
MacroCatalog::default(),
)
.await;
assert!(failure.contains("within 1s"), "{failure}");
assert!(!failure.contains("within 5s"), "{failure}");
}
#[tokio::test(start_paused = true)]
async fn a_modifier_survives_an_ordinary_action() {
let failure = echo_failure(
r#"
Feature: eventual assertion
Scenario: action
Given I expect the next assertion to pass within "1" seconds
And set variable "state" to "pending"
Then the echo counter should reach 1000
"#,
MacroCatalog::default(),
)
.await;
assert!(failure.contains("did not pass within 1s"), "{failure}");
}
#[tokio::test(start_paused = true)]
async fn a_modifier_survives_a_macro_boundary() {
let macros = macro_catalog(
"eventual",
r#"
- step: I check the counter through a macro
do:
- set variable "macro_ran" to "yes"
- the echo counter should reach 1000
"#,
);
let failure = echo_failure(
r#"
Feature: eventual assertion
Scenario: macro
Given I expect the next assertion to pass within "1" seconds
Then I check the counter through a macro
"#,
macros,
)
.await;
assert!(failure.contains("did not pass within 1s"), "{failure}");
}
#[tokio::test(start_paused = true)]
async fn the_first_assertion_consumes_the_modifier() {
let started = Instant::now();
let failure = echo_failure(
r#"
Feature: eventual assertion
Scenario: one shot
Given set variable "state" to "ready"
And I expect the next assertion to pass within "1" seconds
Then variable "state" should be equal to "ready"
And the echo counter should reach 1000
"#,
MacroCatalog::default(),
)
.await;
assert!(!failure.contains("did not pass"), "{failure}");
assert!(failure.contains("counter is 1 of 1000"), "{failure}");
assert_eq!(Instant::now(), started);
}
#[tokio::test(start_paused = true)]
async fn a_dangling_modifier_passes_silently() {
let result = run(
r#"
Feature: eventual assertion
Scenario: dangling
Given I expect the next assertion to pass eventually
"#,
Registry::new().unwrap(),
)
.await;
assert_eq!(result.failed(), 0, "{:?}", result.scenarios[0].failure);
}
#[tokio::test(start_paused = true)]
async fn a_scenario_reset_discards_a_dangling_modifier() {
let started = Instant::now();
let result = run_with_echo(
r#"
Feature: eventual assertion
Scenario: arm
Given I expect the next assertion to pass within "1" seconds
Scenario: assert
Then the echo counter should reach 1000
"#,
MacroCatalog::default(),
)
.await;
assert!(result.scenarios[0].failure.is_none());
let failure = result.scenarios[1]
.failure
.as_deref()
.expect("second scenario fails");
assert!(!failure.contains("did not pass"), "{failure}");
assert!(failure.contains("counter is 1 of 1000"), "{failure}");
assert_eq!(Instant::now(), started);
}
#[tokio::test(start_paused = true)]
async fn zero_polling_values_are_rejected_when_arming() {
let failure = failure(
r#"
Feature: eventual assertion
Scenario: zero
Given I expect the next assertion to pass within "0" seconds
"#,
Registry::new().unwrap(),
)
.await;
assert!(failure.contains("positive"), "{failure}");
}
#[tokio::test(start_paused = true)]
async fn an_explicit_interval_cannot_exceed_the_explicit_timeout() {
let failure = failure(
r#"
Feature: eventual assertion
Scenario: invalid interval
Given I expect the next assertion to pass within "1" seconds, checking every "1001" milliseconds
"#,
Registry::new().unwrap(),
)
.await;
assert!(failure.contains("must not exceed"), "{failure}");
}
}
#[test]
fn a_zero_concurrency_config_still_runs_with_one_worker() {
assert_eq!(worker_count(0, 5), 1);
}
#[test]
fn workers_never_outnumber_the_work() {
assert_eq!(worker_count(8, 3), 3);
}
#[test]
fn a_concurrency_below_the_unit_count_is_respected() {
assert_eq!(worker_count(4, 10), 4);
}
mod chains {
use super::super::*;
use crate::feature::parse_str;
use std::path::PathBuf;
use std::sync::Arc;
pub(super) fn file(path: &str, tags: &str) -> Arc<LoadedFeature> {
let src =
format!("{tags}Feature: f\n Scenario: s\n Then the response code is 200\n");
Arc::new(LoadedFeature {
path: PathBuf::from(path),
feature: parse_str(&src).expect("gherkin parses"),
})
}
#[test]
fn untagged_files_become_one_chain_each() {
let chains = build_chains(vec![file("a.feature", ""), file("b.feature", "")])
.expect("no tags — no errors");
assert_eq!(chains.len(), 2);
assert!(chains.iter().all(|c| c.files.len() == 1));
}
#[test]
fn files_sharing_a_name_end_up_in_one_chain() {
let chains = build_chains(vec![
file("a.feature", "@serial(x)\n"),
file("b.feature", ""),
file("c.feature", "@serial(x)\n"),
])
.expect("tags are valid");
assert_eq!(chains.len(), 2, "two chains: x and the lone b");
let x = chains
.iter()
.find(|c| c.files.len() == 2)
.expect("chain x exists");
assert_eq!(x.files[0].path, PathBuf::from("a.feature"));
assert_eq!(x.files[1].path, PathBuf::from("c.feature"));
}
#[test]
fn different_names_stay_in_different_chains() {
let chains = build_chains(vec![
file("a.feature", "@serial(x)\n"),
file("b.feature", "@serial(y)\n"),
])
.expect("tags are valid");
assert_eq!(chains.len(), 2);
}
#[test]
fn chain_order_is_reproducible_across_runs() {
let build = || {
build_chains(vec![
file("b.feature", "@serial(zeta)\n"),
file("a.feature", ""),
file("c.feature", "@serial(alpha)\n"),
])
.expect("tags are valid")
.into_iter()
.map(|c| c.files[0].path.clone())
.collect::<Vec<_>>()
};
assert_eq!(build(), build());
}
#[test]
fn a_broken_serial_tag_stops_chain_building() {
let err = build_chains(vec![file("a.feature", "@serial()\n")])
.expect_err("an empty chain name is an error");
assert!(err.contains("a.feature"), "{err}");
}
#[test]
fn higher_priority_chains_come_first() {
let chains = build_chains(vec![
file("a.feature", ""),
file("b.feature", "@priority(-1)\n"),
file("c.feature", "@priority(5)\n"),
])
.expect("tags are valid");
let order: Vec<_> = chains.iter().map(|c| c.files[0].path.clone()).collect();
assert_eq!(
order,
vec![
PathBuf::from("c.feature"),
PathBuf::from("a.feature"),
PathBuf::from("b.feature"),
]
);
}
#[test]
fn a_chain_takes_the_highest_priority_of_its_files() {
let chains = build_chains(vec![
file("a.feature", "@serial(x)\n"),
file("b.feature", "@serial(x)\n@priority(7)\n"),
file("c.feature", "@priority(3)\n"),
])
.expect("tags are valid");
assert_eq!(chains[0].priority, 7, "chain x must go first");
assert_eq!(chains[0].files.len(), 2);
}
#[test]
fn files_inside_a_chain_are_ordered_by_priority_too() {
let chains = build_chains(vec![
file("a.feature", "@serial(x)\n"),
file("b.feature", "@serial(x)\n@priority(7)\n"),
])
.expect("tags are valid");
assert_eq!(chains[0].files[0].path, PathBuf::from("b.feature"));
}
#[test]
fn equal_priorities_keep_alphabetical_order() {
let chains = build_chains(vec![
file("b.feature", "@priority(1)\n"),
file("a.feature", "@priority(1)\n"),
])
.expect("tags are valid");
assert_eq!(chains[0].files[0].path, PathBuf::from("a.feature"));
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_panicking_file_becomes_a_failed_result_instead_of_killing_the_run() {
let error = tokio::spawn(async { panic!("deliberate") })
.await
.expect_err("the task must panic");
let result = panicked_file(&chains::file("bad.feature", ""), &error);
assert_eq!(result.failed(), 1);
assert!(
result.scenarios[0]
.failure
.as_deref()
.is_some_and(|f| f.contains("panic")),
"{:?}",
result.scenarios[0].failure
);
}
}