use crate::timing::{SharedTimer, TimerPhase};
use livesplit_auto_splitting::{
CreationError, InterruptHandle, Runtime as ScriptRuntime, SettingsStore,
Timer as AutoSplitTimer, TimerState,
};
use snafu::Snafu;
use std::{fmt, fs, io, path::PathBuf, thread, time::Duration};
use tokio::{
runtime,
sync::{mpsc, oneshot, watch},
time::{timeout_at, Instant},
};
#[derive(Debug, Snafu)]
pub enum Error {
ThreadStopped,
LoadFailed {
source: CreationError,
},
ReadFileFailed {
source: io::Error,
},
}
pub struct Runtime {
interrupt_receiver: watch::Receiver<Option<InterruptHandle>>,
sender: mpsc::UnboundedSender<Request>,
}
impl Drop for Runtime {
fn drop(&mut self) {
if let Some(handle) = &*self.interrupt_receiver.borrow() {
handle.interrupt();
}
}
}
impl Runtime {
pub fn new(timer: SharedTimer) -> Self {
let (sender, receiver) = mpsc::unbounded_channel();
let (interrupt_sender, interrupt_receiver) = watch::channel(None);
let (timeout_sender, timeout_receiver) = watch::channel(None);
thread::Builder::new()
.name("Auto Splitting Runtime".into())
.spawn(move || {
runtime::Builder::new_current_thread()
.enable_time()
.build()
.unwrap()
.block_on(run(receiver, timer, timeout_sender, interrupt_sender))
})
.unwrap();
thread::Builder::new()
.name("Auto Splitting Watchdog".into())
.spawn({
let interrupt_receiver = interrupt_receiver.clone();
move || {
runtime::Builder::new_current_thread()
.enable_time()
.build()
.unwrap()
.block_on(watchdog(timeout_receiver, interrupt_receiver))
}
})
.unwrap();
Self {
interrupt_receiver,
sender,
}
}
pub async fn load_script(&self, script: PathBuf) -> Result<(), Error> {
let (sender, receiver) = oneshot::channel();
let script = fs::read(script).map_err(|e| Error::ReadFileFailed { source: e })?;
self.sender
.send(Request::LoadScript(script, sender))
.map_err(|_| Error::ThreadStopped)?;
receiver.await.map_err(|_| Error::ThreadStopped)??;
Ok(())
}
pub fn load_script_blocking(&self, script: PathBuf) -> Result<(), Error> {
runtime::Builder::new_current_thread()
.enable_time()
.build()
.unwrap()
.block_on(self.load_script(script))
}
pub async fn unload_script(&self) -> Result<(), Error> {
let (sender, receiver) = oneshot::channel();
self.sender
.send(Request::UnloadScript(sender))
.map_err(|_| Error::ThreadStopped)?;
receiver.await.map_err(|_| Error::ThreadStopped)
}
pub fn unload_script_blocking(&self) -> Result<(), Error> {
runtime::Builder::new_current_thread()
.enable_time()
.build()
.unwrap()
.block_on(self.unload_script())
}
}
enum Request {
LoadScript(Vec<u8>, oneshot::Sender<Result<(), Error>>),
UnloadScript(oneshot::Sender<()>),
}
struct Timer(SharedTimer);
impl AutoSplitTimer for Timer {
fn state(&self) -> TimerState {
match self.0.read().unwrap().current_phase() {
TimerPhase::NotRunning => TimerState::NotRunning,
TimerPhase::Running => TimerState::Running,
TimerPhase::Paused => TimerState::Paused,
TimerPhase::Ended => TimerState::Ended,
}
}
fn start(&mut self) {
self.0.write().unwrap().start()
}
fn split(&mut self) {
self.0.write().unwrap().split()
}
fn reset(&mut self) {
self.0.write().unwrap().reset(true)
}
fn set_game_time(&mut self, time: time::Duration) {
self.0.write().unwrap().set_game_time(time.into());
}
fn pause_game_time(&mut self) {
self.0.write().unwrap().pause_game_time()
}
fn resume_game_time(&mut self) {
self.0.write().unwrap().resume_game_time()
}
fn set_variable(&mut self, name: &str, value: &str) {
self.0.write().unwrap().set_custom_variable(name, value)
}
fn log(&mut self, message: fmt::Arguments<'_>) {
log::info!(target: "Auto Splitter", "{message}");
}
}
async fn run(
mut receiver: mpsc::UnboundedReceiver<Request>,
timer: SharedTimer,
timeout_sender: watch::Sender<Option<Instant>>,
interrupt_sender: watch::Sender<Option<InterruptHandle>>,
) {
'back_to_not_having_a_runtime: loop {
interrupt_sender.send(None).ok();
timeout_sender.send(None).ok();
let mut runtime = loop {
match receiver.recv().await {
Some(Request::LoadScript(script, ret)) => {
match ScriptRuntime::new(&script, Timer(timer.clone()), SettingsStore::new()) {
Ok(r) => {
ret.send(Ok(())).ok();
break r;
}
Err(source) => {
ret.send(Err(Error::LoadFailed { source })).ok();
}
};
}
Some(Request::UnloadScript(ret)) => {
log::warn!(target: "Auto Splitter", "Attempted to unload already unloaded script");
ret.send(()).ok();
}
None => {
return;
}
};
};
log::info!(target: "Auto Splitter", "Loaded script");
let mut next_step = Instant::now();
interrupt_sender.send(Some(runtime.interrupt_handle())).ok();
timeout_sender.send(Some(next_step)).ok();
loop {
match timeout_at(next_step, receiver.recv()).await {
Ok(Some(request)) => match request {
Request::LoadScript(script, ret) => {
match ScriptRuntime::new(
&script,
Timer(timer.clone()),
SettingsStore::new(),
) {
Ok(r) => {
ret.send(Ok(())).ok();
runtime = r;
log::info!(target: "Auto Splitter", "Reloaded script");
}
Err(source) => {
ret.send(Err(Error::LoadFailed { source })).ok();
log::info!(target: "Auto Splitter", "Failed to load");
}
}
}
Request::UnloadScript(ret) => {
ret.send(()).ok();
log::info!(target: "Auto Splitter", "Unloaded script");
continue 'back_to_not_having_a_runtime;
}
},
Ok(None) => return,
Err(_) => match runtime.update() {
Ok(tick_rate) => {
next_step += tick_rate;
timeout_sender.send(Some(next_step)).ok();
}
Err(e) => {
log::error!(target: "Auto Splitter", "Unloaded due to failure: {:?}", e);
continue 'back_to_not_having_a_runtime;
}
},
}
}
}
}
async fn watchdog(
mut timeout_receiver: watch::Receiver<Option<Instant>>,
interrupt_receiver: watch::Receiver<Option<InterruptHandle>>,
) {
const TIMEOUT: Duration = Duration::from_secs(5);
loop {
let instant = *timeout_receiver.borrow();
match instant {
Some(time) => match timeout_at(time + TIMEOUT, timeout_receiver.changed()).await {
Ok(Ok(_)) => {}
Ok(Err(_)) => return,
Err(_) => {
if let Some(handle) = &*interrupt_receiver.borrow() {
handle.interrupt();
}
}
},
None => {
if timeout_receiver.changed().await.is_err() {
return;
}
}
}
}
}