#![expect(
unsafe_code,
reason = "BH-3 sanctioned-unsafe module (decision-log 2026-07-16): the parallel Step phase drives access-disjoint flows on ComputeTaskPool through an UnsafeWorldCell, bevy's own multi-threaded-executor primitive. The workspace-wide unsafe_code deny stands everywhere else; every unsafe block carries a SAFETY argument grounded in per-flow entity/Access-set disjointness."
)]
use bevy_asset::{AssetId, Assets};
use bevy_ecs::change_detection::DetectChanges;
use bevy_ecs::entity::Entity;
use bevy_ecs::query::Access;
use bevy_ecs::system::Commands;
use bevy_ecs::world::{CommandQueue, World};
use bevy_log::warn;
use bevy_tasks::{ComputeTaskPool, TaskPool};
use brink_format::LineEntry;
use brink_runtime::{Program, World as BrinkWorld};
use super::{
BatchApplyResult, BrinkBatchReport, FlowBatchOutcome, aggregate_access, apply_batch_writes,
flush_deferred, homes_any_local, record_wake_delta, step_one,
};
use crate::asset::{BrinkProgram, LineTablesAsset, ProgramAsset};
use crate::bindings::BrinkBindings;
use crate::capability::CapabilityTable;
use crate::flow::BrinkFlow;
use crate::globals::{BrinkGlobals, BrinkWorldPolicy};
use crate::line_tables::BrinkLocale;
use crate::sleep::FlowSleep;
use crate::wake_delta::BrinkWorldDelta;
struct Prep {
entity: Entity,
story: AssetId<ProgramAsset>,
line_tables: AssetId<LineTablesAsset>,
is_local: bool,
access: Option<Access>,
}
struct Job<'w> {
entity: Entity,
story: AssetId<ProgramAsset>,
program: &'w Program,
tables: &'w [Vec<LineEntry>],
access: Option<Access>,
}
pub fn advance_batch_parallel<M: Send + Sync + 'static>(world: &mut World) {
let globals_changed_on_entry = {
let Some(globals) = world.get_resource_ref::<BrinkGlobals<M>>() else {
return;
};
globals.is_changed()
};
let mut query = world.query::<(
Entity,
&BrinkFlow<M>,
&BrinkProgram<M>,
&BrinkLocale<M>,
Option<&FlowSleep<M>>,
)>();
let policy_local = world
.get_resource::<BrinkWorldPolicy<M>>()
.is_some_and(|p| homes_any_local(&p.policy));
let programs = world.resource::<Assets<ProgramAsset>>();
let line_tables_assets = world.resource::<Assets<LineTablesAsset>>();
let cap_table = world.resource::<CapabilityTable<M>>();
let mut preps: Vec<Prep> = query
.iter(world)
.filter(|(_, flow, _, _, sleep)| {
!flow.inner.has_pending_external() && sleep.is_none_or(FlowSleep::wants_collect)
})
.filter_map(|(entity, _, program_ref, locale, _)| {
let story = program_ref.handle.id();
let program_asset = programs.get(story)?;
let line_tables = locale.handle.id();
line_tables_assets.get(line_tables)?;
let access = cap_table.access_for(story).map(aggregate_access);
let is_local = policy_local || program_asset.program.has_local_defaults();
Some(Prep {
entity,
story,
line_tables,
is_local,
access,
})
})
.collect();
preps.sort_unstable_by_key(|p| p.entity);
let frame_start: BrinkWorld = world.resource::<BrinkGlobals<M>>().inner.clone();
let outcomes = parallel_step::<M>(world, &preps, &frame_start);
let (result, deferred) = if outcomes.is_empty() {
(BatchApplyResult::default(), Vec::new())
} else {
let mut globals = world.resource_mut::<BrinkGlobals<M>>();
apply_batch_writes(outcomes, &mut globals.inner)
};
let globals_tick = world
.get_resource_ref::<BrinkGlobals<M>>()
.map(|globals| globals.last_changed());
world.increment_change_tick();
if let Some(mut ledger) = world.get_resource_mut::<BrinkWorldDelta<M>>() {
record_wake_delta(&mut ledger, &result, globals_changed_on_entry, globals_tick);
}
let mut queue = CommandQueue::default();
{
let mut commands = Commands::new(&mut queue, world);
flush_deferred::<M>(deferred, &mut commands);
}
queue.apply(world);
if let Some(mut report) = world.get_resource_mut::<BrinkBatchReport<M>>() {
report.record(result);
}
}
fn parallel_step<M: Send + Sync + 'static>(
world: &mut World,
preps: &[Prep],
frame_start: &BrinkWorld,
) -> Vec<FlowBatchOutcome> {
let mut outcomes: Vec<FlowBatchOutcome> = Vec::with_capacity(preps.len());
for prep in preps.iter().filter(|p| p.is_local) {
warn!(
"parallel batch skipping Local-policy flow {:?} (story {:?}): \
batch mode routes only shared World state — keep it on the serial API",
prep.entity, prep.story
);
outcomes.push(FlowBatchOutcome::skipped_local(
prep.entity,
prep.story,
prep.access.clone(),
));
}
let cell = world.as_unsafe_world_cell();
let programs = unsafe { cell.get_resource::<Assets<ProgramAsset>>() };
let line_tables_assets = unsafe { cell.get_resource::<Assets<LineTablesAsset>>() };
let bindings = unsafe { cell.get_resource::<BrinkBindings<M>>() };
let (Some(programs), Some(line_tables_assets)) = (programs, line_tables_assets) else {
return finish(outcomes);
};
let jobs: Vec<Job> = preps
.iter()
.filter(|p| !p.is_local)
.filter_map(|prep| {
let program = &programs.get(prep.story)?.program;
let tables = &line_tables_assets.get(prep.line_tables)?.tables;
Some(Job {
entity: prep.entity,
story: prep.story,
program,
tables,
access: prep.access.clone(),
})
})
.collect();
let pool = ComputeTaskPool::get_or_init(TaskPool::default);
let stepped: Vec<FlowBatchOutcome> = pool
.scope(|scope| {
for job in &jobs {
scope.spawn(async move {
let entity_cell = cell.get_entity(job.entity).ok()?;
let mut flow = unsafe { entity_cell.get_mut::<BrinkFlow<M>>() }?;
Some(step_one::<M>(
frame_start,
&mut flow.inner,
job.program,
job.tables,
bindings,
job.story,
job.entity,
job.access.clone(),
))
});
}
})
.into_iter()
.flatten()
.collect();
outcomes.extend(stepped);
finish(outcomes)
}
fn finish(mut outcomes: Vec<FlowBatchOutcome>) -> Vec<FlowBatchOutcome> {
outcomes.sort_unstable_by_key(FlowBatchOutcome::entity);
outcomes
}