use std::future::Future;
use std::pin::Pin;
use bevy_asset::Assets;
use bevy_ecs::entity::Entity;
use bevy_ecs::system::{Query, Res, ResMut, SystemState};
use bevy_ecs::world::World;
use bevy_log::warn;
use bevy_tasks::{AsyncComputeTaskPool, TaskPool};
use brink_format::Value;
use brink_runtime::{FastRng, FlowInstance, Program, RuntimeError, Step, StepOutcome};
#[cfg(feature = "dev")]
use brink_runtime::{RecordingHandler, ReplayRecorder};
use thiserror::Error;
use crate::asset::{BrinkProgram, LineTablesAsset, ProgramAsset};
use crate::async_bind::{BrinkAwaiting, BrinkExternalAwaited, BrinkPendingTask};
use crate::flow::BrinkFlow;
use crate::globals::BrinkContext;
use crate::line_tables::BrinkLocale;
use super::registration::{AsyncKind, BrinkBindings, BrinkHandler, QuerySystemId, TriggerFn};
#[derive(Debug, Error)]
pub enum BrinkCallError {
#[error("entity is not a fulfilled brink flow")]
NotAFlow,
#[error("program asset not loaded")]
ProgramNotLoaded,
#[error("line tables asset not loaded")]
LineTablesNotLoaded,
#[error("function '{0}' not found")]
FunctionNotFound(String),
#[error("no query binding registered for external '{0}'")]
UnknownQuery(String),
#[error("external '{0}' is async; drive the flow via step_one, not the exclusive driver")]
AsyncExternalUnsupported(String),
#[error("query binding system failed: {0}")]
QueryFailed(String),
#[error("system param validation failed: {0}")]
SystemParamInvalid(String),
#[error(transparent)]
Runtime(#[from] RuntimeError),
}
enum NextStep {
Done(Value),
RunQuery {
system: QuerySystemId,
qargs: Vec<Value>,
},
}
type EvalSystemState<M> = SystemState<(
Query<
'static,
'static,
(
&'static BrinkProgram<M>,
&'static BrinkLocale<M>,
&'static mut BrinkFlow<M>,
&'static mut BrinkContext<M>,
),
>,
Option<ResMut<'static, crate::BrinkGlobals<M>>>,
Res<'static, Assets<ProgramAsset>>,
Res<'static, Assets<LineTablesAsset>>,
Res<'static, BrinkBindings<M>>,
)>;
fn drive_function_eval_to_done<M: Send + Sync + 'static>(
world: &mut World,
entity: Entity,
state: &mut EvalSystemState<M>,
mut next: NextStep,
triggers: &mut Vec<TriggerFn>,
) -> Result<Value, BrinkCallError> {
loop {
match next {
NextStep::Done(value) => return Ok(value),
NextStep::RunQuery { system, qargs } => {
let value = world
.run_system_with(system, (entity, qargs))
.map_err(|e| BrinkCallError::QueryFailed(format!("{e:?}")))?;
next = {
let (mut flows, globals, programs, tables, bindings) = state
.get_mut(world)
.map_err(|e| BrinkCallError::SystemParamInvalid(e.to_string()))?;
let mut globals = globals.ok_or(BrinkCallError::NotAFlow)?;
let (prog_c, loc_c, mut flow, mut ctx) = flows
.get_mut(entity)
.map_err(|_| BrinkCallError::NotAFlow)?;
let program = &programs
.get(&prog_c.handle)
.ok_or(BrinkCallError::ProgramNotLoaded)?
.program;
let line_tables = &tables
.get(&loc_c.handle)
.ok_or(BrinkCallError::LineTablesNotLoaded)?
.tables;
let handler = bindings.eval_handler();
flow.inner.resolve_external(value);
let mut view = crate::globals::flow_context_view(&mut globals, &mut ctx);
let outcome = flow.inner.resume_function_eval::<FastRng>(
program,
line_tables,
&mut view,
&handler,
None,
)?;
triggers.extend(handler.take_queued());
classify_eval(&flow.inner, program, &bindings, outcome)?
};
}
}
}
}
fn flush_eval_triggers(world: &mut World, triggers: Vec<TriggerFn>) {
for trigger in triggers {
trigger(world);
}
}
fn classify_eval<M: Send + Sync + 'static>(
flow: &FlowInstance,
program: &Program,
bindings: &BrinkBindings<M>,
outcome: brink_runtime::FunctionEval,
) -> Result<NextStep, BrinkCallError> {
match outcome {
brink_runtime::FunctionEval::Returned(value) => Ok(NextStep::Done(value)),
brink_runtime::FunctionEval::AwaitingExternal => {
let name = flow
.pending_external_name(program)
.unwrap_or_default()
.to_owned();
if bindings.async_bindings.contains_key(&name) {
return Err(BrinkCallError::AsyncExternalUnsupported(name));
}
let system = bindings
.query(&name)
.ok_or(BrinkCallError::UnknownQuery(name))?;
let qargs = flow.pending_external_args().to_vec();
Ok(NextStep::RunQuery { system, qargs })
}
}
}
pub fn call_ink_function<M: Send + Sync + 'static>(
world: &mut World,
entity: Entity,
name: &str,
args: &[Value],
) -> Result<Value, BrinkCallError> {
let mut state: EvalSystemState<M> = SystemState::new(world);
let mut triggers: Vec<TriggerFn> = Vec::new();
let next = begin_eval_by_name(world, entity, &mut state, name, args, &mut triggers)?;
let value = drive_function_eval_to_done(world, entity, &mut state, next, &mut triggers)?;
flush_eval_triggers(world, triggers);
Ok(value)
}
fn begin_eval_by_name<M: Send + Sync + 'static>(
world: &mut World,
entity: Entity,
state: &mut EvalSystemState<M>,
name: &str,
args: &[Value],
triggers: &mut Vec<TriggerFn>,
) -> Result<NextStep, BrinkCallError> {
let (mut flows, globals, programs, tables, bindings) = state
.get_mut(world)
.map_err(|e| BrinkCallError::SystemParamInvalid(e.to_string()))?;
let mut globals = globals.ok_or(BrinkCallError::NotAFlow)?;
let (prog_c, loc_c, mut flow, mut ctx) = flows
.get_mut(entity)
.map_err(|_| BrinkCallError::NotAFlow)?;
let program = &programs
.get(&prog_c.handle)
.ok_or(BrinkCallError::ProgramNotLoaded)?
.program;
let line_tables = &tables
.get(&loc_c.handle)
.ok_or(BrinkCallError::LineTablesNotLoaded)?
.tables;
let idx = program
.find_address(name)
.ok_or_else(|| BrinkCallError::FunctionNotFound(name.to_owned()))?
.0;
let handler = bindings.eval_handler();
let mut view = crate::globals::flow_context_view(&mut globals, &mut ctx);
let outcome = flow.inner.begin_function_eval::<FastRng>(
program,
line_tables,
&mut view,
&handler,
idx,
args,
None,
)?;
triggers.extend(handler.take_queued());
classify_eval(&flow.inner, program, &bindings, outcome)
}
pub fn call_ink_functions<M, N, A>(
world: &mut World,
entity: Entity,
calls: impl IntoIterator<Item = (N, A)>,
) -> Vec<Result<Value, BrinkCallError>>
where
M: Send + Sync + 'static,
N: AsRef<str>,
A: AsRef<[Value]>,
{
let mut state: EvalSystemState<M> = SystemState::new(world);
calls
.into_iter()
.map(|(name, args)| {
let mut triggers: Vec<TriggerFn> = Vec::new();
let next = begin_eval_by_name(
world,
entity,
&mut state,
name.as_ref(),
args.as_ref(),
&mut triggers,
)?;
let value =
drive_function_eval_to_done(world, entity, &mut state, next, &mut triggers)?;
flush_eval_triggers(world, triggers);
Ok(value)
})
.collect()
}
pub fn call_ink_function_value<M: Send + Sync + 'static>(
world: &mut World,
entity: Entity,
callee: &Value,
args: &[Value],
) -> Result<Value, BrinkCallError> {
let mut state: EvalSystemState<M> = SystemState::new(world);
let mut triggers: Vec<TriggerFn> = Vec::new();
let next = {
let (mut flows, globals, programs, tables, bindings) = state
.get_mut(world)
.map_err(|e| BrinkCallError::SystemParamInvalid(e.to_string()))?;
let mut globals = globals.ok_or(BrinkCallError::NotAFlow)?;
let (prog_c, loc_c, mut flow, mut ctx) = flows
.get_mut(entity)
.map_err(|_| BrinkCallError::NotAFlow)?;
let program = &programs
.get(&prog_c.handle)
.ok_or(BrinkCallError::ProgramNotLoaded)?
.program;
let line_tables = &tables
.get(&loc_c.handle)
.ok_or(BrinkCallError::LineTablesNotLoaded)?
.tables;
let handler = bindings.eval_handler();
let mut view = crate::globals::flow_context_view(&mut globals, &mut ctx);
let outcome = flow.inner.begin_function_value_eval::<FastRng>(
program,
line_tables,
&mut view,
&handler,
callee,
args,
None,
)?;
triggers.extend(handler.take_queued());
classify_eval(&flow.inner, program, &bindings, outcome)?
};
let value = drive_function_eval_to_done(world, entity, &mut state, next, &mut triggers)?;
flush_eval_triggers(world, triggers);
Ok(value)
}
enum FlowStep {
Line(Step),
Query {
system: QuerySystemId,
qargs: Vec<Value>,
#[cfg(feature = "dev")]
name: String,
},
}
fn emit_line_event_world<M: Send + Sync + 'static>(world: &mut World, entity: Entity, step: &Step) {
use crate::event::{BrinkChoicesPresented, BrinkLineDelivered, BrinkStoryEnded, BrinkTurnDone};
match step {
Step::Line(line) => {
world
.entity_mut(entity)
.trigger(|e| BrinkLineDelivered::<M>::new(e, line.text.clone(), line.tags.clone()));
}
Step::Choices(choices) => {
world.entity_mut(entity).trigger(|e| {
BrinkChoicesPresented::<M>::new(e, String::new(), Vec::new(), choices.clone())
});
}
Step::Done | Step::Suspended => {
world
.entity_mut(entity)
.trigger(|e| BrinkTurnDone::<M>::new(e, String::new(), Vec::new()));
}
Step::End => {
world
.entity_mut(entity)
.trigger(|e| BrinkStoryEnded::<M>::new(e, String::new(), Vec::new()));
}
}
}
fn advance_recording<M: Send + Sync + 'static>(
flow: &mut FlowInstance,
program: &Program,
line_tables: &[Vec<brink_format::LineEntry>],
context: &mut (impl brink_runtime::ContextAccess + ?Sized),
handler: &BrinkHandler<'_, M>,
#[cfg(feature = "dev")] recorder: Option<&mut ReplayRecorder>,
) -> Result<StepOutcome, RuntimeError> {
#[cfg(feature = "dev")]
if let Some(rec) = recorder {
let recording = RecordingHandler::new(handler, rec);
return flow.advance::<FastRng>(program, line_tables, context, &recording, None);
}
flow.advance::<FastRng>(program, line_tables, context, handler, None)
}
#[expect(
clippy::too_many_lines,
reason = "the SystemState re-borrow dance around run_system_with doesn't split cleanly"
)]
pub fn advance_flow<M: Send + Sync + 'static>(
world: &mut World,
entity: Entity,
) -> Result<Step, BrinkCallError> {
#[expect(
clippy::type_complexity,
reason = "SystemState param tuple for the flow components + assets + bindings"
)]
let mut state: SystemState<(
Query<(
&BrinkProgram<M>,
&BrinkLocale<M>,
&mut BrinkFlow<M>,
&mut BrinkContext<M>,
)>,
Option<ResMut<crate::BrinkGlobals<M>>>,
Res<Assets<ProgramAsset>>,
Res<Assets<LineTablesAsset>>,
Res<BrinkBindings<M>>,
)> = SystemState::new(world);
let mut triggers: Vec<TriggerFn> = Vec::new();
#[cfg(feature = "dev")]
let mut recorder: Option<ReplayRecorder> = crate::replay::take_recorder::<M>(world, entity);
let mut budget = FlowInstance::LINE_LIMIT;
loop {
if budget == 0 {
return Err(RuntimeError::LineLimitExceeded(FlowInstance::LINE_LIMIT).into());
}
budget -= 1;
let step = {
let (mut flows, globals, programs, tables, bindings) = state
.get_mut(world)
.map_err(|e| BrinkCallError::SystemParamInvalid(e.to_string()))?;
let mut globals = globals.ok_or(BrinkCallError::NotAFlow)?;
let (prog_c, loc_c, mut flow, mut ctx) = flows
.get_mut(entity)
.map_err(|_| BrinkCallError::NotAFlow)?;
let program = &programs
.get(&prog_c.handle)
.ok_or(BrinkCallError::ProgramNotLoaded)?
.program;
let line_tables = &tables
.get(&loc_c.handle)
.ok_or(BrinkCallError::LineTablesNotLoaded)?
.tables;
let handler = bindings.handler();
let mut view = crate::globals::flow_context_view(&mut globals, &mut ctx);
let outcome = advance_recording(
&mut flow.inner,
program,
line_tables,
&mut view,
&handler,
#[cfg(feature = "dev")]
recorder.as_mut(),
)?;
triggers.extend(handler.take_queued());
match outcome {
StepOutcome::Step(step) => FlowStep::Line(step),
StepOutcome::AwaitingExternal => {
let name = flow
.inner
.pending_external_name(program)
.unwrap_or_default()
.to_owned();
if bindings.async_bindings.contains_key(&name) {
return Err(BrinkCallError::AsyncExternalUnsupported(name));
}
let system = bindings
.query(&name)
.ok_or_else(|| BrinkCallError::UnknownQuery(name.clone()))?;
let qargs = flow.inner.pending_external_args().to_vec();
FlowStep::Query {
system,
qargs,
#[cfg(feature = "dev")]
name,
}
}
}
};
match step {
FlowStep::Line(line) => {
for trigger in triggers {
trigger(world);
}
emit_line_event_world::<M>(world, entity, &line);
#[cfg(feature = "dev")]
if let Some(rec) = recorder {
crate::replay::put_recorder::<M>(world, entity, rec);
}
return Ok(line);
}
FlowStep::Query {
system,
qargs,
#[cfg(feature = "dev")]
name,
} => {
let value = world
.run_system_with(system, (entity, qargs.clone()))
.map_err(|e| BrinkCallError::QueryFailed(format!("{e:?}")))?;
#[cfg(feature = "dev")]
if let Some(rec) = recorder.as_mut() {
rec.record(&name, &qargs, &value);
}
let (mut flows, ..) = state
.get_mut(world)
.map_err(|e| BrinkCallError::SystemParamInvalid(e.to_string()))?;
let (_, _, mut flow, _) = flows
.get_mut(entity)
.map_err(|_| BrinkCallError::NotAFlow)?;
flow.inner.resolve_external(value);
}
}
}
}
enum Dispatch {
Nothing,
Query {
system: QuerySystemId,
qargs: Vec<Value>,
#[cfg(any(feature = "dev", feature = "effect-trace"))]
name: String,
},
FireEvent { name: String, qargs: Vec<Value> },
SpawnTask {
fut: Pin<Box<dyn Future<Output = Value> + Send>>,
#[cfg(feature = "dev")]
name: String,
#[cfg(feature = "dev")]
qargs: Vec<Value>,
},
}
#[expect(
clippy::type_complexity,
reason = "SystemState param tuple for the flow component (+ dispatch markers) + assets + bindings"
)]
fn decide_dispatch<M: Send + Sync + 'static>(world: &mut World, entity: Entity) -> Dispatch {
let mut state: SystemState<(
Query<(
&BrinkProgram<M>,
&BrinkFlow<M>,
Option<&BrinkAwaiting<M>>,
Option<&BrinkPendingTask<M>>,
)>,
Res<Assets<ProgramAsset>>,
Res<BrinkBindings<M>>,
)> = SystemState::new(world);
let Ok((flows, programs, bindings)) = state.get(world) else {
return Dispatch::Nothing;
};
let Ok((prog_c, flow, awaiting, pending_task)) = flows.get(entity) else {
return Dispatch::Nothing;
};
if !flow.inner.has_pending_external() {
return Dispatch::Nothing;
}
let Some(program) = programs.get(&prog_c.handle) else {
return Dispatch::Nothing;
};
let program = &program.program;
let name = flow
.inner
.pending_external_name(program)
.unwrap_or_default()
.to_owned();
let qargs = flow.inner.pending_external_args().to_vec();
if let Some(system) = bindings.query(&name) {
Dispatch::Query {
system,
qargs,
#[cfg(any(feature = "dev", feature = "effect-trace"))]
name,
}
} else if let Some(kind) = bindings.async_bindings.get(&name) {
match kind {
AsyncKind::Event if awaiting.is_some() => Dispatch::Nothing, AsyncKind::Event => Dispatch::FireEvent { name, qargs },
AsyncKind::Task(_) if pending_task.is_some() => Dispatch::Nothing, AsyncKind::Task(factory) => {
#[cfg(feature = "dev")]
let fut = factory(qargs.clone());
#[cfg(not(feature = "dev"))]
let fut = factory(qargs);
Dispatch::SpawnTask {
fut,
#[cfg(feature = "dev")]
name,
#[cfg(feature = "dev")]
qargs,
}
}
}
} else {
warn!("brink: flow {entity:?} parked on unbound external '{name}'");
Dispatch::Nothing
}
}
fn dispatch_one_external<M: Send + Sync + 'static>(world: &mut World, entity: Entity) {
match decide_dispatch::<M>(world, entity) {
Dispatch::Nothing => {}
Dispatch::Query {
system,
qargs,
#[cfg(any(feature = "dev", feature = "effect-trace"))]
name,
} => match world.run_system_with(system, (entity, qargs.clone())) {
Ok(value) => {
#[cfg(feature = "dev")]
crate::replay::record_external::<M>(world, entity, &name, &qargs, &value);
#[cfg(feature = "effect-trace")]
if let Some(access) = world
.get_resource::<BrinkBindings<M>>()
.and_then(|b| b.query_access(&name))
.cloned()
{
crate::ground_truth::record::<M>(world, entity, &name, access);
}
let mut flows = world.query::<&mut BrinkFlow<M>>();
if let Ok(mut flow) = flows.get_mut(world, entity) {
flow.inner.resolve_external(value);
}
}
Err(err) => warn!("brink: query binding failed on {entity:?}: {err:?}"),
},
Dispatch::FireEvent { name, qargs } => {
world
.entity_mut(entity)
.insert(BrinkAwaiting::<M>::new(name.clone()));
world
.entity_mut(entity)
.trigger(|e| BrinkExternalAwaited::<M>::new(e, name, qargs));
}
Dispatch::SpawnTask {
fut,
#[cfg(feature = "dev")]
name,
#[cfg(feature = "dev")]
qargs,
} => {
let task = AsyncComputeTaskPool::get_or_init(TaskPool::default).spawn(fut);
world.entity_mut(entity).insert(BrinkPendingTask::<M>::new(
task,
#[cfg(feature = "dev")]
name,
#[cfg(feature = "dev")]
qargs,
));
}
}
}
#[must_use]
pub fn any_flow_awaiting_external<M: Send + Sync + 'static>(flows: Query<&BrinkFlow<M>>) -> bool {
flows.iter().any(|f| f.inner.has_pending_external())
}
pub fn resolve_pending_externals<M: Send + Sync + 'static>(world: &mut World) {
let paused: Vec<Entity> = {
let mut flows = world.query::<(Entity, &BrinkFlow<M>)>();
flows
.iter(world)
.filter(|(_, f)| f.inner.has_pending_external())
.map(|(e, _)| e)
.collect()
};
for entity in paused {
dispatch_one_external::<M>(world, entity);
}
}