use std::marker::PhantomData;
use bevy_asset::{AssetId, Assets};
use bevy_ecs::change_detection::DetectChanges;
use bevy_ecs::change_detection::Tick;
use bevy_ecs::entity::Entity;
use bevy_ecs::query::Access;
use bevy_ecs::resource::Resource;
use bevy_ecs::system::{Commands, Query, Res, ResMut};
use bevy_log::warn;
use brink_format::{DefinitionId, LineEntry, Value};
use brink_runtime::{
ContextAccess, DriveOutcome, ExternalFnHandler, FallbackHandler, FastRng, FlowInstance,
FrameStartView, Program, RuntimeError, Scope, Step, World, WorldPolicy, WriteObserver,
};
use crate::asset::{BrinkProgram, LineTablesAsset, ProgramAsset};
use crate::bindings::{BrinkBindings, TriggerFn};
use crate::capability::{CapabilityTable, ContainerAccessTable};
use crate::flow::{BrinkFlow, emit_event};
use crate::globals::{BrinkGlobals, BrinkWorldPolicy};
use crate::line_tables::BrinkLocale;
use crate::sleep::FlowSleep;
use crate::wake_delta::{BrinkWorldDelta, WorldDelta};
pub mod parallel;
pub(crate) fn homes_any_local(policy: &WorldPolicy) -> bool {
policy.default == Scope::Local
|| policy.turn_index == Scope::Local
|| policy.rng == Scope::Local
|| policy.overrides.values().any(|s| *s == Scope::Local)
}
#[derive(Debug, Clone)]
pub(crate) enum WorldWrite {
Global(u32, Value),
VisitCount(DefinitionId, u32),
TurnCount(DefinitionId, u32),
TurnIndex(u32),
RngSeed(i32),
PreviousRandom(i32),
}
impl WorldWrite {
fn apply(&self, target: &mut World) {
match *self {
WorldWrite::Global(idx, ref value) => target.set_global(idx, value.clone()),
WorldWrite::VisitCount(id, count) => target.set_visit_count(id, count),
WorldWrite::TurnCount(id, turn) => target.set_turn_count(id, turn),
WorldWrite::TurnIndex(index) => target.set_turn_index(index),
WorldWrite::RngSeed(seed) => target.set_rng_seed(seed),
WorldWrite::PreviousRandom(val) => target.set_previous_random(val),
}
}
}
#[derive(Default)]
pub(crate) struct WriteBuffer {
writes: Vec<WorldWrite>,
}
impl WriteBuffer {
fn apply_to(&self, target: &mut World) {
for w in &self.writes {
w.apply(target);
}
}
fn record_into(&self, delta: &mut WorldDelta) {
for w in &self.writes {
match *w {
WorldWrite::Global(idx, _) => delta.note_global(idx),
WorldWrite::VisitCount(..)
| WorldWrite::TurnCount(..)
| WorldWrite::TurnIndex(_)
| WorldWrite::RngSeed(_)
| WorldWrite::PreviousRandom(_) => delta.note_bookkeeping(),
}
}
}
#[cfg(test)]
pub(crate) fn len(&self) -> usize {
self.writes.len()
}
}
impl WriteObserver for WriteBuffer {
fn on_set_global(&mut self, idx: u32, value: &Value) {
self.writes.push(WorldWrite::Global(idx, value.clone()));
}
fn on_increment_visit(&mut self, id: DefinitionId, new_count: u32) {
self.writes.push(WorldWrite::VisitCount(id, new_count));
}
fn on_set_visit_count(&mut self, id: DefinitionId, count: u32) {
self.writes.push(WorldWrite::VisitCount(id, count));
}
fn on_set_turn_count(&mut self, id: DefinitionId, turn: u32) {
self.writes.push(WorldWrite::TurnCount(id, turn));
}
fn on_increment_turn_index(&mut self, new_value: u32) {
self.writes.push(WorldWrite::TurnIndex(new_value));
}
fn on_set_turn_index(&mut self, index: u32) {
self.writes.push(WorldWrite::TurnIndex(index));
}
fn on_set_rng_seed(&mut self, new_seed: i32) {
self.writes.push(WorldWrite::RngSeed(new_seed));
}
fn on_set_previous_random(&mut self, new_val: i32) {
self.writes.push(WorldWrite::PreviousRandom(new_val));
}
}
pub(crate) struct FlowBatchOutcome {
entity: Entity,
story: AssetId<ProgramAsset>,
writes: WriteBuffer,
triggers: Vec<TriggerFn>,
lines: Vec<Step>,
awaiting: bool,
errored: bool,
skipped_local: bool,
access: Option<Access>,
}
impl FlowBatchOutcome {
pub(crate) fn skipped_local(
entity: Entity,
story: AssetId<ProgramAsset>,
access: Option<Access>,
) -> Self {
Self {
entity,
story,
writes: WriteBuffer::default(),
triggers: Vec::new(),
lines: Vec::new(),
awaiting: false,
errored: false,
skipped_local: true,
access,
}
}
pub(crate) fn entity(&self) -> Entity {
self.entity
}
}
#[expect(
clippy::too_many_arguments,
reason = "each argument is a distinct, already-resolved Step input (frame-start, flow, program, tables, bindings, story, entity, access); bundling them into a struct would just relocate the same fields with no clarity gain and force both call sites to build it"
)]
pub(crate) fn step_one<M: Send + Sync + 'static>(
frame_start: &World,
flow_inner: &mut FlowInstance,
program: &Program,
tables: &[Vec<LineEntry>],
bindings: Option<&BrinkBindings<M>>,
story: AssetId<ProgramAsset>,
entity: Entity,
access: Option<Access>,
) -> FlowBatchOutcome {
let mut buf = WriteBuffer::default();
let handler = bindings.map(BrinkBindings::handler);
let handler_ref: &dyn ExternalFnHandler = match &handler {
Some(h) => h,
None => &FallbackHandler,
};
let (lines, awaiting, error) = step_flow(
frame_start,
flow_inner,
program,
tables,
handler_ref,
&mut buf,
);
let errored = if let Some(err) = &error {
warn!("batch step faulted for flow {entity:?} (story {story:?}): {err}");
true
} else {
false
};
let triggers = handler.map(|h| h.take_queued()).unwrap_or_default();
FlowBatchOutcome {
entity,
story,
writes: buf,
triggers,
lines,
awaiting,
errored,
skipped_local: false,
access,
}
}
pub(crate) fn aggregate_access(table: &ContainerAccessTable) -> Access {
let mut acc = Access::default();
for container in table.values() {
acc.extend(&container.access);
}
acc
}
fn step_flow(
frame_start: &World,
flow: &mut FlowInstance,
program: &Program,
line_tables: &[Vec<LineEntry>],
handler: &dyn ExternalFnHandler,
buf: &mut WriteBuffer,
) -> (Vec<Step>, bool, Option<RuntimeError>) {
let mut scratch = FrameStartView::new(frame_start);
let mut observed = brink_runtime::ObservedContext::new(&mut scratch, buf);
let mut budget = FlowInstance::LINE_LIMIT;
match flow.drive::<FastRng>(
program,
line_tables,
&mut observed,
handler,
None,
&mut budget,
) {
Ok(DriveOutcome::Terminal(lines)) => (lines, false, None),
Ok(DriveOutcome::AwaitingExternal(lines)) => (lines, true, None),
Err(err) => (Vec::new(), false, Some(err)),
}
}
#[derive(Debug, Clone)]
pub struct FlowAccessRecord {
pub entity: Entity,
pub story: AssetId<ProgramAsset>,
pub access: Option<Access>,
pub awaiting: bool,
pub errored: bool,
pub skipped_local: bool,
}
#[derive(Resource)]
pub struct BrinkBatchReport<M: Send + Sync + 'static = ()> {
pub flows: Vec<FlowAccessRecord>,
pub stepped: usize,
pub awaiting: usize,
pub errored: usize,
pub skipped_local: usize,
pub writes_applied: usize,
pub commands_applied: usize,
_marker: PhantomData<fn() -> M>,
}
impl<M: Send + Sync + 'static> Default for BrinkBatchReport<M> {
fn default() -> Self {
Self {
flows: Vec::new(),
stepped: 0,
awaiting: 0,
errored: 0,
skipped_local: 0,
writes_applied: 0,
commands_applied: 0,
_marker: PhantomData,
}
}
}
impl<M: Send + Sync + 'static> BrinkBatchReport<M> {
fn record(&mut self, result: BatchApplyResult) {
self.flows = result.flows;
self.stepped = result.stepped;
self.awaiting = result.awaiting;
self.errored = result.errored;
self.skipped_local = result.skipped_local;
self.writes_applied = result.writes_applied;
self.commands_applied = result.commands_applied;
}
}
#[derive(Default)]
pub(crate) struct BatchApplyResult {
pub flows: Vec<FlowAccessRecord>,
pub stepped: usize,
pub awaiting: usize,
pub errored: usize,
pub skipped_local: usize,
pub writes_applied: usize,
pub commands_applied: usize,
pub changed: WorldDelta,
}
pub(crate) struct DeferredFlush {
entity: Entity,
triggers: Vec<TriggerFn>,
lines: Vec<Step>,
}
pub(crate) fn apply_batch_writes(
outcomes: Vec<FlowBatchOutcome>,
world: &mut World,
) -> (BatchApplyResult, Vec<DeferredFlush>) {
let mut flows = Vec::with_capacity(outcomes.len());
let mut deferred = Vec::with_capacity(outcomes.len());
let mut stepped = 0usize;
let mut awaiting = 0usize;
let mut errored = 0usize;
let mut skipped_local = 0usize;
let mut writes_applied = 0usize;
let mut commands_applied = 0usize;
let mut changed = WorldDelta::default();
for outcome in outcomes {
outcome.writes.apply_to(world);
outcome.writes.record_into(&mut changed);
writes_applied += outcome.writes.writes.len();
commands_applied += outcome.triggers.len();
if outcome.skipped_local {
skipped_local += 1;
} else if outcome.errored {
errored += 1;
} else if outcome.awaiting {
awaiting += 1;
} else {
stepped += 1;
}
flows.push(FlowAccessRecord {
entity: outcome.entity,
story: outcome.story,
access: outcome.access,
awaiting: outcome.awaiting,
errored: outcome.errored,
skipped_local: outcome.skipped_local,
});
deferred.push(DeferredFlush {
entity: outcome.entity,
triggers: outcome.triggers,
lines: outcome.lines,
});
}
(
BatchApplyResult {
flows,
stepped,
awaiting,
errored,
skipped_local,
writes_applied,
commands_applied,
changed,
},
deferred,
)
}
pub(crate) fn record_wake_delta<M: Send + Sync + 'static>(
ledger: &mut BrinkWorldDelta<M>,
result: &BatchApplyResult,
globals_changed_on_entry: bool,
globals_tick: Option<Tick>,
) {
if globals_changed_on_entry {
ledger.note_foreign();
}
ledger.record(&result.changed, globals_tick);
}
pub(crate) fn flush_deferred<M: Send + Sync + 'static>(
deferred: Vec<DeferredFlush>,
commands: &mut Commands,
) {
for flush in deferred {
for trigger in flush.triggers {
commands.queue(trigger);
}
for line in &flush.lines {
emit_event::<M>(line, flush.entity, commands);
}
}
}
#[expect(
clippy::needless_pass_by_value,
clippy::type_complexity,
clippy::too_many_arguments,
reason = "bevy systems take Res/Query by value; the flow query tuple is inherently wide, and the phase inputs (globals, policy, assets, bindings, capability table, report) are each a distinct system param"
)]
pub fn advance_batch<M: Send + Sync + 'static>(
mut flows: Query<(
Entity,
&mut BrinkFlow<M>,
&BrinkProgram<M>,
&BrinkLocale<M>,
Option<&FlowSleep<M>>,
)>,
globals: Option<ResMut<BrinkGlobals<M>>>,
policy: Option<Res<BrinkWorldPolicy<M>>>,
programs: Res<Assets<ProgramAsset>>,
line_tables_assets: Res<Assets<LineTablesAsset>>,
bindings: Option<Res<BrinkBindings<M>>>,
cap_table: Res<CapabilityTable<M>>,
report: Option<ResMut<BrinkBatchReport<M>>>,
wake_delta: Option<ResMut<BrinkWorldDelta<M>>>,
mut commands: Commands,
) {
let Some(mut globals) = globals else {
return;
};
let globals_changed_on_entry = globals.is_changed();
let policy_local = policy.is_some_and(|p| homes_any_local(&p.policy));
let mut collected: Vec<Entity> = flows
.iter()
.filter(|(_, flow, _, _, sleep)| {
!flow.inner.has_pending_external() && sleep.is_none_or(FlowSleep::wants_collect)
})
.map(|(e, _, _, _, _)| e)
.collect();
collected.sort_unstable();
let frame_start = globals.inner.clone();
let mut outcomes: Vec<FlowBatchOutcome> = Vec::with_capacity(collected.len());
for &entity in &collected {
let Ok((_, mut flow, program_ref, locale, _)) = flows.get_mut(entity) else {
continue;
};
let Some(program_asset) = programs.get(&program_ref.handle) else {
continue;
};
let Some(lt_asset) = line_tables_assets.get(&locale.handle) else {
continue;
};
let story = program_ref.handle.id();
let access = cap_table.access_for(story).map(aggregate_access);
if policy_local || program_asset.program.has_local_defaults() {
warn!(
"batch skipping Local-policy flow {entity:?} (story {story:?}): \
batch mode routes only shared World state — keep it on the serial API"
);
outcomes.push(FlowBatchOutcome::skipped_local(entity, story, access));
continue;
}
outcomes.push(step_one::<M>(
&frame_start,
&mut flow.inner,
&program_asset.program,
<_asset.tables,
bindings.as_deref(),
story,
entity,
access,
));
}
let (result, deferred) = if outcomes.is_empty() {
(BatchApplyResult::default(), Vec::new())
} else {
apply_batch_writes(outcomes, &mut globals.inner)
};
flush_deferred::<M>(deferred, &mut commands);
if let Some(mut ledger) = wake_delta {
record_wake_delta(
&mut ledger,
&result,
globals_changed_on_entry,
Some(globals.last_changed()),
);
}
if let Some(mut report) = report {
report.record(result);
}
}
#[cfg(test)]
#[expect(clippy::panic, reason = "tests assert via panic on the error arm")]
mod tests;