use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
use crate::cluster::framing::SpawnRequest;
use tokio::sync::mpsc;
pub(crate) type SpawnItem = (String, SpawnRequest, Instant);
pub(crate) struct SpawnSender {
tx: mpsc::UnboundedSender<SpawnItem>,
queue_depth: Arc<AtomicUsize>,
}
impl SpawnSender {
pub fn new(tx: mpsc::UnboundedSender<SpawnItem>, queue_depth: Arc<AtomicUsize>) -> Self {
Self { tx, queue_depth }
}
pub fn send_spawn(&self, target_node_id: &str, request: SpawnRequest) {
self.queue_depth.fetch_add(1, Ordering::Relaxed);
if self
.tx
.send((target_node_id.to_string(), request, Instant::now()))
.is_err()
{
self.queue_depth.fetch_sub(1, Ordering::Relaxed);
tracing::warn!("Spawn sender channel closed — spawn request dropped");
}
}
}