use std::collections::HashMap;
use std::future::Future;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use mongreldb_core::{
ClassConfig, GovernorAction, HierarchicalScheduler, MemoryClass, MemoryError, MemoryGovernor,
NodeMemoryGovernor, NodePressureInputs, Reservation, ResourceGroupRegistry, SchedulerError,
SchedulerStats, WorkItem, WorkloadClass,
};
use mongreldb_types::ids::QueryId;
use tokio::sync::oneshot;
#[derive(Debug, Clone)]
pub struct AdmitRequest<'a> {
pub tenant: &'a str,
pub class: WorkloadClass,
pub priority: u8,
pub deadline: Option<Duration>,
pub query_id: Option<QueryId>,
pub tag: &'a str,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[allow(dead_code)] pub struct AdmissionMetrics {
pub running_by_class: std::collections::BTreeMap<String, usize>,
pub queued_by_class: std::collections::BTreeMap<String, usize>,
pub parent_reserved_bytes: u64,
pub child_reserved_bytes: u64,
pub open_parents: usize,
pub open_children: usize,
}
struct ParentBudget {
budget_bytes: u64,
used_bytes: u64,
children: usize,
}
#[derive(Clone)]
pub struct NodeAdmissionController {
scheduler: SchedulerAdmission,
memory: MemoryGovernor,
parents: Arc<Mutex<HashMap<u64, ParentBudget>>>,
parent_reserved_bytes: Arc<AtomicU64>,
child_reserved_bytes: Arc<AtomicU64>,
open_children: Arc<AtomicU64>,
}
impl NodeAdmissionController {
pub fn new(groups: &ResourceGroupRegistry, memory: MemoryGovernor) -> Self {
Self {
scheduler: SchedulerAdmission::from_resource_groups(groups),
memory,
parents: Arc::new(Mutex::new(HashMap::new())),
parent_reserved_bytes: Arc::new(AtomicU64::new(0)),
child_reserved_bytes: Arc::new(AtomicU64::new(0)),
open_children: Arc::new(AtomicU64::new(0)),
}
}
pub fn scheduler(&self) -> &SchedulerAdmission {
&self.scheduler
}
#[allow(dead_code)] pub fn memory(&self) -> &MemoryGovernor {
&self.memory
}
#[allow(dead_code)] pub fn metrics(&self) -> AdmissionMetrics {
let stats = self.scheduler.stats();
let mut running_by_class = std::collections::BTreeMap::new();
let mut queued_by_class = std::collections::BTreeMap::new();
for (name, class_stats) in &stats.per_class {
running_by_class.insert(name.clone(), class_stats.running);
queued_by_class.insert(name.clone(), class_stats.queued);
}
let parents = self
.parents
.lock()
.unwrap_or_else(|error| error.into_inner());
AdmissionMetrics {
running_by_class,
queued_by_class,
parent_reserved_bytes: self.parent_reserved_bytes.load(Ordering::Relaxed),
child_reserved_bytes: self.child_reserved_bytes.load(Ordering::Relaxed),
open_parents: parents.len(),
open_children: self.open_children.load(Ordering::Relaxed) as usize,
}
}
pub async fn admit<C>(
&self,
req: AdmitRequest<'_>,
cancel: C,
) -> Result<AdmittedWork, AdmitError>
where
C: Future<Output = ()>,
{
self.scheduler.submit_and_wait(req, cancel).await
}
pub async fn admit_parent<C>(
&self,
req: AdmitRequest<'_>,
memory_class: MemoryClass,
budget_bytes: u64,
cancel: C,
) -> Result<ParentAdmission, AdmitError>
where
C: Future<Output = ()>,
{
let work = self.admit(req, cancel).await?;
let reservation = self
.memory
.try_reserve(budget_bytes, memory_class)
.map_err(AdmitError::Memory)?;
self.parent_reserved_bytes
.fetch_add(budget_bytes, Ordering::Relaxed);
let work_id = work.work_id();
self.parents
.lock()
.unwrap_or_else(|error| error.into_inner())
.insert(
work_id,
ParentBudget {
budget_bytes,
used_bytes: 0,
children: 0,
},
);
Ok(ParentAdmission {
work,
reservation,
controller: self.clone(),
work_id,
budget_bytes,
})
}
pub fn reserve_child(
&self,
parent: &ParentAdmission,
memory_class: MemoryClass,
bytes: u64,
) -> Result<ChildReservation, AdmitError> {
{
let mut parents = self
.parents
.lock()
.unwrap_or_else(|error| error.into_inner());
let budget = parents
.get_mut(&parent.work_id)
.ok_or(AdmitError::UnknownParent {
work_id: parent.work_id,
})?;
let next = budget.used_bytes.saturating_add(bytes);
if next > budget.budget_bytes {
return Err(AdmitError::ChildExceedsParent {
requested: bytes,
parent_remaining: budget.budget_bytes.saturating_sub(budget.used_bytes),
});
}
budget.used_bytes = next;
budget.children = budget.children.saturating_add(1);
}
self.child_reserved_bytes
.fetch_add(bytes, Ordering::Relaxed);
self.open_children.fetch_add(1, Ordering::Relaxed);
Ok(ChildReservation {
controller: self.clone(),
parent_work_id: parent.work_id,
memory_class,
bytes,
released: false,
})
}
pub fn admit_child(
&self,
parent: &ParentAdmission,
memory_class: MemoryClass,
bytes: u64,
) -> Result<ChildReservation, AdmitError> {
self.reserve_child(parent, memory_class, bytes)
}
fn release_child(&self, parent_work_id: u64, bytes: u64) {
self.child_reserved_bytes
.fetch_sub(bytes, Ordering::Relaxed);
self.open_children.fetch_sub(1, Ordering::Relaxed);
if let Some(budget) = self
.parents
.lock()
.unwrap_or_else(|error| error.into_inner())
.get_mut(&parent_work_id)
{
budget.used_bytes = budget.used_bytes.saturating_sub(bytes);
budget.children = budget.children.saturating_sub(1);
}
}
fn release_parent(&self, work_id: u64, budget_bytes: u64) {
self.parents
.lock()
.unwrap_or_else(|error| error.into_inner())
.remove(&work_id);
self.parent_reserved_bytes
.fetch_sub(budget_bytes, Ordering::Relaxed);
}
}
pub struct ParentAdmission {
work: AdmittedWork,
reservation: Reservation,
controller: NodeAdmissionController,
work_id: u64,
budget_bytes: u64,
}
impl ParentAdmission {
#[allow(dead_code)] pub fn work_id(&self) -> u64 {
self.work_id
}
#[allow(dead_code)] pub fn remaining_bytes(&self) -> u64 {
let parents = self
.controller
.parents
.lock()
.unwrap_or_else(|error| error.into_inner());
parents
.get(&self.work_id)
.map(|b| b.budget_bytes.saturating_sub(b.used_bytes))
.unwrap_or(0)
}
#[allow(dead_code)] pub fn child_used_bytes(&self) -> u64 {
let parents = self
.controller
.parents
.lock()
.unwrap_or_else(|error| error.into_inner());
parents
.get(&self.work_id)
.map(|b| b.used_bytes)
.unwrap_or(0)
}
}
impl Drop for ParentAdmission {
fn drop(&mut self) {
self.controller
.release_parent(self.work_id, self.budget_bytes);
let _ = &self.reservation;
let _ = &self.work;
}
}
#[allow(dead_code)] pub struct ChildReservation {
controller: NodeAdmissionController,
parent_work_id: u64,
memory_class: MemoryClass,
bytes: u64,
released: bool,
}
impl ChildReservation {
#[allow(dead_code)]
pub fn bytes(&self) -> u64 {
self.bytes
}
#[allow(dead_code)]
pub fn memory_class(&self) -> MemoryClass {
self.memory_class
}
#[allow(dead_code)] pub fn release(mut self) {
self.finish();
}
fn finish(&mut self) {
if self.released {
return;
}
self.released = true;
self.controller
.release_child(self.parent_work_id, self.bytes);
}
}
impl Drop for ChildReservation {
fn drop(&mut self) {
self.finish();
}
}
#[derive(Clone)]
pub struct SchedulerAdmission {
inner: Arc<SchedulerAdmissionInner>,
pressure: Arc<PressureControl>,
}
struct SchedulerAdmissionInner {
state: Mutex<AdmissionState>,
}
struct AdmissionState {
scheduler: HierarchicalScheduler,
waiters: HashMap<u64, oneshot::Sender<WorkItem>>,
}
#[derive(Debug, Clone)]
struct ClassBaselines {
interactive: ClassConfig,
ai: ClassConfig,
analytics: ClassConfig,
}
#[derive(Debug)]
pub struct PressureControl {
reject_ai: AtomicBool,
reduced: AtomicBool,
move_tablet_noops: AtomicU64,
last_evict_bytes: AtomicU64,
evaluate_count: AtomicU64,
baselines: Mutex<ClassBaselines>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PressureSnapshot {
pub reject_ai: bool,
pub reduced_admission: bool,
pub move_tablet_noops: u64,
pub last_evict_bytes: u64,
pub evaluate_count: u64,
}
impl PressureControl {
fn new(baselines: ClassBaselines) -> Self {
Self {
reject_ai: AtomicBool::new(false),
reduced: AtomicBool::new(false),
move_tablet_noops: AtomicU64::new(0),
last_evict_bytes: AtomicU64::new(0),
evaluate_count: AtomicU64::new(0),
baselines: Mutex::new(baselines),
}
}
pub fn snapshot(&self) -> PressureSnapshot {
PressureSnapshot {
reject_ai: self.reject_ai.load(Ordering::Relaxed),
reduced_admission: self.reduced.load(Ordering::Relaxed),
move_tablet_noops: self.move_tablet_noops.load(Ordering::Relaxed),
last_evict_bytes: self.last_evict_bytes.load(Ordering::Relaxed),
evaluate_count: self.evaluate_count.load(Ordering::Relaxed),
}
}
pub fn reject_ai(&self) -> bool {
self.reject_ai.load(Ordering::Relaxed)
}
fn lock_baselines(&self) -> std::sync::MutexGuard<'_, ClassBaselines> {
self.baselines
.lock()
.unwrap_or_else(|error| error.into_inner())
}
}
#[derive(Debug, Clone)]
pub struct PressureInputSources {
pub db_reserved_bytes: u64,
pub db_max_bytes: u64,
pub node_configured_max_bytes: u64,
pub tablet_reserved_bytes: u64,
pub ai_capacity: usize,
pub ai_available: usize,
pub process_rss_bytes: Option<u64>,
}
pub fn process_rss_bytes() -> Option<u64> {
let status = std::fs::read_to_string("/proc/self/status").ok()?;
for line in status.lines() {
let Some(rest) = line.strip_prefix("VmRSS:") else {
continue;
};
let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
return Some(kb.saturating_mul(1024));
}
None
}
pub fn build_pressure_inputs(src: &PressureInputSources) -> NodePressureInputs {
let configured_max_bytes = src.db_max_bytes.max(src.node_configured_max_bytes).max(1);
let query_reserved_bytes = src
.db_reserved_bytes
.saturating_add(src.tablet_reserved_bytes);
let ai_util = if src.ai_capacity == 0 {
0.0
} else {
let used = src.ai_capacity.saturating_sub(src.ai_available);
(used as f64 / src.ai_capacity as f64).clamp(0.0, 1.0)
};
let (physical_memory_bytes, rss_pressure) = match src.process_rss_bytes {
Some(rss) if rss > 0 => {
let physical = rss.saturating_mul(4).max(configured_max_bytes);
let p = (rss as f64 / physical as f64).clamp(0.0, 1.0);
(physical, p)
}
_ => (NodePressureInputs::default().physical_memory_bytes, 0.0),
};
let mut os_pressure = rss_pressure.max(ai_util);
if let Ok(forced) = std::env::var("MONGRELDB_NODE_GOVERNOR_FORCE_OS_PRESSURE") {
if let Ok(value) = forced.parse::<f64>() {
os_pressure = os_pressure.max(value.clamp(0.0, 1.0));
}
}
NodePressureInputs {
physical_memory_bytes,
configured_max_bytes,
os_pressure,
cache_hit_rate: 1.0,
query_reserved_bytes,
compaction_backlog_bytes: 0,
replication_backlog_bytes: 0,
}
}
pub fn refresh_pressure(
governor: &mut NodeMemoryGovernor,
inputs: &NodePressureInputs,
admission: &SchedulerAdmission,
cache_governor: Option<&MemoryGovernor>,
) -> Vec<GovernorAction> {
let actions = governor.evaluate(inputs);
apply_governor_actions(admission, &actions, cache_governor);
admission
.pressure
.evaluate_count
.fetch_add(1, Ordering::Relaxed);
actions
}
pub fn apply_governor_actions(
admission: &SchedulerAdmission,
actions: &[GovernorAction],
cache_governor: Option<&MemoryGovernor>,
) {
let reject_ai = actions
.iter()
.any(|a| matches!(a, GovernorAction::RejectOversizedAi));
let reduce = actions
.iter()
.any(|a| matches!(a, GovernorAction::ReduceAdmission));
let evict = actions
.iter()
.any(|a| matches!(a, GovernorAction::EvictCaches));
let move_leaders = actions
.iter()
.any(|a| matches!(a, GovernorAction::MoveTabletLeaders { .. }));
admission
.pressure
.reject_ai
.store(reject_ai, Ordering::Relaxed);
if reduce {
apply_reduce_admission(admission);
} else {
clear_reduce_admission(admission);
}
if evict {
if let Some(cache) = cache_governor {
let budget = cache.reclaimable_bytes().max(cache.max_bytes() / 16).max(1);
let freed = cache.evict_reclaimable(budget);
admission
.pressure
.last_evict_bytes
.store(freed, Ordering::Relaxed);
}
}
if move_leaders {
let n = admission
.pressure
.move_tablet_noops
.fetch_add(1, Ordering::Relaxed)
+ 1;
if n == 1 || n.is_multiple_of(100) {
eprintln!(
"[node_governor] MoveTabletLeaders requested; no-op outside cluster mode (count={n})"
);
}
}
}
fn apply_reduce_admission(admission: &SchedulerAdmission) {
if admission.pressure.reduced.swap(true, Ordering::SeqCst) {
return;
}
let (interactive, ai) = {
let baselines = admission.pressure.lock_baselines();
let mut interactive = baselines.interactive.clone();
let mut ai = baselines.ai.clone();
interactive.max_concurrency = (interactive.max_concurrency / 2).max(1);
ai.max_concurrency = (ai.max_concurrency / 2).max(1);
(interactive, ai)
};
{
let mut state = admission.inner.lock();
state
.scheduler
.set_class_config(WorkloadClass::InteractiveSql, interactive);
state
.scheduler
.set_class_config(WorkloadClass::AiRetrieval, ai);
}
}
fn clear_reduce_admission(admission: &SchedulerAdmission) {
if !admission.pressure.reduced.swap(false, Ordering::SeqCst) {
return;
}
let (interactive, ai) = {
let baselines = admission.pressure.lock_baselines();
(baselines.interactive.clone(), baselines.ai.clone())
};
let mut state = admission.inner.lock();
state
.scheduler
.set_class_config(WorkloadClass::InteractiveSql, interactive);
state
.scheduler
.set_class_config(WorkloadClass::AiRetrieval, ai);
}
pub struct AdmittedWork {
work_id: u64,
inner: Arc<SchedulerAdmissionInner>,
completed: bool,
}
impl std::fmt::Debug for AdmittedWork {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AdmittedWork")
.field("work_id", &self.work_id)
.field("completed", &self.completed)
.finish()
}
}
pub struct SqlAdmissionGuard {
_permit: tokio::sync::OwnedSemaphorePermit,
#[allow(dead_code)] parent: ParentAdmission,
}
impl SqlAdmissionGuard {
pub fn new(permit: tokio::sync::OwnedSemaphorePermit, parent: ParentAdmission) -> Self {
Self {
_permit: permit,
parent,
}
}
#[allow(dead_code)] pub fn parent(&self) -> &ParentAdmission {
&self.parent
}
#[allow(dead_code)] pub fn work_id(&self) -> u64 {
self.parent.work_id()
}
}
impl AdmittedWork {
#[allow(dead_code)] pub fn work_id(&self) -> u64 {
self.work_id
}
#[allow(dead_code)] pub fn complete(mut self) {
self.finish();
}
fn finish(&mut self) {
if self.completed {
return;
}
self.completed = true;
let mut state = self.inner.lock();
let _ = state.scheduler.complete(self.work_id);
dispatch_ready(&mut state);
}
}
impl Drop for AdmittedWork {
fn drop(&mut self) {
self.finish();
}
}
impl SchedulerAdmission {
pub fn new() -> Self {
Self::from_resource_groups(&ResourceGroupRegistry::with_defaults())
}
pub fn from_resource_groups(groups: &ResourceGroupRegistry) -> Self {
let mut scheduler = HierarchicalScheduler::new();
apply_resource_groups(&mut scheduler, groups);
let baselines = ClassBaselines {
interactive: resolved_class_config(groups, WorkloadClass::InteractiveSql),
ai: resolved_class_config(groups, WorkloadClass::AiRetrieval),
analytics: resolved_class_config(groups, WorkloadClass::Analytics),
};
Self {
inner: Arc::new(SchedulerAdmissionInner {
state: Mutex::new(AdmissionState {
scheduler,
waiters: HashMap::new(),
}),
}),
pressure: Arc::new(PressureControl::new(baselines)),
}
}
pub fn pressure(&self) -> &PressureControl {
&self.pressure
}
pub fn check_ai_admitted(&self) -> Result<(), AdmitError> {
if self.pressure.reject_ai() {
Err(AdmitError::PressureRejected {
resource: "ai_memory_pressure",
})
} else {
Ok(())
}
}
#[allow(dead_code)] pub fn set_class_config(&self, class: WorkloadClass, config: ClassConfig) {
{
let mut state = self.inner.lock();
state.scheduler.set_class_config(class, config.clone());
}
if !self.pressure.reduced.load(Ordering::Relaxed) {
let mut baselines = self.pressure.lock_baselines();
match class {
WorkloadClass::InteractiveSql => baselines.interactive = config,
WorkloadClass::AiRetrieval => baselines.ai = config,
WorkloadClass::Analytics => baselines.analytics = config,
_ => {}
}
}
}
pub fn stats(&self) -> SchedulerStats {
self.inner.lock().scheduler.stats()
}
pub async fn submit_and_wait<C>(
&self,
req: AdmitRequest<'_>,
cancel: C,
) -> Result<AdmittedWork, AdmitError>
where
C: Future<Output = ()>,
{
if matches!(
req.class,
WorkloadClass::AiRetrieval | WorkloadClass::Analytics
) {
self.check_ai_admitted()?;
}
let (work_id, rx) = {
let mut state = self.inner.lock();
let work_id = state
.scheduler
.submit(
req.tenant,
req.class,
req.priority,
req.deadline,
req.query_id,
req.tag,
)
.map_err(AdmitError::Rejected)?;
let (tx, rx) = oneshot::channel();
state.waiters.insert(work_id, tx);
dispatch_ready(&mut state);
(work_id, rx)
};
tokio::pin!(cancel);
tokio::select! {
biased;
result = rx => {
match result {
Ok(_item) => Ok(AdmittedWork {
work_id,
inner: Arc::clone(&self.inner),
completed: false,
}),
Err(_) => {
self.cancel_work(work_id);
Err(AdmitError::Cancelled)
}
}
}
_ = &mut cancel => {
self.cancel_work(work_id);
Err(AdmitError::Cancelled)
}
}
}
pub fn cancel_work(&self, work_id: u64) {
let mut state = self.inner.lock();
state.waiters.remove(&work_id);
let _ = state.scheduler.cancel(work_id);
let _ = state.scheduler.complete(work_id);
dispatch_ready(&mut state);
}
}
impl Default for SchedulerAdmission {
fn default() -> Self {
Self::new()
}
}
impl SchedulerAdmissionInner {
fn lock(&self) -> std::sync::MutexGuard<'_, AdmissionState> {
self.state.lock().unwrap_or_else(|error| error.into_inner())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AdmitError {
Rejected(SchedulerError),
Cancelled,
PressureRejected {
resource: &'static str,
},
Memory(MemoryError),
ChildExceedsParent {
requested: u64,
parent_remaining: u64,
},
UnknownParent { work_id: u64 },
}
pub fn scheduler_error_to_query(error: SchedulerError) -> mongreldb_query::MongrelQueryError {
let (resource, requested, limit) = match &error {
SchedulerError::QueueFull { depth, max, .. } => ("scheduler_queue", *depth + 1, *max),
SchedulerError::TenantQuota { .. } => ("tenant_quota", 1, 0),
SchedulerError::UnknownWork(_) => ("scheduler", 1, 0),
};
mongreldb_query::MongrelQueryError::Core(mongreldb_core::MongrelError::ResourceLimitExceeded {
resource,
requested,
limit,
})
}
pub fn admit_error_to_query(error: AdmitError) -> mongreldb_query::MongrelQueryError {
match error {
AdmitError::Rejected(error) => scheduler_error_to_query(error),
AdmitError::Cancelled => mongreldb_query::MongrelQueryError::InvalidQueryState(
"SQL admission cancelled while waiting for scheduler slot".into(),
),
AdmitError::PressureRejected { resource } => mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource,
requested: 1,
limit: 0,
},
),
AdmitError::Memory(MemoryError::Exhausted {
requested,
available,
..
}) => mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "node_memory",
requested: requested as usize,
limit: available as usize,
},
),
AdmitError::Memory(MemoryError::LowPriorityRejected { .. }) => {
mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "node_memory_low_priority",
requested: 1,
limit: 0,
},
)
}
AdmitError::Memory(MemoryError::InvalidConfig(_)) => {
mongreldb_query::MongrelQueryError::Core(mongreldb_core::MongrelError::Other(
"invalid memory governor configuration".into(),
))
}
AdmitError::ChildExceedsParent {
requested,
parent_remaining,
} => mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "parent_memory_budget",
requested: requested as usize,
limit: parent_remaining as usize,
},
),
AdmitError::UnknownParent { work_id } => mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::Other(format!("unknown parent admission {work_id}")),
),
}
}
pub fn admit_error_to_core(error: AdmitError) -> mongreldb_core::MongrelError {
match error {
AdmitError::Rejected(error) => match scheduler_error_to_query(error) {
mongreldb_query::MongrelQueryError::Core(error) => error,
error => mongreldb_core::MongrelError::Other(error.to_string()),
},
AdmitError::Cancelled => mongreldb_core::MongrelError::Cancelled,
AdmitError::PressureRejected { resource } => {
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource,
requested: 1,
limit: 0,
}
}
other => match admit_error_to_query(other) {
mongreldb_query::MongrelQueryError::Core(error) => error,
error => mongreldb_core::MongrelError::Other(error.to_string()),
},
}
}
pub fn priority_for_class(groups: &ResourceGroupRegistry, class: WorkloadClass) -> u8 {
groups
.get(class.name())
.map(|g| g.priority)
.unwrap_or_else(|| match class {
WorkloadClass::Control => 255,
WorkloadClass::Replication => 254,
WorkloadClass::Oltp => 200,
WorkloadClass::InteractiveSql => 180,
WorkloadClass::AiRetrieval => 150,
WorkloadClass::Analytics => 100,
WorkloadClass::Maintenance => 50,
WorkloadClass::Backup => 40,
})
}
fn resolved_class_config(groups: &ResourceGroupRegistry, class: WorkloadClass) -> ClassConfig {
let mut config = ClassConfig::for_class(class);
if let Some(group) = groups.get(class.name()) {
config.max_concurrency = group.max_concurrency;
config.max_queue = group.max_queue;
config.weight = group.cpu_weight.max(1);
if class.has_reserved_capacity() {
config.reserved_slots = config.reserved_slots.max(1).min(config.max_concurrency);
}
}
if class == WorkloadClass::InteractiveSql {
if let Some(v) = positive_env_usize("MONGRELDB_SCHEDULER_INTERACTIVE_SQL_MAX_QUEUE") {
config.max_queue = v;
}
if let Some(v) = positive_env_usize("MONGRELDB_SCHEDULER_INTERACTIVE_SQL_MAX_CONCURRENCY") {
config.max_concurrency = v;
}
}
config
}
fn apply_resource_groups(scheduler: &mut HierarchicalScheduler, groups: &ResourceGroupRegistry) {
for class in WorkloadClass::ALL {
scheduler.set_class_config(class, resolved_class_config(groups, class));
}
}
fn positive_env_usize(name: &str) -> Option<usize> {
std::env::var(name)
.ok()
.and_then(|value| value.parse().ok())
.filter(|value| *value > 0)
}
fn dispatch_ready(state: &mut AdmissionState) {
loop {
let ready = state.scheduler.poll(32);
if ready.is_empty() {
break;
}
for item in ready {
let work_id = item.work_id;
match state.waiters.remove(&work_id) {
Some(tx) => {
if tx.send(item).is_err() {
let _ = state.scheduler.complete(work_id);
}
}
None => {
let _ = state.scheduler.complete(work_id);
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
fn tiny_interactive() -> ClassConfig {
ClassConfig {
max_queue: 1,
weight: 64,
reserved_slots: 0,
max_concurrency: 1,
}
}
fn sql_req(tag: &str) -> AdmitRequest<'_> {
AdmitRequest {
tenant: "t",
class: WorkloadClass::InteractiveSql,
priority: 180,
deadline: None,
query_id: None,
tag,
}
}
#[tokio::test]
async fn queue_full_is_resource_exhausted_mapping() {
let admission = SchedulerAdmission::new();
admission.set_class_config(WorkloadClass::InteractiveSql, tiny_interactive());
let never = std::future::pending::<()>();
let first = admission
.submit_and_wait(sql_req("a"), never)
.await
.expect("first admits");
let admission2 = admission.clone();
let second = tokio::spawn(async move {
admission2
.submit_and_wait(sql_req("b"), std::future::pending::<()>())
.await
});
let queued = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let stats = admission.stats();
let sql = stats.per_class.get("interactive_sql").unwrap();
if sql.queued >= 1 {
break;
}
tokio::task::yield_now().await;
}
})
.await;
assert!(
queued.is_ok(),
"second request must enqueue behind the holder"
);
let err = admission
.submit_and_wait(sql_req("c"), std::future::pending::<()>())
.await
.expect_err("third must be rejected");
let rejected = match err {
AdmitError::Rejected(e) => e,
other => panic!("expected QueueFull, got {other:?}"),
};
assert!(matches!(rejected, SchedulerError::QueueFull { max: 1, .. }));
let mapped = scheduler_error_to_query(rejected);
assert!(matches!(
mapped,
mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "scheduler_queue",
..
}
)
));
assert_eq!(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "scheduler_queue",
requested: 1,
limit: 1,
}
.category(),
mongreldb_types::errors::ErrorCategory::ResourceExhausted
);
drop(first);
let second = tokio::time::timeout(Duration::from_secs(2), second)
.await
.expect("second must admit after holder completes")
.expect("join")
.expect("second admit");
drop(second);
}
#[tokio::test]
async fn control_admits_when_interactive_sql_saturated() {
let admission = SchedulerAdmission::new();
admission.set_class_config(WorkloadClass::InteractiveSql, tiny_interactive());
admission.set_class_config(
WorkloadClass::Control,
ClassConfig {
max_queue: 8,
weight: 256,
reserved_slots: 2,
max_concurrency: 8,
},
);
let _sql = admission
.submit_and_wait(sql_req("sql"), std::future::pending::<()>())
.await
.unwrap();
let control = admission
.submit_and_wait(
AdmitRequest {
tenant: "system",
class: WorkloadClass::Control,
priority: 255,
deadline: None,
query_id: None,
tag: "ctl",
},
std::future::pending::<()>(),
)
.await
.expect("control must admit under interactive saturation");
assert!(control.work_id() > 0);
control.complete();
}
#[tokio::test]
async fn cancel_while_waiting_frees_queue_slot() {
let admission = SchedulerAdmission::new();
admission.set_class_config(WorkloadClass::InteractiveSql, tiny_interactive());
let _holder = admission
.submit_and_wait(sql_req("hold"), std::future::pending::<()>())
.await
.unwrap();
let cancelled = AtomicBool::new(false);
let (cancel_tx, cancel_rx) = oneshot::channel::<()>();
let admission2 = admission.clone();
let waiter = tokio::spawn(async move {
let result = admission2
.submit_and_wait(sql_req("wait"), async {
let _ = cancel_rx.await;
})
.await;
cancelled.store(true, Ordering::SeqCst);
result
});
tokio::time::sleep(Duration::from_millis(20)).await;
let _ = cancel_tx.send(());
let result = waiter.await.unwrap();
assert!(matches!(result, Err(AdmitError::Cancelled)));
let stats = admission.stats();
let sql = stats.per_class.get("interactive_sql").unwrap();
assert_eq!(sql.queued, 0, "cancelled work must leave the queue");
assert_eq!(sql.running, 1, "holder still running");
}
#[tokio::test]
async fn concurrent_waiters_receive_own_work_ids() {
let admission = SchedulerAdmission::new();
admission.set_class_config(
WorkloadClass::InteractiveSql,
ClassConfig {
max_queue: 16,
weight: 64,
reserved_slots: 0,
max_concurrency: 2,
},
);
let a = admission.clone();
let b = admission.clone();
let (wa, wb) = tokio::join!(
a.submit_and_wait(sql_req("a"), std::future::pending::<()>()),
b.submit_and_wait(sql_req("b"), std::future::pending::<()>()),
);
let wa = wa.unwrap();
let wb = wb.unwrap();
assert_ne!(wa.work_id(), wb.work_id());
drop(wa);
drop(wb);
}
fn high_pressure_inputs(max: u64) -> NodePressureInputs {
NodePressureInputs {
configured_max_bytes: max,
query_reserved_bytes: (max as f64 * 0.95) as u64,
os_pressure: 0.95,
..NodePressureInputs::default()
}
}
fn calm_inputs(max: u64) -> NodePressureInputs {
NodePressureInputs {
configured_max_bytes: max,
query_reserved_bytes: max / 100,
os_pressure: 0.0,
..NodePressureInputs::default()
}
}
#[test]
fn high_pressure_rejects_ai_and_reduces_admission() {
let admission = SchedulerAdmission::new();
admission.set_class_config(
WorkloadClass::InteractiveSql,
ClassConfig {
max_queue: 64,
weight: 64,
reserved_slots: 0,
max_concurrency: 16,
},
);
admission.set_class_config(
WorkloadClass::AiRetrieval,
ClassConfig {
max_queue: 64,
weight: 32,
reserved_slots: 0,
max_concurrency: 16,
},
);
let mut gov = NodeMemoryGovernor::new(
mongreldb_core::MemoryGovernor::new(mongreldb_core::GovernorConfig::new(1_000_000))
.unwrap(),
);
let actions =
refresh_pressure(&mut gov, &high_pressure_inputs(1_000_000), &admission, None);
assert!(
actions
.iter()
.any(|a| matches!(a, GovernorAction::RejectOversizedAi)),
"actions={actions:?}"
);
assert!(
actions
.iter()
.any(|a| matches!(a, GovernorAction::ReduceAdmission)),
"actions={actions:?}"
);
assert!(
actions
.iter()
.any(|a| matches!(a, GovernorAction::MoveTabletLeaders { .. })),
"actions={actions:?}"
);
let snap = admission.pressure().snapshot();
assert!(snap.reject_ai);
assert!(snap.reduced_admission);
assert!(
snap.move_tablet_noops >= 1,
"tablet move must be no-op counted"
);
assert!(snap.evaluate_count >= 1);
assert!(matches!(
admission.check_ai_admitted(),
Err(AdmitError::PressureRejected {
resource: "ai_memory_pressure"
})
));
let calm_actions = refresh_pressure(&mut gov, &calm_inputs(1_000_000), &admission, None);
assert!(!calm_actions
.iter()
.any(|a| matches!(a, GovernorAction::RejectOversizedAi)));
let snap = admission.pressure().snapshot();
assert!(!snap.reject_ai);
assert!(!snap.reduced_admission);
assert!(admission.check_ai_admitted().is_ok());
}
#[tokio::test]
async fn high_pressure_submit_and_wait_rejects_ai_class() {
let admission = SchedulerAdmission::new();
let mut gov = NodeMemoryGovernor::new(
mongreldb_core::MemoryGovernor::new(mongreldb_core::GovernorConfig::new(1_000_000))
.unwrap(),
);
refresh_pressure(&mut gov, &high_pressure_inputs(1_000_000), &admission, None);
let err = admission
.submit_and_wait(
AdmitRequest {
tenant: "t",
class: WorkloadClass::AiRetrieval,
priority: 150,
deadline: None,
query_id: None,
tag: "ai",
},
std::future::pending::<()>(),
)
.await
.expect_err("AI class must be rejected under RejectOversizedAi");
assert!(matches!(
err,
AdmitError::PressureRejected {
resource: "ai_memory_pressure"
}
));
let mapped = admit_error_to_query(err);
assert!(matches!(
mapped,
mongreldb_query::MongrelQueryError::Core(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "ai_memory_pressure",
..
}
)
));
assert_eq!(
mongreldb_core::MongrelError::ResourceLimitExceeded {
resource: "ai_memory_pressure",
requested: 1,
limit: 0,
}
.category(),
mongreldb_types::errors::ErrorCategory::ResourceExhausted
);
let sql = admission
.submit_and_wait(sql_req("sql"), std::future::pending::<()>())
.await
.expect("InteractiveSql must still admit under ReduceAdmission");
drop(sql);
}
#[test]
fn reduce_admission_halves_interactive_concurrency() {
let admission = SchedulerAdmission::new();
admission.set_class_config(
WorkloadClass::InteractiveSql,
ClassConfig {
max_queue: 4,
weight: 64,
reserved_slots: 0,
max_concurrency: 4,
},
);
let mut gov = NodeMemoryGovernor::new(
mongreldb_core::MemoryGovernor::new(mongreldb_core::GovernorConfig::new(1_000_000))
.unwrap(),
);
let inputs = NodePressureInputs {
configured_max_bytes: 1_000_000,
query_reserved_bytes: 820_000,
os_pressure: 0.82,
..NodePressureInputs::default()
};
let actions = refresh_pressure(&mut gov, &inputs, &admission, None);
assert!(actions
.iter()
.any(|a| matches!(a, GovernorAction::ReduceAdmission)));
assert!(admission.pressure().snapshot().reduced_admission);
clear_reduce_admission(&admission);
assert!(!admission.pressure().snapshot().reduced_admission);
refresh_pressure(&mut gov, &inputs, &admission, None);
assert!(admission.pressure().snapshot().reduced_admission);
refresh_pressure(&mut gov, &calm_inputs(1_000_000), &admission, None);
assert!(!admission.pressure().snapshot().reduced_admission);
}
#[test]
fn build_pressure_inputs_uses_db_and_ai_proxy() {
let inputs = build_pressure_inputs(&PressureInputSources {
db_reserved_bytes: 500,
db_max_bytes: 1000,
node_configured_max_bytes: 2000,
tablet_reserved_bytes: 50,
ai_capacity: 4,
ai_available: 0,
process_rss_bytes: None,
});
assert_eq!(inputs.query_reserved_bytes, 550);
assert_eq!(inputs.configured_max_bytes, 2000);
assert!((inputs.os_pressure - 1.0).abs() < f64::EPSILON);
}
#[test]
fn evict_caches_drives_memory_governor() {
use mongreldb_core::memory::{GovernorConfig, MemoryClass, MemoryGovernor};
use std::sync::atomic::AtomicU64 as StdAtomicU64;
struct FakeCache {
reclaimable: StdAtomicU64,
}
impl mongreldb_core::memory::Reclaimable for FakeCache {
fn evict_reclaimable(&self, budget: u64) -> u64 {
let have = self.reclaimable.load(Ordering::Relaxed);
let take = have.min(budget);
self.reclaimable.fetch_sub(take, Ordering::Relaxed);
take
}
fn reclaimable_bytes(&self) -> u64 {
self.reclaimable.load(Ordering::Relaxed)
}
}
let gov =
MemoryGovernor::new(GovernorConfig::new(1_000_000).with_reserved_floor(0)).unwrap();
let cache = Arc::new(FakeCache {
reclaimable: StdAtomicU64::new(10_000),
});
gov.register_reclaimable(&cache);
let _res = gov.try_reserve(100, MemoryClass::PageCache).unwrap();
let admission = SchedulerAdmission::new();
let mut node = NodeMemoryGovernor::new(gov.clone());
let actions = refresh_pressure(
&mut node,
&high_pressure_inputs(1_000_000),
&admission,
Some(&gov),
);
assert!(actions
.iter()
.any(|a| matches!(a, GovernorAction::EvictCaches)));
assert!(
admission.pressure().snapshot().last_evict_bytes > 0,
"evict must free reclaimable bytes"
);
assert_eq!(cache.reclaimable.load(Ordering::Relaxed), 0);
}
fn test_controller() -> NodeAdmissionController {
use mongreldb_core::{GovernorConfig, MemoryGovernor};
let memory =
MemoryGovernor::new(GovernorConfig::new(1_000_000).with_reserved_floor(100_000))
.unwrap();
NodeAdmissionController::new(&ResourceGroupRegistry::with_defaults(), memory)
}
#[test]
fn p11_x3_compaction_priority_yields_to_oltp() {
let groups = ResourceGroupRegistry::with_defaults();
let oltp = priority_for_class(&groups, WorkloadClass::Oltp);
let maintenance = priority_for_class(&groups, WorkloadClass::Maintenance);
let backup = priority_for_class(&groups, WorkloadClass::Backup);
assert!(
oltp > maintenance,
"OLTP priority {oltp} must exceed maintenance/compaction {maintenance}"
);
assert!(
maintenance >= backup,
"maintenance should not rank below backup: {maintenance} vs {backup}"
);
assert!(mongreldb_core::MemoryClass::Compaction.is_low_priority());
}
#[tokio::test]
async fn ai_overload_does_not_consume_control_reserve() {
let controller = test_controller();
controller.scheduler().set_class_config(
WorkloadClass::AiRetrieval,
ClassConfig {
max_queue: 1,
weight: 32,
reserved_slots: 0,
max_concurrency: 1,
},
);
controller.scheduler().set_class_config(
WorkloadClass::Control,
ClassConfig {
max_queue: 8,
weight: 256,
reserved_slots: 2,
max_concurrency: 8,
},
);
let _ai = controller
.admit(
AdmitRequest {
tenant: "t",
class: WorkloadClass::AiRetrieval,
priority: 150,
deadline: None,
query_id: None,
tag: "ai-hold",
},
std::future::pending::<()>(),
)
.await
.unwrap();
let controller2 = controller.clone();
let _queued = tokio::spawn(async move {
controller2
.admit(
AdmitRequest {
tenant: "t",
class: WorkloadClass::AiRetrieval,
priority: 150,
deadline: None,
query_id: None,
tag: "ai-queued",
},
std::future::pending::<()>(),
)
.await
});
tokio::time::timeout(Duration::from_secs(2), async {
loop {
let m = controller.metrics();
if m.queued_by_class.get("ai_retrieval").copied().unwrap_or(0) >= 1 {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("ai must queue");
let overflow = controller
.admit(
AdmitRequest {
tenant: "t",
class: WorkloadClass::AiRetrieval,
priority: 150,
deadline: None,
query_id: None,
tag: "ai-overflow",
},
std::future::pending::<()>(),
)
.await;
assert!(matches!(
overflow,
Err(AdmitError::Rejected(SchedulerError::QueueFull { .. }))
));
let control = controller
.admit(
AdmitRequest {
tenant: "system",
class: WorkloadClass::Control,
priority: 255,
deadline: None,
query_id: None,
tag: "ctl",
},
std::future::pending::<()>(),
)
.await
.expect("control reserve must survive AI overload");
let metrics = controller.metrics();
assert_eq!(
metrics
.running_by_class
.get("control")
.copied()
.unwrap_or(0),
1
);
assert_eq!(
metrics
.running_by_class
.get("ai_retrieval")
.copied()
.unwrap_or(0),
1
);
drop(control);
}
#[test]
fn snapshot_install_cannot_exceed_node_reserve() {
use mongreldb_core::{GovernorConfig, MemoryClass, MemoryGovernor};
let memory =
MemoryGovernor::new(GovernorConfig::new(1_000_000).with_reserved_floor(100_000))
.unwrap();
let _hold = memory
.try_reserve(900_000, MemoryClass::AiCandidates)
.expect("non-reserved can fill up to max-floor");
let err = memory
.try_reserve(50_000, MemoryClass::Backup)
.expect_err("snapshot/backup must not exceed non-reserved capacity");
assert!(
matches!(err, mongreldb_core::MemoryError::Exhausted { .. })
|| matches!(err, mongreldb_core::MemoryError::LowPriorityRejected { .. }),
"unexpected: {err:?}"
);
let huge = memory.try_reserve(2_000_000, MemoryClass::Backup);
assert!(huge.is_err(), "snapshot install above node max must fail");
let install = memory
.try_reserve(50_000, MemoryClass::Replication)
.expect("replication reserve survives snapshot pressure");
assert_eq!(install.bytes(), 50_000);
}
#[test]
fn rss_above_configured_maximum_triggers_pressure_actions() {
use mongreldb_core::{
GovernorConfig, MemoryGovernor, NodeMemoryGovernor, NodePressureInputs,
};
let memory = MemoryGovernor::new(GovernorConfig::new(1_000_000)).unwrap();
let mut node = NodeMemoryGovernor::new(memory);
let sources = PressureInputSources {
db_reserved_bytes: 100_000,
db_max_bytes: 1_000_000,
node_configured_max_bytes: 1_000_000,
tablet_reserved_bytes: 0,
ai_capacity: 4,
ai_available: 4,
process_rss_bytes: Some(950_000),
};
let inputs = build_pressure_inputs(&sources);
assert!(
inputs.os_pressure >= 0.20,
"high RSS must raise os_pressure: {inputs:?}"
);
let hot = NodePressureInputs {
physical_memory_bytes: 1_000_000,
configured_max_bytes: 1_000_000,
os_pressure: 0.95,
cache_hit_rate: 1.0,
query_reserved_bytes: 900_000,
compaction_backlog_bytes: 0,
replication_backlog_bytes: 0,
};
let actions = node.evaluate(&hot);
assert!(
actions
.iter()
.any(|a| matches!(a, GovernorAction::RejectOversizedAi)),
"RSS/pressure near max must reject oversized AI: {actions:?}"
);
assert!(
actions
.iter()
.any(|a| matches!(a, GovernorAction::ReduceAdmission)),
"must reduce admission under high RSS: {actions:?}"
);
}
#[tokio::test]
async fn tenant_quota_blocks_noisy_tenant_on_controller() {
use mongreldb_core::TenantQuota;
use std::collections::BTreeMap;
let controller = test_controller();
{
let admission = controller.scheduler();
admission.set_class_config(
WorkloadClass::Analytics,
ClassConfig {
max_queue: 1,
weight: 32,
reserved_slots: 0,
max_concurrency: 1,
},
);
}
let _hold = controller
.admit(
AdmitRequest {
tenant: "noisy",
class: WorkloadClass::Analytics,
priority: 50,
deadline: None,
query_id: None,
tag: "hold",
},
std::future::pending::<()>(),
)
.await
.unwrap();
let controller2 = controller.clone();
let _queued = tokio::spawn(async move {
controller2
.admit(
AdmitRequest {
tenant: "noisy",
class: WorkloadClass::Analytics,
priority: 50,
deadline: None,
query_id: None,
tag: "queued",
},
std::future::pending::<()>(),
)
.await
});
tokio::time::timeout(Duration::from_secs(2), async {
loop {
let m = controller.metrics();
if m.queued_by_class.get("analytics").copied().unwrap_or(0) >= 1 {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("analytics must queue");
let overflow = controller
.admit(
AdmitRequest {
tenant: "noisy",
class: WorkloadClass::Analytics,
priority: 50,
deadline: None,
query_id: None,
tag: "overflow",
},
std::future::pending::<()>(),
)
.await;
assert!(
matches!(overflow, Err(AdmitError::Rejected(_))),
"noisy tenant overflow must be rejected: {overflow:?}"
);
let quiet = controller
.admit(
AdmitRequest {
tenant: "quiet",
class: WorkloadClass::Oltp,
priority: 200,
deadline: None,
query_id: None,
tag: "ok",
},
std::future::pending::<()>(),
)
.await
.expect("other tenant must still admit");
drop(quiet);
let mut sched = mongreldb_core::HierarchicalScheduler::new();
sched.set_tenant_quota(
"noisy",
TenantQuota {
max_running: 1,
max_queued: 2,
per_class_running: BTreeMap::new(),
},
);
sched
.submit("noisy", WorkloadClass::Oltp, 1, None, None, "a")
.unwrap();
sched
.submit("noisy", WorkloadClass::Oltp, 1, None, None, "b")
.unwrap();
let err = sched
.submit("noisy", WorkloadClass::Oltp, 1, None, None, "c")
.unwrap_err();
assert!(matches!(err, SchedulerError::TenantQuota { .. }));
sched
.submit("quiet", WorkloadClass::Oltp, 1, None, None, "d")
.unwrap();
}
#[tokio::test]
async fn fragment_memory_counts_against_parent_and_metrics_match() {
let controller = test_controller();
let parent = controller
.admit_parent(
AdmitRequest {
tenant: "t",
class: WorkloadClass::InteractiveSql,
priority: 180,
deadline: None,
query_id: None,
tag: "coord",
},
MemoryClass::QueryExecution,
10_000,
std::future::pending::<()>(),
)
.await
.unwrap();
let child = controller
.reserve_child(&parent, MemoryClass::AiCandidates, 4_000)
.expect("child within budget");
assert_eq!(child.bytes(), 4_000);
assert_eq!(parent.child_used_bytes(), 4_000);
assert_eq!(parent.remaining_bytes(), 6_000);
let metrics = controller.metrics();
assert_eq!(metrics.parent_reserved_bytes, 10_000);
assert_eq!(metrics.child_reserved_bytes, 4_000);
assert_eq!(metrics.open_parents, 1);
assert_eq!(metrics.open_children, 1);
assert_eq!(
metrics
.running_by_class
.get("interactive_sql")
.copied()
.unwrap_or(0),
1
);
let err = match controller.admit_child(&parent, MemoryClass::AiCandidates, 7_000) {
Ok(_) => panic!("expected ChildExceedsParent"),
Err(error) => error,
};
assert!(matches!(
err,
AdmitError::ChildExceedsParent {
requested: 7_000,
parent_remaining: 6_000
}
));
assert_eq!(controller.metrics().child_reserved_bytes, 4_000);
drop(child);
assert_eq!(parent.child_used_bytes(), 0);
assert_eq!(controller.metrics().child_reserved_bytes, 0);
assert_eq!(controller.metrics().open_children, 0);
drop(parent);
let cleared = controller.metrics();
assert_eq!(cleared.parent_reserved_bytes, 0);
assert_eq!(cleared.open_parents, 0);
assert_eq!(
cleared
.running_by_class
.get("interactive_sql")
.copied()
.unwrap_or(0),
0
);
}
#[tokio::test]
async fn replication_reserve_survives_ai_overload() {
let controller = test_controller();
controller.scheduler().set_class_config(
WorkloadClass::AiRetrieval,
ClassConfig {
max_queue: 1,
weight: 32,
reserved_slots: 0,
max_concurrency: 1,
},
);
controller.scheduler().set_class_config(
WorkloadClass::Replication,
ClassConfig {
max_queue: 8,
weight: 256,
reserved_slots: 2,
max_concurrency: 8,
},
);
let _ai = controller
.admit(
AdmitRequest {
tenant: "t",
class: WorkloadClass::AiRetrieval,
priority: 150,
deadline: None,
query_id: None,
tag: "ai-hold",
},
std::future::pending::<()>(),
)
.await
.unwrap();
let replication = controller
.admit(
AdmitRequest {
tenant: "system",
class: WorkloadClass::Replication,
priority: 254,
deadline: None,
query_id: None,
tag: "repl",
},
std::future::pending::<()>(),
)
.await
.expect("replication reserve must survive AI overload");
assert_eq!(
controller
.metrics()
.running_by_class
.get("replication")
.copied()
.unwrap_or(0),
1
);
drop(replication);
}
}