feat(simulation): background worker pool infrastructure (#843 Part C)
Generic BackgroundWorkerPool<Req, Resp> with crossbeam channels, closure handlers, and 3 delivery strategies (Fallback, GracefulDegrade, ModalLock). Stub workers registered as Bevy resources: - ChunkGenWorker (2 threads) — terrain/props/navmesh generation - NpcPrepWorker (1 thread) — pre-compute NPC state for incoming areas - OffscreenTickWorker (1 thread) — advance NPCs outside active tier Tick loop integration: - PreInput: poll_worker_results drains completed work - PostSnapshot: push_worker_requests queues new work (no-op until Phase 5) Handlers are stubs — real computation plugs in when the phases that need them arrive. The infrastructure (channels, threads, push/poll, shutdown) is real and tested. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,158 @@
|
||||
//! 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
|
||||
//!
|
||||
//! TODO (#843): If a worker thread panics, in-flight requests on that thread
|
||||
//! are lost. The pool does NOT currently re-enqueue them. For Fallback-strategy
|
||||
//! workers (voice, LLM) this 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. Future: add heartbeat monitoring and
|
||||
//! automatic re-enqueue on worker death.
|
||||
|
||||
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<Req> {
|
||||
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<Req, Resp> {
|
||||
tx: Sender<WorkRequest<Req>>,
|
||||
rx: Receiver<(Resp, DeliveryStrategy)>,
|
||||
shutdown: Arc<AtomicBool>,
|
||||
_workers: Vec<JoinHandle<()>>,
|
||||
}
|
||||
|
||||
impl<Req, Resp> BackgroundWorkerPool<Req, Resp>
|
||||
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<H>(count: usize, handler: H) -> Self
|
||||
where
|
||||
H: Fn(Req) -> Resp + Send + Sync + Clone + 'static,
|
||||
{
|
||||
let (req_tx, req_rx) = crossbeam_channel::unbounded::<WorkRequest<Req>>();
|
||||
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::Relaxed) {
|
||||
match rx.recv_timeout(std::time::Duration::from_millis(100)) {
|
||||
Ok(work) => {
|
||||
let strategy = work.strategy;
|
||||
let result = handler(work.payload);
|
||||
if tx.send((result, strategy)).is_err() {
|
||||
break; // response channel closed
|
||||
}
|
||||
}
|
||||
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: workers,
|
||||
}
|
||||
}
|
||||
|
||||
/// Submit a request. Non-blocking — returns immediately.
|
||||
pub fn submit(&self, request: WorkRequest<Req>) {
|
||||
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::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
impl<Req, Resp> Drop for BackgroundWorkerPool<Req, Resp> {
|
||||
fn drop(&mut self) {
|
||||
self.shutdown.store(true, Ordering::Relaxed);
|
||||
// Workers will exit on next recv_timeout cycle (100ms max)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user