use std::sync::Weak;
use std::sync::atomic::Ordering;
use std::time::Duration;
use tokio::sync::oneshot;
use nodedb_types::Surrogate;
use super::super::registry::{RESERVE_BATCH_SIZE, SurrogateRegistry};
use super::core::SurrogateAssigner;
use crate::control::state::SharedState;
const RESERVE_WAIT_TIMEOUT: Duration = Duration::from_secs(10);
const RESERVE_LOW_WATERMARK: u64 = (RESERVE_BATCH_SIZE / 4) as u64;
const REFILL_RETRY_BACKOFF: Duration = Duration::from_millis(100);
impl SurrogateAssigner {
pub(super) fn alloc_locked(
&self,
registry: &SurrogateRegistry,
) -> crate::Result<Option<Surrogate>> {
if self.should_use_reservation() {
Ok(registry.try_alloc_reserved())
} else {
Ok(Some(registry.alloc_one()?))
}
}
pub(super) fn nudge_refill_if_low(&self, registry: &SurrogateRegistry) {
if !self.should_use_reservation() {
return;
}
if registry.remaining_reserved() < RESERVE_LOW_WATERMARK {
self.refill_notify.notify_one();
}
}
pub async fn run_refill_loop(self: std::sync::Arc<Self>, shared: Weak<SharedState>) {
let mut eager = true;
loop {
if !eager {
self.refill_notify.notified().await;
}
eager = false;
if shared.upgrade().is_none() {
tracing::debug!("surrogate refill loop exiting: SharedState dropped");
return;
}
if !self.should_use_reservation() {
continue;
}
let remaining = self
.registry
.read()
.map(|r| r.remaining_reserved())
.unwrap_or_else(|p| p.into_inner().remaining_reserved());
if remaining >= RESERVE_LOW_WATERMARK {
continue;
}
match self.ensure_batch() {
Ok(()) => {}
Err(e) => {
tracing::debug!(error = %e, "surrogate background reservation failed; retrying");
tokio::time::sleep(REFILL_RETRY_BACKOFF).await;
self.refill_notify.notify_one();
}
}
}
}
pub(super) fn should_use_reservation(&self) -> bool {
if self.reservation_latched.load(Ordering::Relaxed) {
return true;
}
let Some(shared) = self.shared.get().and_then(|w| w.upgrade()) else {
return false;
};
if shared.metadata_raft.get().is_none() {
return false;
}
if let Some(topology) = shared.cluster_topology.as_ref() {
let member_count = match topology.read() {
Ok(guard) => guard
.all_nodes()
.filter(|node| node.state.receives_log())
.count(),
Err(_poisoned) => return false,
};
if member_count > 1 {
self.reservation_latched.store(true, Ordering::Relaxed);
return true;
}
}
let Some(routing) = shared.cluster_routing.as_ref() else {
return false;
};
let member_count = match routing.read() {
Ok(guard) => guard
.group_info(nodedb_cluster::METADATA_GROUP_ID)
.map(|info| info.members.len())
.unwrap_or(0),
Err(_poisoned) => return false,
};
let multi = member_count > 1;
if multi {
self.reservation_latched.store(true, Ordering::Relaxed);
}
multi
}
pub fn enable_reservation_mode(&self) {
self.reservation_latched.store(true, Ordering::Relaxed);
}
pub(super) fn ensure_batch(&self) -> crate::Result<()> {
let shared =
self.shared
.get()
.and_then(|w| w.upgrade())
.ok_or_else(|| crate::Error::Internal {
detail: "surrogate reserve: SharedState unavailable in cluster mode".into(),
})?;
let handle = tokio::runtime::Handle::current();
let _gate = tokio::task::block_in_place(|| handle.block_on(self.reserve_gate.lock()));
if self
.registry
.read()
.map(|r| r.has_reserved())
.unwrap_or_else(|p| p.into_inner().has_reserved())
{
return Ok(());
}
let request_id = self.next_request_id.fetch_add(1, Ordering::Relaxed);
let (tx, rx) = oneshot::channel();
if let Ok(mut pending) = self.pending_reservations.lock() {
pending.insert(request_id, tx);
} else {
return Err(crate::Error::Internal {
detail: "surrogate reserve: pending map poisoned".into(),
});
}
let propose_result = crate::control::metadata_proposer::propose_surrogate_reserve(
&shared,
shared.node_id,
request_id,
RESERVE_BATCH_SIZE,
);
if let Err(e) = propose_result {
if let Ok(mut pending) = self.pending_reservations.lock() {
pending.remove(&request_id);
}
return Err(crate::Error::Internal {
detail: format!("surrogate reserve propose failed: {e}"),
});
}
let wait = tokio::task::block_in_place(|| {
handle.block_on(async { tokio::time::timeout(RESERVE_WAIT_TIMEOUT, rx).await })
});
match wait {
Ok(Ok((_start, _end))) => Ok(()),
Ok(Err(_recv_err)) => {
if let Ok(mut pending) = self.pending_reservations.lock() {
pending.remove(&request_id);
}
Err(crate::Error::Internal {
detail: "surrogate reserve: completion signal dropped before apply".into(),
})
}
Err(_timeout) => {
if let Ok(mut pending) = self.pending_reservations.lock() {
pending.remove(&request_id);
}
Err(crate::Error::Internal {
detail: "surrogate reserve: timed out waiting for batch apply".into(),
})
}
}
}
pub fn complete_reservation(&self, request_id: u64, start: u32, end: u32) {
if let Ok(mut pending) = self.pending_reservations.lock()
&& let Some(tx) = pending.remove(&request_id)
{
if let Ok(reg) = self.registry.read() {
reg.set_reserved_batch(start, end);
}
let _ = tx.send((start, end));
}
}
}