mod apply;
mod eval;
mod instance;
mod resolve;
mod state;
mod trace;
#[cfg(test)]
mod bench;
#[cfg(test)]
mod test_world;
#[cfg(test)]
mod tests;
use alloc::boxed::Box;
use alloc::vec::Vec;
use eval::{EvalCtx, PARALLEL_EVAL_MIN_JOBS, Snapshot, eval_one};
use instance::Instance;
use resolve::{Resolved, SourceTicks};
pub use eval::{EvalBucket, EvalScheduler};
pub use state::{BehaviorState, BehaviorStore, def_hash};
use crate::behavior::{Effect, Program, Val, VarTable};
use crate::components::{Behavior, BehaviorSource, InteractEvent, Variables, VolumeEvent};
use crate::ecs::{
Entity, EntityByName, EventCursor, FrameContext, MenuActive, PipelineContext, ScheduleMode,
SimTiming, StepResult, System, TraceRequest, TransientSaves,
};
#[derive(Debug, Default)]
pub struct BehaviorSystem {
programs: Vec<Program>,
instances: Vec<Vec<Instance>>,
vars: Vec<Val>,
var_table: VarTable,
pending: Vec<(usize, Option<Entity>, f32)>,
crossing_cursor: EventCursor,
press_cursor: EventCursor,
crossings: Vec<VolumeEvent>,
presses: Vec<InteractEvent>,
store: Option<Box<dyn BehaviorStore>>,
scheduler: Option<Box<dyn EvalScheduler>>,
transient_saves: bool,
sources: SourceTicks,
trace_frame: u64,
trace_paths_published: bool,
sim_ticks: u64,
populated: bool,
eval_buckets: Vec<EvalBucket>,
jobs: Vec<(usize, Option<Entity>)>,
bindings: Vec<Option<Val>>,
serial_effects: Vec<Effect>,
serial_produced: Vec<(usize, Option<Entity>, usize)>,
snapshot: Snapshot,
tag_scratch: Vec<Entity>,
}
impl BehaviorSystem {
pub fn new() -> Self {
Self::default()
}
pub fn with_store(mut self, store: Box<dyn BehaviorStore>) -> Self {
self.store = Some(store);
self
}
pub fn with_scheduler(mut self, scheduler: Box<dyn EvalScheduler>) -> Self {
self.scheduler = Some(scheduler);
self
}
}
impl System for BehaviorSystem {
fn init(&mut self, ctx: &mut PipelineContext) {
self.reseed(ctx);
self.transient_saves = ctx.resource::<TransientSaves>().is_some_and(|t| t.0);
self.trace_frame = 0;
self.sim_ticks = 0;
let restored = self.restore_state();
tracing::info!(
"BehaviorSystem: {} behavior(s), {} variable(s), restored {}",
self.programs.len(),
self.vars.len(),
restored,
);
}
fn step(&mut self, ctx: &mut PipelineContext) -> StepResult {
self.reseed_if_edited(ctx);
if self.programs.is_empty() {
return StepResult::Continue;
}
if let Some(events) = ctx.events::<VolumeEvent>() {
self.crossings
.extend(events.read(&mut self.crossing_cursor).copied());
}
if let Some(events) = ctx.events::<InteractEvent>() {
self.presses
.extend(events.read(&mut self.press_cursor).copied());
}
let menu_active = ctx.resource::<MenuActive>().map(|m| m.0).unwrap_or(false);
if menu_active {
return StepResult::Continue;
}
let timing = ctx.resource::<SimTiming>().copied().unwrap_or_default();
for _ in 0..timing.ticks {
self.sim_ticks += 1;
let elapsed = (self.sim_ticks as f64 * timing.tick_dt as f64) as f32;
self.tick(ctx, timing.tick_dt, elapsed);
}
StepResult::Continue
}
}
impl BehaviorSystem {
fn reseed(&mut self, ctx: &PipelineContext) {
let resolved = resolve::resolve(
ctx.query::<Variables>().as_slice(),
ctx.query::<Behavior>().as_slice(),
);
self.sources = SourceTicks::of(ctx);
self.adopt(resolved);
}
fn reseed_if_edited(&mut self, ctx: &PipelineContext) {
if SourceTicks::of(ctx) != self.sources {
self.reseed(ctx);
}
}
fn adopt(&mut self, resolved: Resolved) {
let programs = core::mem::take(&mut self.programs);
let instances = core::mem::take(&mut self.instances);
let carried = resolve::carry_instances(&programs, instances, &resolved.programs);
let moved = carried.moved;
self.pending.retain_mut(
|(program, _, _)| match moved.get(*program).copied().flatten() {
Some(next) => {
*program = next;
true
}
None => false,
},
);
self.vars = resolve::carry_vars(&self.var_table, &self.vars, &resolved.var_table);
self.instances = carried.instances;
self.programs = resolved.programs;
self.var_table = resolved.var_table;
self.trace_paths_published = false;
}
fn restore_state(&mut self) -> usize {
if self.transient_saves || !self.programs.iter().any(|p| p.def.saves_state()) {
return 0;
}
let Some(state) = self.store.as_ref().and_then(|store| store.read()) else {
return 0;
};
let mut restored = 0usize;
for (name, value) in &state.vars {
let Some(slot) = self.var_table.slot_of(name) else {
continue;
};
let value = Val::from_literal(value);
if self.vars[slot as usize].same_type(value) {
self.vars[slot as usize] = value;
restored += 1;
}
}
for (id, hash) in state.fired {
if let Some(i) = self
.programs
.iter()
.position(|p| p.def.asset_id.0 == id && def_hash(&p.def) == hash)
{
if !self.programs[i].is_scoped() {
let mut instance = Instance::new(None, Vec::new(), false);
instance.fired_once = true;
self.instances[i].clear();
self.instances[i].push(instance);
}
}
}
restored
}
fn resync_instances(&mut self, snapshot: &Snapshot, frame: FrameContext) {
let baselines = frame.collect(self.programs.iter().map(|p| {
match &p.def.on {
BehaviorSource::Variable(name) => self
.var_table
.slot_of(name)
.and_then(|s| self.vars.get(s as usize))
.copied()
.unwrap_or(Val::Int(0)),
_ => Val::Int(0),
}
}));
for (i, program) in self.programs.iter().enumerate() {
if !program.is_scoped() {
if self.instances[i].is_empty() {
let mut instance = Instance::new(None, Vec::new(), false);
instance.last_value = baselines[i];
self.instances[i].push(instance);
}
continue;
}
let matched = &snapshot.scoped[i];
self.instances[i].retain(|inst| {
inst.entity.is_some_and(|e| {
matched
.binary_search_by_key(&e.to_bits(), |o| o.to_bits())
.is_ok()
})
});
let before = self.instances[i].len();
let mut j = 0;
for entity in matched {
if j < before && self.instances[i][j].entity == Some(*entity) {
j += 1;
continue;
}
let mut instance =
Instance::new(Some(*entity), program.local_inits.clone(), self.populated);
instance.last_value = baselines[i];
self.instances[i].push(instance);
}
if self.instances[i].len() != before {
self.instances[i].sort_by_key(|inst| inst.entity.map(|e| e.to_bits()));
}
}
}
fn tick(&mut self, ctx: &mut PipelineContext, dt: f32, elapsed: f32) {
let request = ctx.resource::<TraceRequest>().cloned();
let tracing = request.is_some();
let mut fired: Vec<(usize, Vec<u32>)> = Vec::new();
let frame = ctx.frame;
let mut snapshot = core::mem::take(&mut self.snapshot);
self.gather(ctx, &mut snapshot);
self.resync_instances(&snapshot, frame);
self.populated = true;
let bound: usize = self.instances.iter().map(Vec::len).sum();
let mut runs = frame.vec::<(usize, Option<Entity>)>(bound);
for i in 0..self.programs.len() {
let var_slot = match &self.programs[i].def.on {
BehaviorSource::Variable(name) => self.var_table.slot_of(name),
_ => None,
};
let def = &self.programs[i].def;
for instance in &mut self.instances[i] {
if instance.due(
def,
&self.vars,
var_slot,
dt,
&self.crossings,
&self.presses,
) {
runs.push((i, instance.entity));
}
}
}
self.crossings.clear();
self.presses.clear();
let mut jobs = core::mem::take(&mut self.jobs);
jobs.clear();
let mut idx = 0;
while idx < self.pending.len() {
self.pending[idx].2 -= dt;
if self.pending[idx].2 <= 0.0 {
let (i, entity, _) = self.pending.swap_remove(idx);
jobs.push((i, entity));
} else {
idx += 1;
}
}
for &(i, entity) in runs.iter() {
let delay = self.programs[i].def.delay;
if delay > 0.0 {
self.pending.push((i, entity, delay));
} else {
jobs.push((i, entity));
}
}
let parallel = self.scheduler.is_some()
&& jobs.len() >= PARALLEL_EVAL_MIN_JOBS
&& ScheduleMode::current(ctx.resources) == ScheduleMode::Parallel;
let scheduler = self.scheduler.as_deref().filter(|_| parallel);
let mut effects = core::mem::take(&mut self.serial_effects);
let mut produced = core::mem::take(&mut self.serial_produced);
effects.clear();
produced.clear();
let mut buckets = core::mem::take(&mut self.eval_buckets);
let mut serial_bindings = core::mem::take(&mut self.bindings);
{
let ec = EvalCtx {
components: ctx.components,
names: ctx.resource::<EntityByName>(),
snapshot: &snapshot,
programs: &self.programs,
instances: &self.instances,
vars: &self.vars,
dt,
elapsed,
tracing,
};
if let Some(scheduler) = scheduler {
let workers = scheduler.workers().max(1);
while buckets.len() < workers {
buckets.push(EvalBucket::default());
}
let chunk = jobs.len().div_ceil(buckets.len()).max(1);
for (b, bucket) in buckets.iter_mut().enumerate() {
bucket.jobs = (b * chunk).min(jobs.len())..((b + 1) * chunk).min(jobs.len());
}
let jobs = &jobs;
let ec = &ec;
scheduler.run(&mut buckets, &|bucket| {
bucket.effects.clear();
bucket.produced.clear();
bucket.fired.clear();
for &(i, entity) in &jobs[bucket.jobs.clone()] {
if let Some((count, nodes)) =
eval_one(ec, &mut bucket.bindings, i, entity, &mut bucket.effects)
{
bucket.produced.push((i, entity, count));
if ec.tracing {
bucket.fired.push((i, nodes));
}
}
}
});
} else {
for &(i, entity) in &jobs {
if let Some((count, nodes)) =
eval_one(&ec, &mut serial_bindings, i, entity, &mut effects)
{
produced.push((i, entity, count));
if tracing {
fired.push((i, nodes));
}
}
}
}
}
let mut save_requested = false;
if parallel {
for bucket in &mut buckets {
let mut recorded = bucket.effects.drain(..);
for k in 0..bucket.produced.len() {
let (i, entity, count) = bucket.produced[k];
save_requested |= self.apply(ctx, i, entity, recorded.by_ref().take(count));
}
if tracing {
fired.append(&mut bucket.fired);
}
}
} else {
let mut recorded = effects.drain(..);
for &(i, entity, count) in &produced {
save_requested |= self.apply(ctx, i, entity, recorded.by_ref().take(count));
}
}
self.eval_buckets = buckets;
self.jobs = jobs;
self.bindings = serial_bindings;
self.serial_effects = effects;
self.serial_produced = produced;
self.snapshot = snapshot;
if save_requested {
self.write_state();
}
if let Some(request) = request {
self.publish_trace(ctx, &request, &fired);
}
}
fn write_state(&self) {
if self.transient_saves {
return;
}
let Some(store) = self.store.as_ref() else {
return;
};
store.write(&BehaviorState {
vars: self
.var_table
.names()
.iter()
.zip(&self.vars)
.map(|(name, value)| (name.clone(), value.to_literal()))
.collect(),
fired: self
.programs
.iter()
.enumerate()
.filter(|(i, p)| {
p.def.once
&& !p.is_scoped()
&& self.instances[*i].iter().any(|inst| inst.fired_once)
})
.map(|(_, p)| (p.def.asset_id.0, def_hash(&p.def)))
.collect(),
});
}
}