Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions nodedb-mem/src/governor/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,25 @@

use std::sync::atomic::Ordering;

use nodedb_types::DatabaseId;

use super::core::MemoryGovernor;
use crate::engine::EngineId;
use crate::pressure::{PressureLevel, PressureThresholds};

impl MemoryGovernor {
/// Current allocated memory in bytes for a database.
///
/// Returns 0 if the database has no scoped budget or no active allocations.
pub fn database_usage_bytes(&self, db: DatabaseId) -> usize {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker. database_budgets has an entry only after ALTER DATABASE … SET QUOTA. try_reserve credits a database only when it has a budget.

A database with no quota reports 0 bytes forever. The doc line above states this, but the gauge still exports 0 as if it were real.

Track per-database allocation for every database, independent of whether a budget is set.

self.database_budgets
.read()
.unwrap_or_else(|p| p.into_inner())
.get(&db)
.map(|b| b.allocated.load(Ordering::Relaxed))
.unwrap_or(0)
}

/// Total memory allocated across all engines (engine-layer sum). A
/// separate aggregate from [`global_utilization_percent`](Self::global_utilization_percent),
/// which reads the global counter admission enforces the ceiling against.
Expand Down Expand Up @@ -164,4 +178,16 @@ mod tests {
assert_eq!(gov.engine_pressure(EngineId::Vector), PressureLevel::Normal);
assert_eq!(gov.worst_engine_pressure(), PressureLevel::Critical);
}

#[test]
fn database_usage_bytes_tracks_allocation() {
let gov = MemoryGovernor::new(test_config()).unwrap();
gov.set_database_budget(db(), 10_000);
assert_eq!(gov.database_usage_bytes(db()), 0);

let _tok = gov
.try_reserve(db(), tenant(), EngineId::Vector, 512)
.unwrap();
assert_eq!(gov.database_usage_bytes(db()), 512);
}
}
100 changes: 78 additions & 22 deletions nodedb-mem/src/governor/reserve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ use crate::error::{MemError, Result};
use crate::over_release::ReleaseIdentity;
use crate::reservation_token::{ReservationParams, ReservationToken};
use crate::reserve_scope::{ReserveScope, ReservedLayers};
use crate::scoped_budget::ScopedBudget;

/// Build the token both entry points return. They differ in how they reach
/// a committed [`ReservedLayers`], never in what they build from one.
Expand Down Expand Up @@ -62,21 +63,34 @@ impl MemoryGovernor {
scope.try_credit_global()?;

{
let map = self
.database_budgets
.read()
.unwrap_or_else(|p| p.into_inner());
if let Some(budget) = map.get(&db) {
match budget.try_reserve(size) {
Ok(arc) => scope.credit_database(arc),
Err(denied) => {
return Err(MemError::DatabaseBudgetExhausted {
db,
requested: size,
available: budget.available(),
limit: denied.limit,
});
}
let budget = {
let map = self
.database_budgets
.read()
.unwrap_or_else(|p| p.into_inner());
map.get(&db).cloned()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The read lock is now released before try_reserve. The budget is a clone, so its limit is a snapshot. Before this change, a quota change under the write lock excluded any reservation in progress. Now a reservation can pass against a quota that was lowered after the snapshot.

Hold the read lock across try_reserve for the existing-entry path, as before. Take the write lock only to insert a missing entry. Then drop Clone from ScopedBudget, because it breaks the "mutate limit in place" rule its own doc states.

};
let budget = match budget {
Some(b) => b,
None => {
let mut map = self
.database_budgets
.write()
.unwrap_or_else(|p| p.into_inner());
map.entry(db)
.or_insert_with(|| ScopedBudget::new(None))
.clone()
}
};
match budget.try_reserve(size) {
Ok(arc) => scope.credit_database(arc),
Err(denied) => {
return Err(MemError::DatabaseBudgetExhausted {
db,
requested: size,
available: budget.available(),
limit: denied.limit,
});
}
}
}
Expand Down Expand Up @@ -138,13 +152,26 @@ impl MemoryGovernor {
scope.credit_global_unchecked();

{
let map = self
.database_budgets
.read()
.unwrap_or_else(|p| p.into_inner());
if let Some(budget) = map.get(&db) {
scope.credit_database(budget.credit(size));
}
let budget = {
let map = self
.database_budgets
.read()
.unwrap_or_else(|p| p.into_inner());
map.get(&db).cloned()
};
let budget = match budget {
Some(b) => b,
None => {
let mut map = self
.database_budgets
.write()
.unwrap_or_else(|p| p.into_inner());
map.entry(db)
.or_insert_with(|| ScopedBudget::new(None))
.clone()
}
};
scope.credit_database(budget.credit(size));
}

{
Expand Down Expand Up @@ -493,4 +520,33 @@ mod tests {
tok.err()
);
}

#[test]
fn uncapped_database_usage_tracks_allocation() {
let gov = MemoryGovernor::new(test_config()).unwrap();
assert_eq!(gov.database_usage_bytes(db()), 0);

let tok = gov
.try_reserve(db(), tenant(), EngineId::Vector, 512)
.unwrap();
assert_eq!(gov.database_usage_bytes(db()), 512);

drop(tok);
assert_eq!(gov.database_usage_bytes(db()), 0);
}

#[test]
fn unbudgeted_database_tracks_8192_bytes_and_releases_on_drop() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is uncapped_database_usage_tracks_allocation again with a different size and id. Remove one of them.

let gov = MemoryGovernor::new(test_config()).unwrap();
let unbudgeted_db = DatabaseId::new(42);
assert_eq!(gov.database_usage_bytes(unbudgeted_db), 0);

let token = gov
.try_reserve(unbudgeted_db, tenant(), EngineId::Vector, 8192)
.unwrap();
assert_eq!(gov.database_usage_bytes(unbudgeted_db), 8192);

drop(token);
assert_eq!(gov.database_usage_bytes(unbudgeted_db), 0);
}
}
4 changes: 2 additions & 2 deletions nodedb-mem/src/scoped_budget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use std::sync::atomic::{AtomicUsize, Ordering};
///
/// Quota changes mutate `limit` in place so live tokens keep decrementing the
/// same counter.
#[derive(Debug)]
#[derive(Debug, Clone)]
pub(crate) struct ScopedBudget {
/// `None` means uncapped, still counted.
pub(crate) limit: Option<usize>,
Expand All @@ -25,7 +25,7 @@ pub(crate) struct BudgetDenied {
}

impl ScopedBudget {
fn new(limit: Option<usize>) -> Self {
pub(crate) fn new(limit: Option<usize>) -> Self {
Self {
limit,
allocated: Arc::new(AtomicUsize::new(0)),
Expand Down
5 changes: 5 additions & 0 deletions nodedb/src/bootstrap/background_loops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,11 @@ pub fn spawn_background_loops(
info!("mirror lag monitor running");
}

// Database metrics sampler (10-second interval).
// Samples live per-database metrics across subsystems and updates
// `DatabaseMetricsRegistry` gauges for Prometheus scraping.
crate::bootstrap::database_metrics_sampler::spawn_database_metrics_sampler(shared);

// Wire stream delivery managers before Event Plane creation can admit
// CREATE CHANGE STREAM delivery tasks.
shared.webhook_manager.set_state(Arc::clone(shared));
Expand Down
Loading