//! Generic background worker pool. //! //! Wraps crossbeam channels + spawned threads. The handler closure captures //! any context it needs (world seed, config, etc.) at construction time. //! //! ## Worker crash behavior //! //! Handler panics are contained per-request (T-1063): the worker thread //! catches the unwind, logs an error, and keeps serving the queue — a //! poisoned request cannot permanently shrink the pool. The in-flight //! request is still LOST: no result is ever delivered for it. //! //! TODO (#843): re-enqueue lost requests (needs heartbeat/ack bookkeeping). //! For Fallback-strategy workers (voice, LLM) the loss is acceptable — the //! tick loop uses the default. For ModalLock-strategy workers (chunk gen on //! gate transit), a lost request means the player waits forever. Before any //! ModalLock consumer ships (Phase 5), implement re-enqueue or downgrade //! ModalLock to an error path. use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::thread::{self, JoinHandle}; use crossbeam_channel::{Receiver, Sender, TryRecvError}; /// How the tick loop handles a not-yet-ready result. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum DeliveryStrategy { /// Use default/cached value. Player never waits. Seamless. Fallback, /// Area loads partially, elements pop in when ready. Visible but non-blocking. GracefulDegrade, /// Tick loop pauses until result arrives. Diegetic buffer UI shown to player. /// Use sparingly — only for gate transit, initial world load, etc. ModalLock, } /// A request with its delivery strategy attached. #[derive(Debug)] pub struct WorkRequest { pub payload: Req, pub strategy: DeliveryStrategy, } /// A generic background worker pool. /// /// `Req` and `Resp` are the request/response types. The handler closure /// runs on worker threads and transforms requests into responses. /// /// ```ignore /// let pool = BackgroundWorkerPool::spawn(2, |req: ChunkGenRequest| { /// // runs on worker thread — has access to captured world_seed /// generate_chunk(req, world_seed) /// }); /// pool.submit(WorkRequest { payload: req, strategy: DeliveryStrategy::GracefulDegrade }); /// for result in pool.drain() { /* incorporate */ } /// ``` pub struct BackgroundWorkerPool { tx: Sender>, rx: Receiver<(Resp, DeliveryStrategy)>, shutdown: Arc, workers: Vec>, } impl BackgroundWorkerPool where Req: Send + 'static, Resp: Send + 'static, { /// Spawn `count` worker threads with the given handler closure. /// /// The handler is cloned to each thread. It should capture any context /// it needs (world seed, config registries, etc.) via `Arc` or `Clone`. pub fn spawn(count: usize, handler: H) -> Self where H: Fn(Req) -> Resp + Send + Sync + Clone + 'static, { let (req_tx, req_rx) = crossbeam_channel::unbounded::>(); let (resp_tx, resp_rx) = crossbeam_channel::unbounded::<(Resp, DeliveryStrategy)>(); let shutdown = Arc::new(AtomicBool::new(false)); let mut workers = Vec::with_capacity(count); for id in 0..count { let rx = req_rx.clone(); let tx = resp_tx.clone(); let shutdown = shutdown.clone(); let handler = handler.clone(); let handle = thread::Builder::new() .name(format!("worker-{id}")) .spawn(move || { while !shutdown.load(Ordering::SeqCst) { match rx.recv_timeout(std::time::Duration::from_millis(100)) { Ok(work) => { let strategy = work.strategy; // Contain handler panics (T-1063): the thread must // survive a poisoned request or the pool permanently // loses capacity. The in-flight request is still lost // (no result delivered) — re-enqueue is future work, // see the module-level TODO (#843). let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { handler(work.payload) })); match result { Ok(resp) => { if tx.send((resp, strategy)).is_err() { break; // response channel closed } } Err(_) => { tracing::error!( strategy = ?strategy, "worker handler panicked — in-flight request lost, \ no result will be delivered (#843)" ); } } } Err(crossbeam_channel::RecvTimeoutError::Timeout) => continue, Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break, } } }) .expect("failed to spawn worker thread"); workers.push(handle); } Self { tx: req_tx, rx: resp_rx, shutdown, workers, } } /// Submit a request. Non-blocking — returns immediately. pub fn submit(&self, request: WorkRequest) { let _ = self.tx.send(request); } /// Poll for a single completed result. Non-blocking. pub fn poll(&self) -> Option<(Resp, DeliveryStrategy)> { match self.rx.try_recv() { Ok(result) => Some(result), Err(TryRecvError::Empty) => None, Err(TryRecvError::Disconnected) => None, } } /// Drain all completed results. Non-blocking. pub fn drain(&self) -> Vec<(Resp, DeliveryStrategy)> { let mut results = Vec::new(); while let Ok(result) = self.rx.try_recv() { results.push(result); } results } /// Check if any results are ready without consuming them. pub fn has_results(&self) -> bool { !self.rx.is_empty() } /// Number of pending requests in the work queue. pub fn pending_count(&self) -> usize { self.tx.len() } /// Signal all workers to shut down gracefully. pub fn shutdown(&self) { self.shutdown.store(true, Ordering::SeqCst); } } impl Drop for BackgroundWorkerPool { fn drop(&mut self) { self.shutdown.store(true, Ordering::SeqCst); // Join all worker threads — wait for clean exit (100ms max per thread) for handle in self.workers.drain(..) { let _ = handle.join(); } } } // --------------------------------------------------------------------------- // Tests (T-1063): spawn/push/poll/shutdown lifecycle + worker-panic behavior. // // Every test also exercises Drop implicitly — a pool whose threads fail to // exit would hang the test binary at scope end. // --------------------------------------------------------------------------- #[cfg(test)] mod tests { use super::*; use std::time::{Duration, Instant}; const WAIT: Duration = Duration::from_secs(5); /// Spin until `cond` holds, panicking after `WAIT`. Worker completion is /// timing-dependent (thread scheduling), so tests wait on observable /// state instead of sleeping fixed amounts. fn wait_for(what: &str, mut cond: impl FnMut() -> bool) { let start = Instant::now(); while !cond() { if start.elapsed() > WAIT { panic!("timed out after {WAIT:?} waiting for: {what}"); } thread::sleep(Duration::from_millis(5)); } } #[test] fn spawn_push_poll_round_trip() { let pool = BackgroundWorkerPool::spawn(2, |x: u32| x * 2); assert!(pool.poll().is_none(), "fresh pool has no results"); assert!(!pool.has_results()); assert_eq!(pool.pending_count(), 0); pool.submit(WorkRequest { payload: 21, strategy: DeliveryStrategy::Fallback, }); wait_for("result ready", || pool.has_results()); let (resp, strategy) = pool.poll().expect("has_results implies poll succeeds"); assert_eq!(resp, 42); assert_eq!(strategy, DeliveryStrategy::Fallback); assert!(pool.poll().is_none(), "single request yields single result"); assert_eq!(pool.pending_count(), 0); } #[test] fn drain_collects_all_results_and_preserves_strategies() { let pool = BackgroundWorkerPool::spawn(2, |x: u32| x + 1); for i in 0..10u32 { pool.submit(WorkRequest { payload: i, strategy: if i % 2 == 0 { DeliveryStrategy::GracefulDegrade } else { DeliveryStrategy::ModalLock }, }); } let mut results: Vec<(u32, DeliveryStrategy)> = Vec::new(); wait_for("all 10 results drained", || { results.extend(pool.drain()); results.len() == 10 }); // Two worker threads — completion order is not guaranteed. Each // request must come back exactly once, paired with its own strategy. results.sort_by_key(|(resp, _)| *resp); for (i, (resp, strategy)) in results.iter().enumerate() { assert_eq!(*resp, i as u32 + 1); let expected = if i % 2 == 0 { DeliveryStrategy::GracefulDegrade } else { DeliveryStrategy::ModalLock }; assert_eq!(*strategy, expected, "strategy must travel with request {i}"); } assert!(!pool.has_results(), "drain leaves the queue empty"); } #[test] fn shutdown_stops_workers_and_drops_later_submissions() { let pool = BackgroundWorkerPool::spawn(2, |x: u32| x); pool.shutdown(); wait_for("all workers exited", || { pool.workers.iter().all(|h| h.is_finished()) }); // Once every worker has exited, all Receiver clones are gone (spawn // drops the original), so the request channel is disconnected: // submit() silently discards the request — by design (`let _ =` on // send). Nothing is queued and no result ever appears. pool.submit(WorkRequest { payload: 7, strategy: DeliveryStrategy::Fallback, }); assert_eq!( pool.pending_count(), 0, "post-shutdown submissions are silently dropped, not queued" ); assert!(pool.poll().is_none()); assert!(!pool.has_results()); } // -- Worker-panic behavior (T-1063) ---------------------------------------- /// Handler that panics on payload 0 and echoes anything else. fn poison_handler(x: u32) -> u32 { if x == 0 { panic!("poisoned request"); } x } #[test] fn handler_panic_does_not_kill_the_worker_thread() { // Single thread: if the panic killed it, the second request would // never be processed and this test would time out. let pool = BackgroundWorkerPool::spawn(1, poison_handler); pool.submit(WorkRequest { payload: 0, // poison strategy: DeliveryStrategy::ModalLock, }); pool.submit(WorkRequest { payload: 5, strategy: DeliveryStrategy::Fallback, }); wait_for("post-panic result", || pool.has_results()); let (resp, _) = pool.poll().expect("worker survived the panic"); assert_eq!(resp, 5); assert!( !pool.workers[0].is_finished(), "worker thread must survive a handler panic" ); } /// KNOWN LIMITATION (#843): the panicked request is silently lost — no /// result is ever delivered for it and it is not re-enqueued. For a /// ModalLock consumer this means the player waits forever; heartbeat + /// re-enqueue (or downgrading ModalLock to an error path) must land /// before any ModalLock consumer ships (Phase 5). This test pins the /// current loss behavior so the eventual fix has to update it visibly. #[test] fn handler_panic_loses_the_inflight_request() { let pool = BackgroundWorkerPool::spawn(1, poison_handler); pool.submit(WorkRequest { payload: 0, // poison — would be the ModalLock request a player waits on strategy: DeliveryStrategy::ModalLock, }); // Sentinel: once its result arrives, the poison request has // definitely been taken off the queue (single worker, FIFO channel). pool.submit(WorkRequest { payload: 99, strategy: DeliveryStrategy::Fallback, }); wait_for("sentinel result", || pool.has_results()); let results = pool.drain(); assert_eq!( results.len(), 1, "only the sentinel completes — the poisoned request produced no result" ); assert_eq!(results[0].0, 99); assert_eq!( pool.pending_count(), 0, "the poisoned request is gone from the queue, not re-enqueued" ); } }