mod api;
mod convert;
mod runtime;
mod worker;
mod world;
pub(crate) use paneru::lua::shared;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use bevy::app::{App, Plugin, PostUpdate, PreUpdate, Update};
use bevy::ecs::message::MessageReader;
use bevy::ecs::resource::Resource;
use bevy::ecs::schedule::IntoScheduleConfigs;
use bevy::ecs::system::{Commands, NonSendMut, Query, Res, ResMut};
use notify::Watcher;
use crate::commands::Command;
use crate::config::Config;
use crate::ecs::params::Windows;
use crate::ecs::script_state::ScriptStateStore;
use crate::ecs::state::QueryStateParams;
use crate::ecs::{SendMessageTrigger, SpawnCommandsExt, apply_config_side_effects};
use crate::events::Event;
use crate::manager::{Application, Display, WindowManager};
use crate::util::symlink_target;
use worker::FromLua;
pub use worker::{LuaSource, LuaWorker};
#[derive(Resource, Debug, Clone)]
pub struct LuaScriptPath(pub PathBuf);
const MISSING_STORE: &str = "the script state store is not available";
pub struct LuaPlugin {}
impl Plugin for LuaPlugin {
fn build(&self, app: &mut App) {
app.add_systems(
PreUpdate,
(
serve_lua_queries.before(crate::ecs::systems::pump_events),
serve_lua_store.before(crate::ecs::systems::pump_events),
drain_lua_outbox,
command_lua_handler,
),
);
app.add_systems(PostUpdate, (serve_lua_queries, serve_lua_store));
app.add_systems(Update, (dispatch_lua_events, lua_reload_system));
}
}
pub fn dispatch_lua_events(worker: Option<Res<LuaWorker>>, mut reader: MessageReader<Event>) {
let Some(worker) = worker else {
return;
};
let mask = worker.subscribed_event_mask();
if mask == 0 {
for _ in reader.read() {}
return;
}
let events: Vec<convert::LuaEvent> = reader
.read()
.filter_map(|event| {
let lua_event = convert::LuaEvent::try_from(event).ok()?;
(mask & convert::LuaEvent::bit_for_name(lua_event.name()) != 0).then_some(lua_event)
})
.collect();
if events.is_empty() {
return;
}
worker.send_events(events);
}
pub fn command_lua_handler(worker: Option<Res<LuaWorker>>, mut reader: MessageReader<Event>) {
let Some(worker) = worker else {
return;
};
let ids: Vec<u32> = reader
.read()
.filter_map(|event| match event {
Event::Command {
command: Command::Lua(id),
} => Some(*id),
_ => None,
})
.collect();
if ids.is_empty() {
return;
}
worker.send_binds(ids);
}
pub fn serve_lua_queries(worker: Option<Res<LuaWorker>>, state: QueryStateParams) {
let Some(worker) = worker else {
return;
};
let requests: Vec<_> = worker.pending_world_queries().collect();
if requests.is_empty() {
return;
}
let mut extracted_state = None;
let mut extracted_set = None;
for request in requests {
match request {
worker::WorldRequest::State { reply } => {
let _ = reply.try_send(extract_once(&mut extracted_state, || state.extract()));
}
worker::WorldRequest::WindowSet { reply } => {
let _ = reply.try_send(extract_once(&mut extracted_set, || {
state.extract_window_set()
}));
}
}
}
}
pub fn serve_lua_store(
worker: Option<Res<LuaWorker>>,
mut script_state: Option<ResMut<ScriptStateStore>>,
) {
let Some(worker) = worker else {
return;
};
for request in worker.pending_store_queries() {
match request {
worker::StoreRequest::Read { reply } => {
let answer = script_state
.as_ref()
.map(|store| store.snapshot())
.ok_or_else(|| MISSING_STORE.to_string());
let _ = reply.try_send(answer);
}
worker::StoreRequest::Write { write, reply } => {
let answer = script_state.as_mut().map_or_else(
|| Err(MISSING_STORE.to_string()),
|store| store.apply(&write),
);
let _ = reply.try_send(answer);
}
}
}
}
fn extract_once<T>(
slot: &mut Option<worker::Shared<T>>,
extract: impl FnOnce() -> crate::errors::Result<T>,
) -> worker::Shared<T> {
slot.get_or_insert_with(|| extract().map(Arc::new).map_err(|err| err.to_string()))
.clone()
}
pub fn drain_lua_outbox(
worker: Option<Res<LuaWorker>>,
config: Option<Res<Config>>,
mut displays: Query<&mut Display>,
windows: Windows,
applications: Query<&Application>,
mut commands: Commands,
) {
let Some(worker) = worker else {
return;
};
for effect in worker.drain_outbox() {
match effect {
FromLua::Command(command) => {
commands.trigger(SendMessageTrigger(Event::Command { command }));
}
FromLua::Flash { message, duration } => commands.flash_message(message, duration),
FromLua::ConfigChanged => {
if let (Some(config), Some(built)) = (config.as_ref(), worker.built_config()) {
config.replace_inner_from(&built);
apply_config_side_effects(config, &mut displays, &windows, &applications);
}
}
}
}
}
pub fn lua_reload_system(
worker: Option<Res<LuaWorker>>,
script_path: Option<Res<LuaScriptPath>>,
mut reader: MessageReader<Event>,
window_manager: Res<WindowManager>,
mut watcher: Option<NonSendMut<Box<dyn Watcher>>>,
) {
let (Some(worker), Some(script_path)) = (worker, script_path) else {
return;
};
let path = &script_path.0;
let mut should_reload = false;
for event in reader.read() {
let Event::ConfigRefresh(event) = event else {
continue;
};
if event.paths.iter().any(|changed| paths_match(changed, path)) {
if let (Some(watcher), Some(_symlink)) = (watcher.as_mut(), symlink_target(path))
&& let Some(new_watcher) = crate::ecs::rewatch_configs(&window_manager, path)
{
**watcher = new_watcher;
}
should_reload = true;
}
}
if !should_reload {
return;
}
worker.send_reload(path.clone());
}
fn paths_match(changed: &Path, script: &Path) -> bool {
changed == script || changed.file_name() == script.file_name()
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::Cell;
#[test]
fn one_extraction_answers_every_waiter() {
let reads = Cell::new(0);
let mut slot = None;
let answers: Vec<_> = (0..5)
.map(|_| {
extract_once(&mut slot, || {
reads.set(reads.get() + 1);
Ok(7_u32)
})
})
.collect();
assert_eq!(reads.get(), 1, "the world should be read once for all five");
for answer in &answers {
assert_eq!(*answer.as_ref().expect("a successful read"), Arc::new(7));
}
let first = answers[0].as_ref().expect("a successful read");
assert!(
answers[1..]
.iter()
.all(|other| Arc::ptr_eq(first, other.as_ref().expect("a successful read"))),
"every waiter should hold the same extraction"
);
}
#[test]
fn a_failed_extraction_is_shared_not_retried() {
let reads = Cell::new(0);
let mut slot = None;
let answers: Vec<_> = (0..3)
.map(|_| {
extract_once(&mut slot, || -> crate::errors::Result<u32> {
reads.set(reads.get() + 1);
Err(crate::errors::Error::InvalidInput("no world".to_string()))
})
})
.collect();
assert_eq!(reads.get(), 1, "a failure should not be retried per waiter");
assert!(answers.iter().all(std::result::Result::is_err));
}
}