use std::path::Path;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use crate::error::BundleError;
use crate::event::{AppEvent, BundleLoadResult};
use crate::replay::{BundleReader, ResourceCeilings};
#[must_use = "register the returned JoinHandle with the Supervisor so it has a shutdown path"]
pub fn spawn_bundle_load(
dir: impl Into<std::path::PathBuf>,
ceilings: ResourceCeilings,
tx_events: mpsc::Sender<AppEvent>,
cancel: CancellationToken,
) -> JoinHandle<()> {
let dir = dir.into();
tokio::task::spawn_blocking(move || {
if cancel.is_cancelled() {
return;
}
if let Some(result) = load_bundle(&dir, ceilings, &cancel) {
let _ = tx_events.blocking_send(AppEvent::BundleLoaded(result));
}
})
}
fn load_bundle(
dir: &Path,
ceilings: ResourceCeilings,
cancel: &CancellationToken,
) -> Option<BundleLoadResult> {
let reader = match BundleReader::open_with_ceilings(dir, ceilings) {
Ok(reader) => reader,
Err(BundleError::Cancelled) => return None,
Err(error) => return Some(BundleLoadResult::Failed(error.to_string())),
};
match reader.load_cancellable(&|| cancel.is_cancelled()) {
Ok(bundle) => Some(BundleLoadResult::Loaded(Box::new(bundle))),
Err(BundleError::Cancelled) => None,
Err(error) => Some(BundleLoadResult::Failed(error.to_string())),
}
}
#[cfg(test)]
mod tests {
use std::fs;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use super::{load_bundle, spawn_bundle_load};
use crate::event::{AppEvent, BundleLoadResult};
use crate::replay::ResourceCeilings;
#[test]
fn test_load_bundle_missing_dir_is_failed_not_panic() {
let cancel = CancellationToken::new();
let dir = std::env::temp_dir().join("chainview-nonexistent-bundle-xyz-34");
let _ = fs::remove_dir_all(&dir);
match load_bundle(&dir, ResourceCeilings::default(), &cancel) {
Some(BundleLoadResult::Failed(message)) => {
assert!(!message.is_empty(), "the failure carries a message");
}
other => panic!("expected Failed for a missing dir, got {other:?}"),
}
}
#[test]
fn test_load_bundle_pre_cancelled_returns_none() {
let cancel = CancellationToken::new();
cancel.cancel();
let dir = std::env::temp_dir().join("chainview-any-bundle-dir");
let _ = load_bundle(&dir, ResourceCeilings::default(), &cancel);
}
#[tokio::test]
async fn test_spawn_bundle_load_pre_cancelled_emits_no_event() {
let (tx, mut rx) = mpsc::channel::<AppEvent>(8);
let cancel = CancellationToken::new();
cancel.cancel();
let handle = spawn_bundle_load(
std::env::temp_dir().join("chainview-bundle"),
ResourceCeilings::default(),
tx,
cancel,
);
match handle.await {
Ok(()) => {}
Err(e) => panic!("load worker join failed: {e}"),
}
assert!(
rx.try_recv().is_err(),
"a pre-cancelled load emits no event"
);
}
#[tokio::test]
async fn test_spawn_bundle_load_missing_dir_emits_failed() {
let (tx, mut rx) = mpsc::channel::<AppEvent>(8);
let cancel = CancellationToken::new();
let dir = std::env::temp_dir().join("chainview-missing-bundle-xyz-34b");
let _ = fs::remove_dir_all(&dir);
let handle = spawn_bundle_load(dir, ResourceCeilings::default(), tx, cancel);
match handle.await {
Ok(()) => {}
Err(e) => panic!("load worker join failed: {e}"),
}
match rx.try_recv() {
Ok(AppEvent::BundleLoaded(BundleLoadResult::Failed(_))) => {}
other => panic!("expected BundleLoaded(Failed), got {other:?}"),
}
}
}