diff --git a/src/worker.rs b/src/worker.rs index 6c99426..e5a465f 100644 --- a/src/worker.rs +++ b/src/worker.rs @@ -15,6 +15,29 @@ use crate::Error; use crate::extension::WorkerExtension; +/// Route a message: try local pool first, then address_map for cross-worker, +/// then inbox_registry for external receivers. +fn route_to_pool_or_remote( + pool: &mut ActorPool, + tc: &TickContext, + dest: ActorAddress, + msg: Box, +) { + if pool.contains(&dest) { + pool.deliver(&dest, msg); + } else { + match tc.address_map.lookup(&dest) { + Some(wid) => { + tc.transfer_txs[wid.as_usize()].send(Envelope::new(dest, msg)); + crate::runtime::notify_worker(tc.worker_threads, wid.as_usize()); + } + None => { + let _ = tc.inbox_registry.try_deliver(dest, msg); + } + } + } +} + // ─── Worker ───────────────────────────────────────────────────────────────── /// A worker owns a set of actors and runs them in a loop. @@ -83,23 +106,12 @@ impl Worker { let t2 = Instant::now(); // 2.5. Fire per-worker extension (e.g., timers) → deliver before tick_all - if let Some(ext) = &mut self.worker_ext { - for (dest, msg) in ext.on_tick() { - if self.pool.contains(&dest) { - self.pool.deliver(&dest, msg); - } else { - match tc.address_map.lookup(&dest) { - Some(wid) => { - tc.transfer_txs[wid.as_usize()].send(Envelope::new(dest, msg)); - crate::runtime::notify_worker(tc.worker_threads, wid.as_usize()); - } - None => { - let _ = tc.inbox_registry.try_deliver(dest, msg); - } - } - } - did_work = true; - } + let ext_msgs: Vec<_> = self.worker_ext.as_mut() + .map(|ext| ext.on_tick()) + .unwrap_or_default(); + for (dest, msg) in ext_msgs { + route_to_pool_or_remote(&mut self.pool, tc, dest, msg); + did_work = true; } // 3. Tick all actors with WorkerContext @@ -233,22 +245,9 @@ impl Worker { let dead_addrs: Vec<_> = dead.iter().map(|(a, _)| *a).collect(); ext.cleanup_dead(&dead_addrs); - // Deliver Down notifications through normal routing + // Deliver death notifications through normal routing for (dest, msg) in notifications { - if self.pool.contains(&dest) { - self.pool.deliver(&dest, msg); - } else { - match tc.address_map.lookup(&dest) { - Some(wid) => { - tc.transfer_txs[wid.as_usize()] - .send(Envelope::new(dest, msg)); - crate::runtime::notify_worker(tc.worker_threads, wid.as_usize()); - } - None => { - let _ = tc.inbox_registry.try_deliver(dest, msg); - } - } - } + route_to_pool_or_remote(&mut self.pool, tc, dest, msg); } }