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_log_types::TimeInt;
use parking_lot::Mutex;
use tokio::sync::Notify;
use crate::metrics_capture::QueryMetrics;
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 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,
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 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)
})
}
impl PipelineBudget {
#[cfg(test)]
pub(crate) fn new(total_uncompressed_estimate: usize, num_partitions: usize) -> Self {
Self::new_impl(total_uncompressed_estimate, num_partitions, None)
}
pub(crate) fn new_with_metrics(
total_uncompressed_estimate: usize,
num_partitions: usize,
metrics: Arc<QueryMetrics>,
) -> Self {
Self::new_impl(total_uncompressed_estimate, num_partitions, Some(metrics))
}
fn new_impl(
total_uncompressed_estimate: usize,
num_partitions: usize,
metrics: Option<Arc<QueryMetrics>>,
) -> Self {
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));
Self::with_exact_budget_and_metrics(budget, metrics)
}
#[cfg(test)]
fn with_exact_budget(budget: usize) -> Self {
Self::with_exact_budget_and_metrics(budget, None)
}
fn with_exact_budget_and_metrics(budget: usize, metrics: Option<Arc<QueryMetrics>>) -> Self {
if let Some(metrics) = &metrics {
metrics
.pipeline_budget_bytes
.fetch_max(u64::try_from(budget).unwrap_or(u64::MAX), Relaxed);
}
Self {
budget,
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),
}
}
#[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 > MAX_CONCURRENT_SEGMENTS {
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(u64::try_from(current).unwrap_or(u64::MAX), 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 > MAX_CONCURRENT_SEGMENTS;
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 > MAX_CONCURRENT_SEGMENTS {
re_log::warn_once!(
"Single fetch reservation spans {distinct_segments} distinct segments, \
exceeding the concurrent-segment cap ({MAX_CONCURRENT_SEGMENTS}) — allowing \
it through to avoid deadlock.",
);
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();
}
}
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
.force_overcommit
.compare_exchange(false, true, AcqRel, Relaxed)
.is_ok()
{
if let Some(metrics) = &self.metrics {
metrics
.pipeline_stall_breaker_activations
.fetch_add(1, Relaxed);
}
re_log::info!(
"PipelineBudget stall detected: {count} consecutive empty emits with \
budget {:.0}% saturated — enabling force_overcommit until next progress",
saturation * 100.0,
);
self.wake_next();
}
}
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;