use std::cmp::Reverse;
use std::collections::{BinaryHeap, HashSet};
use std::sync::Arc;
use std::sync::atomic::{
AtomicBool, AtomicU32, AtomicU64, AtomicUsize,
Ordering::{AcqRel, Acquire, Relaxed, Release},
};
use re_int::SaturatingCast as _;
use re_log_types::TimeInt;
use parking_lot::Mutex;
use tokio::sync::Notify;
use crate::metrics_capture::{
QueryMetrics, SegmentAdmissionCandidateReason, SegmentAdmissionSource,
};
const BUDGET_FRACTION: f64 = 1.0;
pub(crate) const MIN_BUDGET_PER_PARTITION: usize = 4 * 1024 * 1024 * 1024;
pub(crate) const MAX_BUDGET_PER_PARTITION: usize = 1024 * 1024 * 1024 * 1024;
const ENV_BUDGET_MIN: &str = "RERUN_PIPELINE_BUDGET_MIN";
const ENV_BUDGET_MAX: &str = "RERUN_PIPELINE_BUDGET_MAX";
const ENV_BUDGET_FRACTION: &str = "RERUN_PIPELINE_BUDGET_FRACTION";
const INITIAL_ESTIMATE_MULTIPLIER: f64 = 1.5;
const ESTIMATE_EMA_ALPHA: f64 = 0.2;
const MIN_ESTIMATE_MULTIPLIER: f64 = 1.0;
const MAX_ESTIMATE_MULTIPLIER: f64 = 3.0;
pub(crate) const MAX_CONCURRENT_SEGMENTS: usize = 3;
const ADAPTIVE_SEGMENT_ADMISSION_CAP: usize = 16;
const ADAPTIVE_SMALL_SEGMENT_THRESHOLD_BYTES: u64 = 8 * 1024 * 1024;
const ADAPTIVE_DECODE_MULTIPLIER: u64 = MAX_ESTIMATE_MULTIPLIER as u64;
const MAX_EXPERIMENTAL_SEGMENT_ADMISSION_CAP: usize = 1024;
const ENV_SEGMENT_ADMISSION_CAP: &str = "RERUN_SEGMENT_ADMISSION_CAP";
const ENV_ADAPTIVE_SEGMENT_ADMISSION: &str = "RERUN_ADAPTIVE_SEGMENT_ADMISSION";
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct SegmentAdmissionProfile {
segment_count: usize,
complete_segment_sizes: Option<Vec<u64>>,
}
impl SegmentAdmissionProfile {
pub(crate) fn new(segment_count: usize, complete_segment_sizes: Option<Vec<u64>>) -> Self {
Self {
segment_count,
complete_segment_sizes,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct SegmentAdmissionPolicy {
candidate_limit: usize,
effective_limit: usize,
source: SegmentAdmissionSource,
candidate_reason: SegmentAdmissionCandidateReason,
adaptive_enabled: bool,
profile_segment_count: usize,
profile_complete: bool,
p95_segment_bytes: u64,
max_segment_bytes: u64,
largest_window_bytes: u64,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum ExactSegmentAdmissionOverride {
Unset,
Valid(usize),
Invalid,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum AdaptiveSegmentAdmissionOverride {
Default,
Enabled,
Disabled,
Invalid,
}
impl AdaptiveSegmentAdmissionOverride {
fn is_enabled(self) -> bool {
matches!(self, Self::Default | Self::Enabled)
}
}
impl SegmentAdmissionPolicy {
pub(crate) fn from_env(
profile: &SegmentAdmissionProfile,
pipeline_budget_bytes: usize,
) -> Self {
let exact_override = match std::env::var(ENV_SEGMENT_ADMISSION_CAP) {
Ok(raw) => resolve_exact_segment_admission_override(Some(&raw)),
Err(std::env::VarError::NotPresent) => ExactSegmentAdmissionOverride::Unset,
Err(std::env::VarError::NotUnicode(_)) => {
re_log::error!(
"{ENV_SEGMENT_ADMISSION_CAP} is not valid Unicode; \
falling back to {MAX_CONCURRENT_SEGMENTS}",
);
ExactSegmentAdmissionOverride::Invalid
}
};
let adaptive_override = match std::env::var(ENV_ADAPTIVE_SEGMENT_ADMISSION) {
Ok(raw) => resolve_adaptive_segment_admission_override(Some(&raw)),
Err(std::env::VarError::NotPresent) => AdaptiveSegmentAdmissionOverride::Default,
Err(std::env::VarError::NotUnicode(_)) => {
re_log::error!(
"{ENV_ADAPTIVE_SEGMENT_ADMISSION} is not valid Unicode; \
disabling adaptive segment admission and falling back to \
{MAX_CONCURRENT_SEGMENTS}",
);
AdaptiveSegmentAdmissionOverride::Invalid
}
};
Self::resolve_with_override(
profile,
pipeline_budget_bytes,
exact_override,
adaptive_override,
)
}
#[cfg(test)]
fn resolve(
profile: &SegmentAdmissionProfile,
pipeline_budget_bytes: usize,
exact_raw: Option<&str>,
adaptive_raw: Option<&str>,
) -> Self {
Self::resolve_with_override(
profile,
pipeline_budget_bytes,
resolve_exact_segment_admission_override(exact_raw),
resolve_adaptive_segment_admission_override(adaptive_raw),
)
}
fn resolve_with_override(
profile: &SegmentAdmissionProfile,
pipeline_budget_bytes: usize,
exact_override: ExactSegmentAdmissionOverride,
adaptive_override: AdaptiveSegmentAdmissionOverride,
) -> Self {
let profile_complete = profile
.complete_segment_sizes
.as_ref()
.is_some_and(|sizes| {
sizes.len() == profile.segment_count && sizes.iter().all(|size| *size > 0)
});
let mut sorted_sizes = profile
.complete_segment_sizes
.as_deref()
.filter(|_| profile_complete)
.unwrap_or_default()
.to_vec();
sorted_sizes.sort_unstable();
let p95_segment_bytes = nearest_rank_p95(&sorted_sizes).unwrap_or(0);
let max_segment_bytes = sorted_sizes.last().copied().unwrap_or(0);
let largest_window_bytes = sorted_sizes
.iter()
.rev()
.take(ADAPTIVE_SEGMENT_ADMISSION_CAP)
.fold(0u64, |total, size| total.saturating_add(*size));
let decoded_window_bytes = largest_window_bytes.saturating_mul(ADAPTIVE_DECODE_MULTIPLIER);
let candidate_reason = if profile.segment_count <= MAX_CONCURRENT_SEGMENTS {
SegmentAdmissionCandidateReason::InsufficientSegments
} else if !profile_complete {
SegmentAdmissionCandidateReason::IncompleteMetadata
} else if p95_segment_bytes > ADAPTIVE_SMALL_SEGMENT_THRESHOLD_BYTES {
SegmentAdmissionCandidateReason::P95AboveThreshold
} else if decoded_window_bytes > pipeline_budget_bytes.saturating_cast::<u64>() {
SegmentAdmissionCandidateReason::BudgetInsufficient
} else {
SegmentAdmissionCandidateReason::Eligible
};
let candidate_limit = if candidate_reason == SegmentAdmissionCandidateReason::Eligible {
ADAPTIVE_SEGMENT_ADMISSION_CAP
} else {
MAX_CONCURRENT_SEGMENTS
};
let adaptive_enabled = adaptive_override.is_enabled();
let (effective_limit, source) = match exact_override {
ExactSegmentAdmissionOverride::Valid(limit) => {
(limit, SegmentAdmissionSource::ExactOverride)
}
ExactSegmentAdmissionOverride::Invalid => (
MAX_CONCURRENT_SEGMENTS,
SegmentAdmissionSource::InvalidOverride,
),
ExactSegmentAdmissionOverride::Unset if adaptive_enabled => {
(candidate_limit, SegmentAdmissionSource::Adaptive)
}
ExactSegmentAdmissionOverride::Unset => {
(MAX_CONCURRENT_SEGMENTS, SegmentAdmissionSource::MetricsOnly)
}
};
Self {
candidate_limit,
effective_limit,
source,
candidate_reason,
adaptive_enabled,
profile_segment_count: profile.segment_count,
profile_complete,
p95_segment_bytes,
max_segment_bytes,
largest_window_bytes,
}
}
pub(crate) fn effective_limit(self) -> usize {
self.effective_limit
}
pub(crate) fn record_metrics(self, metrics: &QueryMetrics) {
metrics
.segment_admission_candidate_limit
.store(self.candidate_limit.saturating_cast::<u64>(), Relaxed);
metrics
.segment_admission_source
.store(self.source as u64, Relaxed);
metrics
.segment_admission_candidate_reason
.store(self.candidate_reason as u64, Relaxed);
metrics
.segment_admission_adaptive_enabled
.store(u64::from(self.adaptive_enabled), Relaxed);
metrics
.segment_admission_profile_segment_count
.store(self.profile_segment_count.saturating_cast::<u64>(), Relaxed);
metrics
.segment_admission_profile_complete
.store(u64::from(self.profile_complete), Relaxed);
metrics
.segment_admission_p95_segment_bytes
.store(self.p95_segment_bytes, Relaxed);
metrics
.segment_admission_max_segment_bytes
.store(self.max_segment_bytes, Relaxed);
metrics
.segment_admission_largest_window_bytes
.store(self.largest_window_bytes, Relaxed);
metrics
.segment_admission_limit
.store(self.effective_limit.saturating_cast::<u64>(), Relaxed);
}
}
fn nearest_rank_p95(sorted_sizes: &[u64]) -> Option<u64> {
let rank = sorted_sizes.len() - sorted_sizes.len() / 20;
rank.checked_sub(1).map(|index| sorted_sizes[index])
}
const DEFAULT_DIRECT_FETCH_MAX_CONCURRENCY: usize = 128;
const ENV_DIRECT_FETCH_MAX_CONCURRENCY: &str = "RERUN_DIRECT_FETCH_MAX_CONCURRENCY";
const STALL_EMPTY_EMIT_THRESHOLD: u32 = 20;
const STALL_SATURATION_THRESHOLD: f64 = 0.95;
struct PriorityWaiter {
task_time_min: TimeInt,
seq: u64,
notify: Arc<Notify>,
cancelled: Arc<AtomicBool>,
reserved_bytes: usize,
segment_ids: Vec<String>,
}
impl PartialEq for PriorityWaiter {
fn eq(&self, other: &Self) -> bool {
self.task_time_min == other.task_time_min && self.seq == other.seq
}
}
impl Eq for PriorityWaiter {}
impl PartialOrd for PriorityWaiter {
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl Ord for PriorityWaiter {
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
self.task_time_min
.cmp(&other.task_time_min)
.then_with(|| self.seq.cmp(&other.seq))
}
}
#[derive(Default)]
struct SegmentGate {
all: HashSet<String>,
bypass: HashSet<String>,
}
impl SegmentGate {
fn effective_len(&self) -> usize {
self.all.len().saturating_sub(self.bypass.len())
}
}
#[derive(Clone, Copy)]
struct AdmissionBlockers {
bytes: bool,
segments: bool,
}
pub(crate) struct PipelineBudget {
budget: usize,
segment_limit: usize,
current: AtomicUsize,
wait_queue: Mutex<BinaryHeap<Reverse<PriorityWaiter>>>,
wait_seq: AtomicU64,
active_segments: Mutex<SegmentGate>,
empty_emit_count: AtomicU32,
force_overcommit: AtomicBool,
estimate_multiplier: AtomicU64,
peak_current: AtomicUsize,
metrics: Option<Arc<QueryMetrics>>,
total_released_bytes: AtomicUsize,
total_releases: AtomicU64,
#[cfg(test)]
test_pause_hook: parking_lot::Mutex<Option<TestPauseHook>>,
}
#[cfg(test)]
#[derive(Clone)]
struct TestPauseHook {
arrived: Arc<Notify>,
resume: Arc<Notify>,
}
fn read_env_trimmed(key: &str) -> Option<String> {
match std::env::var(key) {
Ok(v) => {
let trimmed = v.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_owned())
}
}
Err(std::env::VarError::NotPresent) => None,
Err(std::env::VarError::NotUnicode(_)) => {
re_log::error!("{key}: value is not valid Unicode; using default");
None
}
}
}
fn parse_bytes_or_default(key: &str, raw: &str, default_bytes: usize) -> usize {
let parsed = re_format::parse_bytes(raw).or_else(|| raw.parse::<i64>().ok());
match parsed {
Some(n) if n > 0 => n as usize,
Some(_) => {
re_log::error!(
"{key}={raw:?} must be > 0; falling back to default {}",
re_format::format_bytes(default_bytes as f64),
);
default_bytes
}
None => {
re_log::error!(
"{key}={raw:?} could not be parsed as a byte size (e.g. \"64MB\", \"1GiB\", \
or a bare integer number of bytes); falling back to default {}",
re_format::format_bytes(default_bytes as f64),
);
default_bytes
}
}
}
fn parse_fraction_or_default(key: &str, raw: &str, default: f64) -> f64 {
match raw.parse::<f64>() {
Ok(f) if f.is_finite() && f > 0.0 && f <= 1.0 => f,
Ok(f) => {
re_log::error!(
"{key}={raw:?} must be a finite value in (0.0, 1.0], got {f}; \
falling back to default {default}",
);
default
}
Err(err) => {
re_log::error!(
"{key}={raw:?} could not be parsed as a float ({err}); \
falling back to default {default}",
);
default
}
}
}
fn parse_usize_or_default(key: &str, raw: &str, default: usize) -> usize {
match raw.parse::<i64>() {
Ok(n) if n > 0 => n as usize,
Ok(_) => {
re_log::error!("{key}={raw:?} must be > 0; falling back to default {default}");
default
}
Err(err) => {
re_log::error!(
"{key}={raw:?} could not be parsed as a positive integer ({err}); \
falling back to default {default}",
);
default
}
}
}
fn resolve_exact_segment_admission_override(raw: Option<&str>) -> ExactSegmentAdmissionOverride {
let Some(raw) = raw.map(str::trim).filter(|raw| !raw.is_empty()) else {
return ExactSegmentAdmissionOverride::Unset;
};
match raw.parse::<usize>() {
Ok(cap)
if (MAX_CONCURRENT_SEGMENTS..=MAX_EXPERIMENTAL_SEGMENT_ADMISSION_CAP)
.contains(&cap) =>
{
ExactSegmentAdmissionOverride::Valid(cap)
}
Ok(cap) => {
re_log::error!(
"{ENV_SEGMENT_ADMISSION_CAP}={cap} must be in \
[{MAX_CONCURRENT_SEGMENTS}, {MAX_EXPERIMENTAL_SEGMENT_ADMISSION_CAP}]; \
falling back to {MAX_CONCURRENT_SEGMENTS}",
);
ExactSegmentAdmissionOverride::Invalid
}
Err(err) => {
re_log::error!(
"{ENV_SEGMENT_ADMISSION_CAP}={raw:?} could not be parsed as an integer ({err}); \
falling back to {MAX_CONCURRENT_SEGMENTS}",
);
ExactSegmentAdmissionOverride::Invalid
}
}
}
#[cfg(test)]
fn resolve_segment_admission_limit(raw: Option<&str>) -> usize {
match resolve_exact_segment_admission_override(raw) {
ExactSegmentAdmissionOverride::Valid(limit) => limit,
ExactSegmentAdmissionOverride::Unset | ExactSegmentAdmissionOverride::Invalid => {
MAX_CONCURRENT_SEGMENTS
}
}
}
fn resolve_adaptive_segment_admission_override(
raw: Option<&str>,
) -> AdaptiveSegmentAdmissionOverride {
let Some(raw) = raw.map(str::trim).filter(|raw| !raw.is_empty()) else {
return AdaptiveSegmentAdmissionOverride::Default;
};
match raw.parse::<bool>() {
Ok(true) => AdaptiveSegmentAdmissionOverride::Enabled,
Ok(false) => AdaptiveSegmentAdmissionOverride::Disabled,
Err(err) => {
re_log::error!(
"{ENV_ADAPTIVE_SEGMENT_ADMISSION}={raw:?} could not be parsed as a boolean ({err}); \
disabling adaptive segment admission and falling back to {MAX_CONCURRENT_SEGMENTS}",
);
AdaptiveSegmentAdmissionOverride::Invalid
}
}
}
fn read_env_bytes(key: &str, default_bytes: usize) -> usize {
match read_env_trimmed(key) {
Some(raw) => parse_bytes_or_default(key, &raw, default_bytes),
None => default_bytes,
}
}
fn read_env_fraction(key: &str, default: f64) -> f64 {
match read_env_trimmed(key) {
Some(raw) => parse_fraction_or_default(key, &raw, default),
None => default,
}
}
fn read_env_usize(key: &str, default: usize) -> usize {
match read_env_trimmed(key) {
Some(raw) => parse_usize_or_default(key, &raw, default),
None => default,
}
}
pub(crate) fn direct_fetch_semaphore() -> &'static tokio::sync::Semaphore {
static SEM: std::sync::OnceLock<tokio::sync::Semaphore> = std::sync::OnceLock::new();
SEM.get_or_init(|| {
let permits = read_env_usize(
ENV_DIRECT_FETCH_MAX_CONCURRENCY,
DEFAULT_DIRECT_FETCH_MAX_CONCURRENCY,
)
.max(1);
re_log::debug!("direct-fetch concurrency cap = {permits}");
tokio::sync::Semaphore::new(permits)
})
}
const DEFAULT_QUERY_DATASET_MAX_CONCURRENCY: usize = 16;
const ENV_QUERY_DATASET_MAX_CONCURRENCY: &str = "RERUN_QUERY_DATASET_MAX_CONCURRENCY";
pub(crate) fn query_dataset_max_concurrency() -> usize {
static N: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
*N.get_or_init(|| {
read_env_usize(
ENV_QUERY_DATASET_MAX_CONCURRENCY,
DEFAULT_QUERY_DATASET_MAX_CONCURRENCY,
)
.max(1)
})
}
pub(crate) fn query_dataset_semaphore() -> &'static tokio::sync::Semaphore {
static SEM: std::sync::OnceLock<tokio::sync::Semaphore> = std::sync::OnceLock::new();
SEM.get_or_init(|| {
let permits = query_dataset_max_concurrency();
re_log::debug!("query_dataset concurrency cap = {permits}");
tokio::sync::Semaphore::new(permits)
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct ResolvedPipelineBudget {
bytes: usize,
}
impl ResolvedPipelineBudget {
pub(crate) fn bytes(self) -> usize {
self.bytes
}
}
impl PipelineBudget {
pub(crate) fn resolve_size(
total_uncompressed_estimate: usize,
num_partitions: usize,
) -> ResolvedPipelineBudget {
let fraction = read_env_fraction(ENV_BUDGET_FRACTION, BUDGET_FRACTION);
let mut min_per_partition = read_env_bytes(ENV_BUDGET_MIN, MIN_BUDGET_PER_PARTITION);
let mut max_per_partition = read_env_bytes(ENV_BUDGET_MAX, MAX_BUDGET_PER_PARTITION);
if min_per_partition > max_per_partition {
re_log::error!(
"{ENV_BUDGET_MIN} ({}) must not exceed {ENV_BUDGET_MAX} ({}); \
falling back to defaults for both.",
re_format::format_bytes(min_per_partition as f64),
re_format::format_bytes(max_per_partition as f64),
);
min_per_partition = MIN_BUDGET_PER_PARTITION;
max_per_partition = MAX_BUDGET_PER_PARTITION;
}
#[expect(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
let per_partition = ((total_uncompressed_estimate as f64 * fraction)
/ num_partitions.max(1) as f64) as usize;
let budget = per_partition.clamp(min_per_partition, max_per_partition) * num_partitions;
re_log::debug!("Pipeline budget: {}MB", budget / (1024 * 1024));
ResolvedPipelineBudget { bytes: budget }
}
#[cfg(test)]
pub(crate) fn new(total_uncompressed_estimate: usize, num_partitions: usize) -> Self {
let resolved = Self::resolve_size(total_uncompressed_estimate, num_partitions);
Self::with_exact_budget_and_metrics(resolved.bytes, MAX_CONCURRENT_SEGMENTS, None)
}
pub(crate) fn new_with_metrics(
resolved: ResolvedPipelineBudget,
policy: SegmentAdmissionPolicy,
metrics: Arc<QueryMetrics>,
) -> Self {
policy.record_metrics(&metrics);
Self::with_exact_budget_and_metrics(resolved.bytes, policy.effective_limit(), Some(metrics))
}
#[cfg(test)]
fn with_exact_budget(budget: usize) -> Self {
Self::with_exact_budget_and_metrics(budget, MAX_CONCURRENT_SEGMENTS, None)
}
fn with_exact_budget_and_metrics(
budget: usize,
segment_limit: usize,
metrics: Option<Arc<QueryMetrics>>,
) -> Self {
assert!(
segment_limit > 0,
"segment admission limit must be positive"
);
if let Some(metrics) = &metrics {
metrics
.pipeline_budget_bytes
.fetch_max(budget.saturating_cast::<u64>(), Relaxed);
}
Self {
budget,
segment_limit,
current: AtomicUsize::new(0),
wait_queue: Mutex::new(BinaryHeap::new()),
wait_seq: AtomicU64::new(0),
active_segments: Mutex::new(SegmentGate::default()),
empty_emit_count: AtomicU32::new(0),
force_overcommit: AtomicBool::new(false),
estimate_multiplier: AtomicU64::new(INITIAL_ESTIMATE_MULTIPLIER.to_bits()),
peak_current: AtomicUsize::new(0),
metrics,
total_released_bytes: AtomicUsize::new(0),
total_releases: AtomicU64::new(0),
#[cfg(test)]
test_pause_hook: parking_lot::Mutex::new(None),
}
}
pub(crate) fn segment_limit(&self) -> usize {
self.segment_limit
}
#[cfg(test)]
pub(crate) fn total_releases(&self) -> u64 {
self.total_releases.load(Acquire)
}
fn current_multiplier(&self) -> f64 {
f64::from_bits(self.estimate_multiplier.load(Acquire))
}
#[cfg(test)]
fn set_multiplier(&self, multiplier: f64) {
self.estimate_multiplier
.store(multiplier.to_bits(), Release);
}
#[cfg(test)]
fn arm_pause_hook(&self) -> TestPauseHook {
let hook = TestPauseHook {
arrived: Arc::new(Notify::new()),
resume: Arc::new(Notify::new()),
};
*self.test_pause_hook.lock() = Some(hook.clone());
hook
}
fn record_actual_sample(&self, estimated: usize, actual: usize) {
if estimated == 0 {
return;
}
let observed = ((actual as f64) / (estimated as f64))
.clamp(MIN_ESTIMATE_MULTIPLIER, MAX_ESTIMATE_MULTIPLIER);
self.estimate_multiplier
.fetch_update(AcqRel, Acquire, |bits| {
let curr = f64::from_bits(bits);
let next = ESTIMATE_EMA_ALPHA * observed + (1.0 - ESTIMATE_EMA_ALPHA) * curr;
let next = next.clamp(MIN_ESTIMATE_MULTIPLIER, MAX_ESTIMATE_MULTIPLIER);
Some(next.to_bits())
})
.expect("closure always returns Some");
}
fn wake_next(&self) {
let mut queue = self.wait_queue.lock();
let segments = self.active_segments.lock();
let force = self.force_overcommit.load(Acquire);
let current = self.current.load(Acquire);
let mut held: Vec<Reverse<PriorityWaiter>> = Vec::new();
let mut woke = false;
while let Some(Reverse(waiter)) = queue.pop() {
if waiter.cancelled.load(Acquire) {
continue;
}
if !woke && self.waiter_admittable(&waiter, force, current, &segments) {
waiter.notify.notify_one();
woke = true;
continue;
}
held.push(Reverse(waiter));
}
for h in held {
queue.push(h);
}
}
fn waiter_admittable(
&self,
waiter: &PriorityWaiter,
force: bool,
current: usize,
segments: &SegmentGate,
) -> bool {
if force {
return true;
}
let new_segments = waiter
.segment_ids
.iter()
.filter(|s| !segments.all.contains(s.as_str()))
.count();
if segments.effective_len() + new_segments > self.segment_limit {
return false;
}
current + waiter.reserved_bytes <= self.budget
}
fn record_peak_active_segments(&self, segments: &SegmentGate) {
if let Some(metrics) = &self.metrics {
metrics
.peak_active_segments
.fetch_max(segments.all.len() as u64, Relaxed);
}
}
fn record_peak_decoded_bytes(&self, current: usize) {
self.peak_current.fetch_max(current, AcqRel);
if let Some(metrics) = &self.metrics {
metrics
.pipeline_peak_decoded_bytes
.fetch_max(current.saturating_cast::<u64>(), Relaxed);
}
}
fn record_initial_wait(&self, blockers: AdmissionBlockers) {
let Some(metrics) = &self.metrics else {
return;
};
if blockers.bytes {
metrics.pipeline_byte_waits.fetch_add(1, Relaxed);
}
if blockers.segments {
metrics.segment_admission_waits.fetch_add(1, Relaxed);
}
}
fn try_admit(
&self,
reserved_bytes: usize,
segment_ids: &[String],
) -> Result<usize, AdmissionBlockers> {
let mut segments = self.active_segments.lock();
if self.force_overcommit.load(Acquire) {
let new_cur = self.current.fetch_add(reserved_bytes, AcqRel) + reserved_bytes;
self.record_peak_decoded_bytes(new_cur);
for s in segment_ids {
if segments.all.insert(s.clone()) {
segments.bypass.insert(s.clone());
}
}
self.record_peak_active_segments(&segments);
return Ok(new_cur);
}
let new_segments = segment_ids
.iter()
.filter(|s| !segments.all.contains(s.as_str()))
.count();
let segments_blocked = segments.effective_len() + new_segments > self.segment_limit;
if segments_blocked {
let bytes_blocked =
self.current.load(Acquire).saturating_add(reserved_bytes) > self.budget;
return Err(AdmissionBlockers {
bytes: bytes_blocked,
segments: true,
});
}
let new_cur = self.try_acquire(reserved_bytes).ok_or(AdmissionBlockers {
bytes: true,
segments: false,
})?;
segments.all.extend(segment_ids.iter().cloned());
self.record_peak_active_segments(&segments);
Ok(new_cur)
}
fn try_acquire(&self, reserved_bytes: usize) -> Option<usize> {
let previous = self
.current
.try_update(AcqRel, Acquire, |cur| {
let next = cur + reserved_bytes;
(next <= self.budget).then_some(next)
})
.ok()?;
let current = previous + reserved_bytes;
self.record_peak_decoded_bytes(current);
Some(current)
}
#[cfg(test)]
pub(crate) async fn reserve(&self, estimated_bytes: usize) -> usize {
self.reserve_with_priority(estimated_bytes, TimeInt::MAX, &[])
.await
}
#[expect(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
pub(crate) async fn reserve_with_priority(
&self,
estimated_bytes: usize,
task_time_min: TimeInt,
segment_ids: &[String],
) -> usize {
let reserved_bytes = ((estimated_bytes as f64) * self.current_multiplier()) as usize;
if reserved_bytes > self.budget {
re_log::warn!(
"Single fetch reservation ({}MB, raw estimate {}MB) exceeds entire \
pipeline budget ({}MB across all partitions) — allowing it through \
to avoid deadlock.",
reserved_bytes / (1024 * 1024),
estimated_bytes / (1024 * 1024),
self.budget / (1024 * 1024),
);
let new_cur = self.current.fetch_add(reserved_bytes, AcqRel) + reserved_bytes;
self.record_peak_decoded_bytes(new_cur);
let mut segments = self.active_segments.lock();
for s in segment_ids {
if segments.all.insert(s.clone()) {
segments.bypass.insert(s.clone());
}
}
self.record_peak_active_segments(&segments);
return reserved_bytes;
}
re_log::debug_assert!(
segment_ids
.iter()
.enumerate()
.all(|(idx, segment_id)| !segment_ids[..idx].contains(segment_id)),
"segment_ids must be distinct"
);
let distinct_segments = segment_ids.len();
if distinct_segments > self.segment_limit {
re_log::warn_once!(
"Single fetch reservation spans {distinct_segments} distinct segments, \
exceeding the concurrent-segment cap ({}) — allowing it through to avoid \
deadlock.",
self.segment_limit,
);
let new_cur = self.current.fetch_add(reserved_bytes, AcqRel) + reserved_bytes;
self.record_peak_decoded_bytes(new_cur);
let mut segments = self.active_segments.lock();
for s in segment_ids {
if segments.all.insert(s.clone()) {
segments.bypass.insert(s.clone());
}
}
self.record_peak_active_segments(&segments);
return reserved_bytes;
}
let mut wait_count: u32 = 0;
loop {
if let Ok(new_cur) = self.try_admit(reserved_bytes, segment_ids) {
if new_cur < self.budget {
self.wake_next();
}
if wait_count > 0 {
re_log::debug!(
"Budget reserve succeeded after {wait_count} waits: \
reserved {}MB, current {}MB / {}MB",
reserved_bytes / (1024 * 1024),
new_cur / (1024 * 1024),
self.budget / (1024 * 1024),
);
}
return reserved_bytes;
}
let notify = Arc::new(Notify::new());
let cancelled = Arc::new(AtomicBool::new(false));
let seq = self.wait_seq.fetch_add(1, AcqRel);
self.wait_queue.lock().push(Reverse(PriorityWaiter {
task_time_min,
seq,
notify: Arc::clone(¬ify),
cancelled: Arc::clone(&cancelled),
reserved_bytes,
segment_ids: segment_ids.to_vec(),
}));
#[cfg(test)]
{
let hook = self.test_pause_hook.lock().clone();
if let Some(hook) = hook {
hook.arrived.notify_one();
hook.resume.notified().await;
}
}
let blockers = match self.try_admit(reserved_bytes, segment_ids) {
Ok(new_cur) => {
cancelled.store(true, Release);
if new_cur < self.budget {
self.wake_next();
}
return reserved_bytes;
}
Err(blockers) => blockers,
};
if wait_count == 0 {
self.record_initial_wait(blockers);
}
wait_count += 1;
if wait_count == 1 || wait_count.is_multiple_of(10) {
let segments = self.active_segments.lock();
re_log::info!(
"Budget backpressure (wait #{wait_count}): want {}MB, \
current {}MB / {}MB budget, active_segments={} (bypass={})",
reserved_bytes / (1024 * 1024),
self.current.load(Acquire) / (1024 * 1024),
self.budget / (1024 * 1024),
segments.all.len(),
segments.bypass.len(),
);
}
notify.notified().await;
}
}
pub(crate) fn adjust_reservation(&self, estimated: usize, reserved: usize, actual: usize) {
if actual > reserved {
let new_cur = self.current.fetch_add(actual - reserved, AcqRel) + (actual - reserved);
self.record_peak_decoded_bytes(new_cur);
} else if reserved > actual {
self.current
.fetch_update(AcqRel, Acquire, |current| {
Some(current.saturating_sub(reserved - actual))
})
.expect("closure always returns Some");
self.empty_emit_count.store(0, Release);
self.force_overcommit.store(false, Release);
self.wake_next();
}
self.record_actual_sample(estimated, actual);
}
pub(crate) fn release(&self, bytes: usize) {
let prev = self
.current
.fetch_update(AcqRel, Acquire, |current| {
Some(current.saturating_sub(bytes))
})
.expect("closure always returns Some");
self.total_released_bytes.fetch_add(bytes, AcqRel);
self.total_releases.fetch_add(1, AcqRel);
self.empty_emit_count.store(0, Release);
self.force_overcommit.store(false, Release);
re_log::debug!(
"Budget release: freed {}MB, {}MB → {}MB / {}MB",
bytes / (1024 * 1024),
prev / (1024 * 1024),
prev.saturating_sub(bytes) / (1024 * 1024),
self.budget / (1024 * 1024),
);
self.wake_next();
}
pub(crate) fn publish_segment_finalized(&self, segment_id: &str) {
let mut segments = self.active_segments.lock();
let was_bypass = segments.bypass.remove(segment_id);
let removed = segments.all.remove(segment_id);
drop(segments);
if removed || was_bypass {
self.wake_next();
}
}
fn arm_stall_breaker(&self, reason: std::fmt::Arguments<'_>, saturation: f64) {
if self
.force_overcommit
.compare_exchange(false, true, AcqRel, Relaxed)
.is_err()
{
return;
}
if let Some(metrics) = &self.metrics {
metrics
.pipeline_stall_breaker_activations
.fetch_add(1, Relaxed);
}
re_log::info!(
"PipelineBudget stall detected: {reason} with budget {:.0}% saturated — enabling \
force_overcommit until next progress",
saturation * 100.0,
);
self.wake_next();
}
pub(crate) fn notify_empty_emit(&self) {
#[expect(clippy::cast_precision_loss)]
let saturation = (self.current.load(Acquire) as f64) / (self.budget.max(1) as f64);
if saturation < STALL_SATURATION_THRESHOLD {
self.empty_emit_count.store(0, Release);
return;
}
let count = self.empty_emit_count.fetch_add(1, AcqRel) + 1;
if count >= STALL_EMPTY_EMIT_THRESHOLD {
self.arm_stall_breaker(format_args!("{count} consecutive empty emits"), saturation);
}
}
pub(crate) fn notify_row_emitted(&self) {
self.empty_emit_count.store(0, Release);
self.force_overcommit.store(false, Release);
}
fn refund_reservation(&self, reserved: usize) {
if reserved == 0 {
return;
}
self.current
.fetch_update(AcqRel, Acquire, |current| {
Some(current.saturating_sub(reserved))
})
.expect("closure always returns Some");
self.wake_next();
}
#[cfg(test)]
pub(crate) async fn reserve_guarded(&self, estimated: usize) -> ReservationGuard<'_> {
self.reserve_guarded_with_priority(estimated, TimeInt::MAX, Vec::new())
.await
}
pub(crate) async fn reserve_guarded_with_priority(
&self,
estimated: usize,
task_time_min: TimeInt,
segment_ids: Vec<String>,
) -> ReservationGuard<'_> {
let reserved = self
.reserve_with_priority(estimated, task_time_min, &segment_ids)
.await;
ReservationGuard {
budget: self,
estimated,
reserved,
segment_ids,
committed: false,
}
}
}
#[must_use = "ReservationGuard returns its bytes and segment slots to the budget on drop; \
call .commit(actual) once the decoded size is known"]
pub(crate) struct ReservationGuard<'a> {
budget: &'a PipelineBudget,
estimated: usize,
reserved: usize,
segment_ids: Vec<String>,
committed: bool,
}
impl ReservationGuard<'_> {
pub(crate) fn commit(mut self, actual: usize) {
self.budget
.adjust_reservation(self.estimated, self.reserved, actual);
self.committed = true;
}
}
impl Drop for ReservationGuard<'_> {
fn drop(&mut self) {
if !self.committed {
self.budget.refund_reservation(self.reserved);
for segment_id in &self.segment_ids {
self.budget.publish_segment_finalized(segment_id);
}
}
}
}
impl Drop for PipelineBudget {
fn drop(&mut self) {
let n_releases = self.total_releases.load(Acquire);
if n_releases == 0 {
return;
}
const MB: usize = 1024 * 1024;
let peak = self.peak_current.load(Acquire);
let total_released = self.total_released_bytes.load(Acquire);
let pct = if self.budget > 0 {
#[expect(clippy::cast_precision_loss)]
let pct = peak as f64 / self.budget as f64 * 100.0;
pct
} else {
0.0
};
re_log::info!(
"PipelineBudget summary: peak={}MB / {}MB ({pct:.0}%), \
released_total={}MB across {n_releases} calls",
peak / MB,
self.budget / MB,
total_released / MB,
);
}
}
impl std::fmt::Debug for PipelineBudget {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PipelineBudget")
.field("budget", &self.budget)
.field("current", &self.current.load(Relaxed))
.finish_non_exhaustive()
}
}
#[cfg(test)]
mod tests;