use super::*;
use std::collections::HashMap;
use bevy_app::{App, Update};
use bevy_asset::AssetPlugin;
use bevy_ecs::event::Event;
use bevy_ecs::observer::On;
use bevy_ecs::resource::Resource;
use bevy_ecs::system::{ResMut, RunSystemOnce as _};
use brink_runtime::Step;
use super::parallel::advance_batch_parallel;
use crate::test_support::{add_story_assets, compile_test_story};
use crate::{
BrinkChoicesPresented, BrinkFlowRequest, BrinkLineDelivered, BrinkStoryEnded, BrinkTurnDone,
};
fn fixture(source: &str) -> (Program, Vec<Vec<LineEntry>>, World) {
compile_test_story(source)
}
fn dump_globals(program: &Program, world: &World) -> Vec<Value> {
(0..program.global_count())
.map(|idx| world.global(idx).clone())
.collect()
}
fn run_batch_permuted(
program: &Program,
tables: &[Vec<LineEntry>],
frame_start: &World,
n: usize,
step_order: &[usize],
) -> (World, Vec<String>) {
let mut flows: Vec<FlowInstance> = (0..n)
.map(|_| FlowInstance::new_at_root(program).0)
.collect();
let mut bufs: Vec<WriteBuffer> = (0..n).map(|_| WriteBuffer::default()).collect();
let mut rendered = vec![String::new(); n];
for &i in step_order {
let (lines, _awaiting, _error) = step_flow(
frame_start,
&mut flows[i],
program,
tables,
&FallbackHandler,
&mut bufs[i],
);
rendered[i] = lines.iter().map(Step::text).collect();
}
let mut world = frame_start.clone();
for buf in &bufs {
buf.apply_to(&mut world);
}
(world, rendered)
}
#[test]
fn reads_pin_to_frame_start_not_a_peers_buffered_write() {
let (program, tables, frame_start) = fixture("VAR g = 0\nVal {g}.\n~ g = g + 1\n-> END\n");
let mut flow_a = FlowInstance::new_at_root(&program).0;
let mut buf_a = WriteBuffer::default();
let (lines_a, _, _) = step_flow(
&frame_start,
&mut flow_a,
&program,
&tables,
&FallbackHandler,
&mut buf_a,
);
let text_a: String = lines_a.iter().map(Step::text).collect();
assert!(
text_a.contains("Val 0."),
"A reads frame-start g=0: {text_a:?}"
);
assert!(buf_a.len() >= 1, "A should have buffered a write to g");
assert_eq!(
frame_start.global(program.global_index("g").expect("g exists")),
&Value::Int(0),
"frame-start world must be untouched by A's buffered write"
);
let mut flow_b = FlowInstance::new_at_root(&program).0;
let mut buf_b = WriteBuffer::default();
let (lines_b, _, _) = step_flow(
&frame_start,
&mut flow_b,
&program,
&tables,
&FallbackHandler,
&mut buf_b,
);
let text_b: String = lines_b.iter().map(Step::text).collect();
assert!(
text_b.contains("Val 0."),
"B must read frame-start g=0, not A's write: {text_b:?}"
);
}
#[test]
fn order_invariance_identical_flows() {
let (program, tables, frame_start) = fixture("VAR g = 0\nVal {g}.\n~ g = g + 1\n-> END\n");
let n = 5;
let permutations: &[&[usize]] = &[
&[0, 1, 2, 3, 4],
&[4, 3, 2, 1, 0],
&[2, 0, 4, 1, 3],
&[1, 3, 0, 4, 2],
&[3, 2, 1, 4, 0],
];
let mut baseline: Option<(Vec<Value>, Vec<String>)> = None;
for perm in permutations {
let (world, rendered) = run_batch_permuted(&program, &tables, &frame_start, n, perm);
let globals = dump_globals(&program, &world);
for (i, text) in rendered.iter().enumerate() {
assert!(
text.contains("Val 0."),
"flow {i} under perm {perm:?} must read frame-start g=0: {text:?}"
);
}
assert_eq!(
world.global(program.global_index("g").expect("g exists")),
&Value::Int(1),
"converged g must be 1 under perm {perm:?}"
);
match &baseline {
None => baseline = Some((globals, rendered)),
Some((base_globals, base_rendered)) => {
assert_eq!(
&globals, base_globals,
"converged world differs under step permutation {perm:?}"
);
assert_eq!(
&rendered, base_rendered,
"per-flow outcomes differ under step permutation {perm:?}"
);
}
}
}
}
#[test]
fn write_write_resolves_by_flow_id_not_step_order() {
let low = fixture("VAR g = 0\n~ g = 10\n-> END\n");
let high = fixture("VAR g = 0\n~ g = 20\n-> END\n");
let g_idx = low.0.global_index("g").expect("g exists");
assert_eq!(g_idx, high.0.global_index("g").expect("g exists"));
let run = |step_first_high: bool| -> Value {
let frame_start = low.2.clone();
let mut flow_low = FlowInstance::new_at_root(&low.0).0;
let mut flow_high = FlowInstance::new_at_root(&high.0).0;
let mut buf_low = WriteBuffer::default();
let mut buf_high = WriteBuffer::default();
if step_first_high {
step_flow(
&frame_start,
&mut flow_high,
&high.0,
&high.1,
&FallbackHandler,
&mut buf_high,
);
step_flow(
&frame_start,
&mut flow_low,
&low.0,
&low.1,
&FallbackHandler,
&mut buf_low,
);
} else {
step_flow(
&frame_start,
&mut flow_low,
&low.0,
&low.1,
&FallbackHandler,
&mut buf_low,
);
step_flow(
&frame_start,
&mut flow_high,
&high.0,
&high.1,
&FallbackHandler,
&mut buf_high,
);
}
let mut world = frame_start.clone();
buf_low.apply_to(&mut world);
buf_high.apply_to(&mut world);
world.global(g_idx).clone()
};
assert_eq!(run(false), Value::Int(20));
assert_eq!(run(true), Value::Int(20));
}
#[test]
fn write_buffer_captures_and_replays_changeset() {
let (program, tables, frame_start) = fixture("VAR g = 0\n~ g = 7\n-> END\n");
let mut flow = FlowInstance::new_at_root(&program).0;
let mut buf = WriteBuffer::default();
step_flow(
&frame_start,
&mut flow,
&program,
&tables,
&FallbackHandler,
&mut buf,
);
assert!(buf.len() >= 1);
let mut world = frame_start.clone();
buf.apply_to(&mut world);
assert_eq!(
world.global(program.global_index("g").expect("g")),
&Value::Int(7)
);
}
#[derive(Default)]
struct Batched;
#[test]
fn advance_batch_system_drives_and_applies() {
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
app.add_systems(Update, advance_batch::<Batched>);
let (program, tables, ctx) = compile_test_story("VAR g = 0\nHi.\n~ g = 3\n-> END\n");
let story = add_story_assets(&mut app, program, tables, ctx);
for _ in 0..2 {
app.world_mut().spawn(
BrinkFlowRequest::<Batched>::builder()
.story(story.clone())
.build(),
);
}
app.update();
app.update();
let report = app.world().resource::<BrinkBatchReport<Batched>>();
assert_eq!(report.stepped, 2, "both flows stepped to terminal");
assert_eq!(report.awaiting, 0);
assert!(
report.writes_applied >= 2,
"each flow buffered a write to g; got {}",
report.writes_applied
);
let globals = app.world().resource::<BrinkGlobals<Batched>>();
assert_eq!(globals.inner.global(0), &Value::Int(3));
}
#[test]
fn step_flow_returns_the_runtime_error_on_fault() {
let (program, tables, frame_start) = fixture("-> spam\n\n=== spam ===\nLine.\n-> spam\n");
let mut flow = FlowInstance::new_at_root(&program).0;
let mut buf = WriteBuffer::default();
let (lines, awaiting, error) = step_flow(
&frame_start,
&mut flow,
&program,
&tables,
&FallbackHandler,
&mut buf,
);
assert!(lines.is_empty(), "a faulted step produces no lines");
assert!(!awaiting, "a faulted step is not a deferred-external park");
match error {
Some(brink_runtime::RuntimeError::LineLimitExceeded(n)) => {
assert_eq!(n, FlowInstance::LINE_LIMIT);
}
other => panic!("expected Some(LineLimitExceeded), got {other:?}"),
}
}
#[test]
fn errored_flow_is_reported_distinctly_not_counted_as_stepped() {
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
app.add_systems(Update, advance_batch::<Batched>);
let (program, tables, ctx) = compile_test_story("-> spam\n\n=== spam ===\nLine.\n-> spam\n");
let story = add_story_assets(&mut app, program, tables, ctx);
app.world_mut().spawn(
BrinkFlowRequest::<Batched>::builder()
.story(story.clone())
.build(),
);
app.update(); app.update();
let report = app.world().resource::<BrinkBatchReport<Batched>>();
assert_eq!(
report.errored, 1,
"the faulted flow must be counted as errored"
);
assert_eq!(
report.stepped, 0,
"a faulted flow must not be counted as stepped"
);
assert_eq!(
report.awaiting, 0,
"a faulted flow is not a deferred-external park"
);
assert_eq!(report.flows.len(), 1);
assert!(
report.flows[0].errored,
"the flow's own FlowAccessRecord must also flag errored"
);
}
#[test]
fn command_triggers_flush_in_flow_id_order() {
use crate::{BrinkBindingsAppExt, BrinkCommand, Value as BValue};
#[derive(Event)]
struct Ping {
who: i32,
}
impl BrinkCommand for Ping {
fn from_ink_args(args: &[BValue]) -> Result<Self, crate::BrinkArgError> {
Ok(Ping {
who: args.first().and_then(BValue::as_int).unwrap_or(-1),
})
}
}
#[derive(Resource, Default)]
struct PingLog(Vec<i32>);
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
app.init_resource::<PingLog>();
app.bind_brink_command::<Batched, Ping>("ping");
app.add_observer(|on: On<Ping>, mut log: ResMut<PingLog>| {
log.0.push(on.event().who);
});
app.add_systems(Update, advance_batch::<Batched>);
let (p0, t0, c0) = compile_test_story("EXTERNAL ping(n)\n~ temp _ = ping(10)\nHi.\n-> END\n");
let (p1, t1, c1) = compile_test_story("EXTERNAL ping(n)\n~ temp _ = ping(20)\nHi.\n-> END\n");
let story0 = add_story_assets(&mut app, p0, t0, c0);
let story1 = add_story_assets(&mut app, p1, t1, c1);
let e0 = app
.world_mut()
.spawn(BrinkFlowRequest::<Batched>::builder().story(story0).build())
.id();
let e1 = app
.world_mut()
.spawn(BrinkFlowRequest::<Batched>::builder().story(story1).build())
.id();
app.update(); app.update(); app.update();
let log = app.world().resource::<PingLog>().0.clone();
assert_eq!(log.len(), 2, "both flows' commands fired; got {log:?}");
let expect_first = if e0 <= e1 { 10 } else { 20 };
let expect_second = if e0 <= e1 { 20 } else { 10 };
assert_eq!(
log,
vec![expect_first, expect_second],
"command triggers must flush in flow-id order (e0={e0:?}, e1={e1:?})"
);
}
#[test]
#[ignore = "scenario harness: wall-clock timing, run explicitly with --ignored --nocapture"]
fn batch_serial_scenario_numbers() {
use std::time::Instant;
let flow_counts = [1usize, 64, 512, 4096];
eprintln!("\n## batch-serial driver — one batch turn (provisional / in-wave)\n");
eprintln!("| flows | batch turn (ms) | writes applied | us/flow |");
eprintln!("|------:|----------------:|---------------:|--------:|");
for &n in &flow_counts {
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
app.add_systems(Update, advance_batch::<Batched>);
let (program, tables, ctx) =
compile_test_story("VAR g = 0\nA line of narration.\n~ g = g + 1\n-> END\n");
let story = add_story_assets(&mut app, program, tables, ctx);
for _ in 0..n {
app.world_mut().spawn(
BrinkFlowRequest::<Batched>::builder()
.story(story.clone())
.build(),
);
}
app.update();
let start = Instant::now();
app.update();
let elapsed = start.elapsed();
let report = app.world().resource::<BrinkBatchReport<Batched>>();
#[expect(
clippy::cast_precision_loss,
reason = "provisional scenario timing; bit-exactness not required"
)]
let us_per_flow = (elapsed.as_secs_f64() * 1_000_000.0) / (n.max(1) as f64);
eprintln!(
"| {n} | {:.3} | {} | {:.2} |",
elapsed.as_secs_f64() * 1000.0,
report.writes_applied,
us_per_flow,
);
assert_eq!(report.stepped, n, "all flows stepped for n={n}");
}
eprintln!();
}
#[test]
#[ignore = "scenario harness: wall-clock timing, run explicitly with --ignored --nocapture"]
fn batch_parallel_scenario_numbers() {
use std::time::Instant;
let flow_counts = [1usize, 64, 512, 4096];
eprintln!("\n## batch-PARALLEL driver — one batch turn (provisional / in-wave)\n");
eprintln!("| flows | batch turn (ms) | writes applied | us/flow |");
eprintln!("|------:|----------------:|---------------:|--------:|");
for &n in &flow_counts {
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
let (program, tables, ctx) =
compile_test_story("VAR g = 0\nA line of narration.\n~ g = g + 1\n-> END\n");
let story = add_story_assets(&mut app, program, tables, ctx);
for _ in 0..n {
app.world_mut().spawn(
BrinkFlowRequest::<Batched>::builder()
.story(story.clone())
.build(),
);
}
app.update();
let start = Instant::now();
advance_batch_parallel::<Batched>(app.world_mut());
let elapsed = start.elapsed();
let report = app.world().resource::<BrinkBatchReport<Batched>>();
#[expect(
clippy::cast_precision_loss,
reason = "provisional scenario timing; bit-exactness not required"
)]
let us_per_flow = (elapsed.as_secs_f64() * 1_000_000.0) / (n.max(1) as f64);
eprintln!(
"| {n} | {:.3} | {} | {:.2} |",
elapsed.as_secs_f64() * 1000.0,
report.writes_applied,
us_per_flow,
);
assert_eq!(report.stepped, n, "all flows stepped for n={n}");
}
eprintln!();
}
#[test]
fn parallel_equals_serial_over_randomized_workloads() {
let base_seed: u64 = std::env::var("BRINK_DETERMINISM_LAW_SEED")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(0);
let iterations: u64 = std::env::var("BRINK_DETERMINISM_LAW_ITERATIONS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(64);
for offset in 0u64..iterations {
let seed = base_seed.wrapping_add(offset);
let mut rng = TestRng::new(seed ^ 0x9E37_79B9_7F4A_7C15);
let flow_count = rng.pick(1, 10);
let assignments: Vec<usize> = (0..flow_count)
.map(|_| rng.pick(0, VARIANTS.len() - 1))
.collect();
let init_g = i32::try_from(rng.pick(0, 40)).unwrap_or(0);
let init_h = i32::try_from(rng.pick(0, 40)).unwrap_or(0);
let rng_seed = i32::try_from(rng.pick(1, 100_000)).unwrap_or(1);
let turns = rng.pick(1, 3);
let serial = run_workload(
&assignments,
init_g,
init_h,
rng_seed,
turns,
Driver::Serial,
);
let parallel = run_workload(
&assignments,
init_g,
init_h,
rng_seed,
turns,
Driver::Parallel,
);
assert_eq!(
serial, parallel,
"determinism law violated: parallel != serial at seed={seed} \
assignments={assignments:?} turns={turns} init=({init_g},{init_h}) rng={rng_seed}"
);
}
}
#[test]
fn advance_batch_parallel_system_drives_and_applies() {
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
app.add_systems(Update, advance_batch_parallel::<Batched>);
let (program, tables, ctx) = compile_test_story("VAR g = 0\nHi.\n~ g = 3\n-> END\n");
let story = add_story_assets(&mut app, program, tables, ctx);
for _ in 0..3 {
app.world_mut().spawn(
BrinkFlowRequest::<Batched>::builder()
.story(story.clone())
.build(),
);
}
app.update(); app.update();
let report = app.world().resource::<BrinkBatchReport<Batched>>();
assert_eq!(report.flows.len(), 3, "all three flows collected");
assert_eq!(report.skipped_local, 0);
assert_eq!(
report.stepped + report.awaiting + report.errored,
3,
"every collected flow is accounted for"
);
let globals = app.world().resource::<BrinkGlobals<Batched>>();
assert_eq!(globals.inner.global(0), &Value::Int(3));
}
#[test]
fn advance_batch_skips_local_policy_flow() {
let (mut app, story) = local_policy_app();
app.world_mut()
.spawn(BrinkFlowRequest::<Batched>::builder().story(story).build());
app.update(); app.world_mut()
.run_system_once(advance_batch::<Batched>)
.expect("serial batch runs");
assert_local_skipped(&app);
}
#[test]
fn advance_batch_parallel_skips_local_policy_flow() {
let (mut app, story) = local_policy_app();
app.world_mut()
.spawn(BrinkFlowRequest::<Batched>::builder().story(story).build());
app.update(); advance_batch_parallel::<Batched>(app.world_mut());
assert_local_skipped(&app);
}
#[derive(Clone, Copy)]
enum Driver {
Serial,
Parallel,
}
const VARIANTS: &[&str] = &[
"VAR g = 0\nVAR h = 0\nGee {g}.\n~ g = g + 3\n-> END\n",
"VAR g = 0\nVAR h = 0\n~ h = g * 2\nAich {h}.\n-> END\n",
"VAR g = 0\nVAR h = 0\n~ g = RANDOM(1, 6)\nRoll {g}.\n-> END\n",
"VAR g = 0\nVAR h = 0\nHello.\n+ [Wait]\n -> END\n",
];
struct TestRng(u64);
impl TestRng {
fn new(seed: u64) -> Self {
Self(seed)
}
fn next(&mut self) -> u64 {
self.0 = self
.0
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1_442_695_040_888_963_407);
self.0 >> 33
}
#[expect(
clippy::cast_possible_truncation,
reason = "ranges here are always small (flow counts, variant indices, seeds); truncation of the high LCG bits is intended and harmless"
)]
fn pick(&mut self, lo: usize, hi: usize) -> usize {
lo + (self.next() as usize) % (hi - lo + 1)
}
}
#[derive(PartialEq, Debug)]
struct RunSnapshot {
globals: Vec<Value>,
visit_counts: HashMap<DefinitionId, u32>,
turn_counts: HashMap<DefinitionId, u32>,
turn_index: u32,
rng_seed: i32,
previous_random: i32,
report: (usize, usize, usize, usize, usize, usize),
flow_flags: Vec<(bool, bool, bool)>,
events: Vec<(u8, String)>,
}
#[derive(Resource, Default)]
struct EventLog(Vec<(u8, String)>);
fn install_event_log(app: &mut App) {
app.init_resource::<EventLog>();
app.add_observer(
|on: On<BrinkLineDelivered<Batched>>, mut log: ResMut<EventLog>| {
log.0.push((0, on.event().text.clone()));
},
);
app.add_observer(
|on: On<BrinkChoicesPresented<Batched>>, mut log: ResMut<EventLog>| {
log.0.push((1, on.event().text.clone()));
},
);
app.add_observer(
|on: On<BrinkTurnDone<Batched>>, mut log: ResMut<EventLog>| {
log.0.push((2, on.event().text.clone()));
},
);
app.add_observer(
|on: On<BrinkStoryEnded<Batched>>, mut log: ResMut<EventLog>| {
log.0.push((3, on.event().text.clone()));
},
);
}
fn run_workload(
assignments: &[usize],
init_g: i32,
init_h: i32,
rng_seed: i32,
turns: usize,
driver: Driver,
) -> RunSnapshot {
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default());
install_event_log(&mut app);
let story_handles: Vec<_> = VARIANTS
.iter()
.map(|src| {
let (p, t, c) = compile_test_story(src);
add_story_assets(&mut app, p, t, c)
})
.collect();
for &v in assignments {
app.world_mut().spawn(
BrinkFlowRequest::<Batched>::builder()
.story(story_handles[v].clone())
.build(),
);
}
app.update();
app.update();
{
let mut globals = app.world_mut().resource_mut::<BrinkGlobals<Batched>>();
globals.inner.set_global(0, Value::Int(init_g));
globals.inner.set_global(1, Value::Int(init_h));
globals.inner.set_rng_seed(rng_seed);
}
for _ in 0..turns {
match driver {
Driver::Serial => {
app.world_mut()
.run_system_once(advance_batch::<Batched>)
.expect("serial batch runs");
}
Driver::Parallel => advance_batch_parallel::<Batched>(app.world_mut()),
}
}
let world = &app.world().resource::<BrinkGlobals<Batched>>().inner;
let report = app.world().resource::<BrinkBatchReport<Batched>>();
RunSnapshot {
globals: world.globals.clone(),
visit_counts: world.visit_counts.clone(),
turn_counts: world.turn_counts.clone(),
turn_index: world.turn_index,
rng_seed: world.rng_seed,
previous_random: world.previous_random,
report: (
report.stepped,
report.awaiting,
report.errored,
report.skipped_local,
report.writes_applied,
report.commands_applied,
),
flow_flags: report
.flows
.iter()
.map(|f| (f.awaiting, f.errored, f.skipped_local))
.collect(),
events: app.world().resource::<EventLog>().0.clone(),
}
}
fn local_policy_app() -> (App, bevy_asset::Handle<crate::asset::BrinkStoryAsset>) {
let policy = WorldPolicy {
overrides: std::iter::once(("g".to_string(), Scope::Local)).collect(),
..Default::default()
};
let mut app = App::new();
app.add_plugins(AssetPlugin::default());
app.add_plugins(crate::BrinkPlugin::<Batched>::default().with_policy(policy));
let (program, tables, ctx) = compile_test_story("VAR g = 0\nHi.\n~ g = 5\n-> END\n");
let story = add_story_assets(&mut app, program, tables, ctx);
(app, story)
}
fn assert_local_skipped(app: &App) {
let report = app.world().resource::<BrinkBatchReport<Batched>>();
assert_eq!(report.skipped_local, 1, "the Local-policy flow was skipped");
assert_eq!(report.stepped, 0, "a skipped flow is not stepped");
assert_eq!(report.awaiting, 0);
assert_eq!(report.errored, 0);
assert_eq!(report.writes_applied, 0, "a skipped flow applies no writes");
assert!(
report.flows.iter().all(|f| f.skipped_local),
"every flow record flags skipped_local"
);
let globals = app.world().resource::<BrinkGlobals<Batched>>();
assert_eq!(globals.inner.global(0), &Value::Int(0));
}