use std::sync::Arc;
use crate::analytics::{QueryInfo, QueryType};
use crate::metrics_capture::QueryMetrics;
use super::*;
fn test_query_metrics() -> Arc<QueryMetrics> {
Arc::new(QueryMetrics::new(QueryInfo {
dataset_id: "test".to_owned(),
query_chunks: 0,
query_segments: 0,
query_layers: 0,
query_columns: 0,
query_entities: 0,
query_bytes: 0,
query_chunks_per_segment_min: 0,
query_chunks_per_segment_max: 0,
query_chunks_per_segment_mean: 0.0,
query_type: QueryType::FullScan,
primary_index_name: None,
time_to_first_chunk_info: None,
trace_id: None,
filters_pushed_down: 0,
filters_applied_client_side: 0,
entity_path_narrowing_applied: false,
filters_total: 0,
filters_signatures: String::new(),
filters_signatures_exact: String::new(),
filters_signatures_inexact: String::new(),
filters_signatures_unsupported: String::new(),
}))
}
#[test]
fn test_budget_clamps_to_min() {
let budget = PipelineBudget::new(10 * 1024 * 1024, 1);
assert_eq!(budget.budget, MIN_BUDGET_PER_PARTITION);
}
#[test]
fn test_budget_clamps_to_max() {
let budget = PipelineBudget::new(8 * 1024 * 1024 * 1024 * 1024 * 1024, 1);
assert_eq!(budget.budget, MAX_BUDGET_PER_PARTITION);
}
#[test]
fn test_budget_scales_with_partitions() {
let budget = PipelineBudget::new(64 * 1024 * 1024 * 1024 * 1024 * 1024, 14);
assert_eq!(budget.budget, MAX_BUDGET_PER_PARTITION * 14);
}
#[test]
fn test_budget_small_data_many_partitions() {
let budget = PipelineBudget::new(100 * 1024 * 1024, 4);
assert_eq!(budget.budget, MIN_BUDGET_PER_PARTITION * 4);
}
#[tokio::test]
async fn test_reserve_blocks_when_budget_exhausted() {
let budget = Arc::new(PipelineBudget::new(0, 1)); budget.set_multiplier(1.0);
let half = budget.budget / 2;
budget.reserve(half).await;
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
half
);
budget.reserve(half).await;
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
half * 2
);
let budget_clone = Arc::clone(&budget);
let handle = tokio::spawn(async move {
budget_clone.reserve(half).await;
});
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
assert!(
!handle.is_finished(),
"reserve should block when budget is exhausted"
);
budget.release(half);
tokio::time::timeout(std::time::Duration::from_secs(1), handle)
.await
.expect("reserve should unblock after release")
.expect("task should not panic");
}
#[tokio::test]
async fn test_adjust_reservation_corrects_estimate() {
let budget = Arc::new(PipelineBudget::new(0, 1)); budget.set_multiplier(1.0);
let estimated = 1000;
let actual = 600;
let reserved = budget.reserve(estimated).await;
assert_eq!(reserved, estimated);
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
estimated
);
budget.adjust_reservation(estimated, reserved, actual);
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
actual
);
}
#[tokio::test]
async fn test_estimate_multiplier_adapts_over_time() {
let budget = Arc::new(PipelineBudget::new(0, 1));
assert_eq!(budget.current_multiplier(), INITIAL_ESTIMATE_MULTIPLIER);
let estimated = 1000;
let actual = 2000;
let reserved1 = budget.reserve(estimated).await;
assert_eq!(reserved1, 1500);
budget.adjust_reservation(estimated, reserved1, actual);
budget.release(actual);
assert!((budget.current_multiplier() - 1.6).abs() < 1e-9);
for _ in 0..40 {
let reserved = budget.reserve(estimated).await;
budget.adjust_reservation(estimated, reserved, actual);
budget.release(actual);
}
assert!(
(budget.current_multiplier() - 2.0).abs() < 0.01,
"multiplier should converge toward 2.0, got {}",
budget.current_multiplier()
);
}
#[tokio::test]
async fn test_estimate_multiplier_is_clamped() {
let budget = Arc::new(PipelineBudget::new(0, 1));
for _ in 0..100 {
budget.record_actual_sample(100, 1000); }
assert!(
(budget.current_multiplier() - MAX_ESTIMATE_MULTIPLIER).abs() < 1e-4,
"multiplier should converge to MAX, got {}",
budget.current_multiplier()
);
for _ in 0..100 {
budget.record_actual_sample(1000, 100); }
assert!(
(budget.current_multiplier() - MIN_ESTIMATE_MULTIPLIER).abs() < 1e-4,
"multiplier should converge to MIN, got {}",
budget.current_multiplier()
);
}
#[tokio::test]
async fn test_peak_current_tracks_high_water_mark() {
let metrics = test_query_metrics();
let budget = Arc::new(PipelineBudget::with_exact_budget_and_metrics(
100 * 1024 * 1024,
Some(Arc::clone(&metrics)),
));
budget.set_multiplier(1.0);
assert_eq!(
metrics.pipeline_budget_bytes.load(Acquire),
100 * 1024 * 1024
);
let r1 = budget.reserve(10 * 1024 * 1024).await;
let r2 = budget.reserve(5 * 1024 * 1024).await;
let peak_before = budget
.peak_current
.load(std::sync::atomic::Ordering::Acquire);
assert_eq!(peak_before, r1 + r2);
assert_eq!(
metrics.pipeline_peak_decoded_bytes.load(Acquire),
u64::try_from(peak_before).unwrap()
);
budget.release(r1 + r2);
let r3 = budget.reserve(1024 * 1024).await;
let peak_after = budget
.peak_current
.load(std::sync::atomic::Ordering::Acquire);
assert_eq!(
peak_after, peak_before,
"peak should not regress after releases",
);
assert_eq!(
metrics.pipeline_peak_decoded_bytes.load(Acquire),
u64::try_from(peak_before).unwrap(),
"published peak should not regress after releases",
);
budget.release(r3);
}
#[tokio::test]
async fn test_reserve_no_thundering_herd_under_contention() {
let budget = Arc::new(PipelineBudget::new(0, 1));
budget.set_multiplier(1.0);
let full = budget.budget;
let n: usize = 32;
let per_task = full / n;
assert!(per_task > 0, "budget too small for {n}-way contention test");
let mut handles = Vec::with_capacity(n);
for _ in 0..n {
let budget = Arc::clone(&budget);
handles.push(tokio::spawn(async move { budget.reserve(per_task).await }));
}
for handle in handles {
tokio::time::timeout(std::time::Duration::from_secs(5), handle)
.await
.expect("reserve hung under contention")
.expect("task panicked");
}
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
per_task * n,
);
}
#[tokio::test]
async fn test_reserve_no_lost_wakeup_in_wait_path() {
let budget = Arc::new(PipelineBudget::new(0, 1));
budget.set_multiplier(1.0);
let full = budget.budget;
budget.reserve(full).await;
let pause = budget.arm_pause_hook();
let reserver = {
let budget = Arc::clone(&budget);
tokio::spawn(async move { budget.reserve(1).await })
};
pause.arrived.notified().await;
*budget.test_pause_hook.lock() = None;
budget.release(full);
pause.resume.notify_one();
tokio::time::timeout(std::time::Duration::from_secs(1), reserver)
.await
.expect("reserve hung — lost-wakeup race in rollback→enqueue gap")
.expect("reserver task panicked");
}
#[tokio::test]
async fn test_reservation_guard_commit_records_actual() {
let budget = Arc::new(PipelineBudget::new(0, 1));
budget.set_multiplier(1.0);
let estimated = 1000;
let actual = 800;
let guard = budget.reserve_guarded(estimated).await;
guard.commit(actual);
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
actual,
"commit should reduce current to the actual decoded size",
);
}
#[tokio::test]
async fn test_reservation_guard_drop_refunds_reservation() {
let budget = Arc::new(PipelineBudget::new(0, 1));
budget.set_multiplier(1.0);
let estimated = 1000;
let multiplier_before = f64::from_bits(
budget
.estimate_multiplier
.load(std::sync::atomic::Ordering::Acquire),
);
{
let _guard = budget.reserve_guarded(estimated).await;
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
estimated,
);
}
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
0,
"dropped guard should refund the entire reservation",
);
let multiplier_after = f64::from_bits(
budget
.estimate_multiplier
.load(std::sync::atomic::Ordering::Acquire),
);
assert!(
(multiplier_after - multiplier_before).abs() < 1e-9,
"guard drop must not fold a (estimated, 0) sample into the EMA \
(before={multiplier_before}, after={multiplier_after})",
);
}
#[tokio::test]
async fn test_reservation_guard_drop_wakes_waiter() {
let budget = Arc::new(PipelineBudget::new(0, 1));
budget.set_multiplier(1.0);
let full = budget.budget;
let blocking_guard = budget.reserve_guarded(full).await;
assert_eq!(
budget.current.load(std::sync::atomic::Ordering::Acquire),
full,
);
let budget_clone = Arc::clone(&budget);
let waiter = tokio::spawn(async move { budget_clone.reserve(1).await });
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
assert!(!waiter.is_finished(), "second reserve should be parked");
drop(blocking_guard);
tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
.await
.expect("guard drop did not wake parked reserver")
.expect("waiter task panicked");
}
#[tokio::test]
async fn test_reservation_guard_drop_vacates_segments() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
{
let mut guards = Vec::new();
for i in 0..MAX_CONCURRENT_SEGMENTS {
guards.push(
budget
.reserve_guarded_with_priority(1, ti(0), vec![format!("seg{i}")])
.await,
);
}
assert_eq!(
budget.active_segments.lock().all.len(),
MAX_CONCURRENT_SEGMENTS,
"all segments should be admitted before the drops",
);
}
assert_eq!(
budget.active_segments.lock().all.len(),
0,
"dropped guards must vacate every segment slot, not just the bytes",
);
let reserved = budget
.reserve_with_priority(1, ti(0), &["fresh".to_owned()])
.await;
assert_eq!(reserved, 1);
assert!(budget.active_segments.lock().all.contains("fresh"));
}
#[tokio::test]
async fn test_reservation_guard_drop_wakes_segment_gated_waiter() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
for i in 0..(MAX_CONCURRENT_SEGMENTS - 1) {
budget
.reserve_with_priority(1, ti(0), &[format!("seg{i}")])
.await;
}
let guard = budget
.reserve_guarded_with_priority(1, ti(0), vec!["doomed".to_owned()])
.await;
let b = Arc::clone(&budget);
let waiter = tokio::spawn(async move {
b.reserve_with_priority(1, ti(0), &["late".to_owned()])
.await
});
for _ in 0..16 {
tokio::task::yield_now().await;
}
assert!(
!waiter.is_finished(),
"control: reserver must park on the segment cap"
);
drop(guard);
tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
.await
.expect("guard drop did not wake the segment-gated reserver")
.expect("waiter task panicked");
}
const DEFAULT_BYTES: usize = 64 * 1024 * 1024;
#[test]
fn test_parse_bytes_accepts_iec_suffix() {
assert_eq!(
parse_bytes_or_default("TEST", "128MiB", DEFAULT_BYTES),
128 * 1024 * 1024,
);
assert_eq!(
parse_bytes_or_default("TEST", "1GiB", DEFAULT_BYTES),
1024 * 1024 * 1024,
);
assert_eq!(
parse_bytes_or_default("TEST", "512KiB", DEFAULT_BYTES),
512 * 1024,
);
}
#[test]
fn test_parse_bytes_accepts_si_suffix() {
assert_eq!(
parse_bytes_or_default("TEST", "100MB", DEFAULT_BYTES),
100_000_000,
);
assert_eq!(
parse_bytes_or_default("TEST", "2GB", DEFAULT_BYTES),
2_000_000_000,
);
}
#[test]
fn test_parse_bytes_accepts_bare_integer_as_bytes() {
assert_eq!(
parse_bytes_or_default("TEST", "67108864", DEFAULT_BYTES),
64 * 1024 * 1024,
);
}
#[test]
fn test_parse_bytes_rejects_zero() {
assert_eq!(
parse_bytes_or_default("TEST", "0", DEFAULT_BYTES),
DEFAULT_BYTES,
);
}
#[test]
fn test_parse_bytes_rejects_negative() {
assert_eq!(
parse_bytes_or_default("TEST", "-1", DEFAULT_BYTES),
DEFAULT_BYTES,
);
assert_eq!(
parse_bytes_or_default("TEST", "-1MB", DEFAULT_BYTES),
DEFAULT_BYTES,
);
}
#[test]
fn test_parse_bytes_rejects_non_numeric() {
assert_eq!(
parse_bytes_or_default("TEST", "not-a-number", DEFAULT_BYTES),
DEFAULT_BYTES,
);
}
#[test]
fn test_parse_bytes_rejects_unknown_suffix() {
assert_eq!(
parse_bytes_or_default("TEST", "10Mb", DEFAULT_BYTES),
DEFAULT_BYTES,
);
}
#[test]
fn test_parse_fraction_accepts_valid_range() {
assert!((parse_fraction_or_default("TEST", "0.5", 0.25) - 0.5).abs() < 1e-12);
assert!((parse_fraction_or_default("TEST", "1.0", 0.25) - 1.0).abs() < 1e-12);
assert!((parse_fraction_or_default("TEST", "0.0001", 0.25) - 0.0001).abs() < 1e-12);
}
#[test]
fn test_parse_fraction_rejects_zero() {
assert!((parse_fraction_or_default("TEST", "0.0", 0.25) - 0.25).abs() < 1e-12);
}
#[test]
fn test_parse_fraction_rejects_above_one() {
assert!((parse_fraction_or_default("TEST", "1.5", 0.25) - 0.25).abs() < 1e-12);
}
#[test]
fn test_parse_fraction_rejects_negative() {
assert!((parse_fraction_or_default("TEST", "-0.5", 0.25) - 0.25).abs() < 1e-12);
}
#[test]
fn test_parse_fraction_rejects_nan_and_inf() {
assert!((parse_fraction_or_default("TEST", "NaN", 0.25) - 0.25).abs() < 1e-12);
assert!((parse_fraction_or_default("TEST", "inf", 0.25) - 0.25).abs() < 1e-12);
}
#[test]
fn test_parse_fraction_rejects_non_numeric() {
assert!((parse_fraction_or_default("TEST", "bogus", 0.25) - 0.25).abs() < 1e-12);
}
#[tokio::test]
async fn test_admission_metrics_record_combined_pressure_after_parking() {
let metrics = test_query_metrics();
let budget = Arc::new(PipelineBudget::with_exact_budget_and_metrics(
100,
Some(Arc::clone(&metrics)),
));
budget.set_multiplier(1.0);
let active = vec!["s1".to_owned(), "s2".to_owned(), "s3".to_owned()];
budget
.reserve_with_priority(100, TimeInt::MAX, &active)
.await;
let waiter_budget = Arc::clone(&budget);
let waiter = tokio::spawn(async move {
waiter_budget
.reserve_with_priority(1, TimeInt::MAX, &["s4".to_owned()])
.await;
});
for _ in 0..100 {
if metrics.segment_admission_waits.load(Acquire) == 1 {
break;
}
tokio::task::yield_now().await;
}
assert!(!waiter.is_finished());
assert_eq!(metrics.pipeline_byte_waits.load(Acquire), 1);
assert_eq!(metrics.segment_admission_waits.load(Acquire), 1);
assert_eq!(metrics.peak_active_segments.load(Acquire), 3);
waiter.abort();
}
#[test]
fn test_parse_usize_accepts_positive() {
assert_eq!(parse_usize_or_default("TEST", "64", 128), 64);
assert_eq!(parse_usize_or_default("TEST", "1", 128), 1);
}
#[test]
fn test_parse_usize_rejects_zero() {
assert_eq!(parse_usize_or_default("TEST", "0", 128), 128);
}
#[test]
fn test_parse_usize_rejects_negative() {
assert_eq!(parse_usize_or_default("TEST", "-1", 128), 128);
}
#[test]
fn test_parse_usize_rejects_non_numeric() {
assert_eq!(parse_usize_or_default("TEST", "not-a-number", 128), 128);
assert_eq!(parse_usize_or_default("TEST", "64MB", 128), 128);
}
fn ti(t: i64) -> TimeInt {
TimeInt::saturated_temporal_i64(t)
}
#[tokio::test]
async fn test_priority_wake_orders_by_task_time_min() {
let half = 1024 * 1024;
let budget = Arc::new(PipelineBudget::with_exact_budget(half * 2));
budget.set_multiplier(1.0);
budget.reserve_with_priority(half * 2, ti(0), &[]).await;
use std::sync::atomic::{AtomicU8, Ordering};
let order = Arc::new(AtomicU8::new(0));
let order_late = Arc::clone(&order);
let b_late = Arc::clone(&budget);
let late = tokio::spawn(async move {
b_late.reserve_with_priority(half, ti(100), &[]).await;
order_late.fetch_or(0b10, Ordering::AcqRel);
});
tokio::task::yield_now().await;
let order_early = Arc::clone(&order);
let b_early = Arc::clone(&budget);
let early = tokio::spawn(async move {
b_early.reserve_with_priority(half, ti(1), &[]).await;
order_early.fetch_or(0b01, Ordering::AcqRel);
});
tokio::task::yield_now().await;
budget.release(half);
early.await.unwrap();
assert_eq!(
order.load(Ordering::Acquire) & 0b01,
0b01,
"early-time waiter must wake first"
);
assert_eq!(
order.load(Ordering::Acquire) & 0b10,
0,
"late-time waiter must still be parked"
);
budget.release(half);
late.await.unwrap();
}
#[tokio::test]
async fn test_segment_count_gate_admits_atomically() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
for i in 0..MAX_CONCURRENT_SEGMENTS {
budget
.reserve_with_priority(1, ti(0), &[format!("seg{i}")])
.await;
}
let freed = format!("seg{}", MAX_CONCURRENT_SEGMENTS - 1);
budget.publish_segment_finalized(&freed);
let b = Arc::clone(&budget);
let three_new = tokio::spawn(async move {
b.reserve_with_priority(
1,
ti(0),
&["new-a".to_owned(), "new-b".to_owned(), "new-c".to_owned()],
)
.await;
});
for _ in 0..16 {
tokio::task::yield_now().await;
}
assert!(
!three_new.is_finished(),
"fetch must park because admitting 3 new segments would push past the cap"
);
for i in 0..(MAX_CONCURRENT_SEGMENTS - 1) {
budget.publish_segment_finalized(&format!("seg{i}"));
}
three_new.await.unwrap();
}
#[tokio::test]
async fn test_segment_count_gate_caps_at_max() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
for i in 0..MAX_CONCURRENT_SEGMENTS {
budget
.reserve_with_priority(1, ti(0), &[format!("seg{i}")])
.await;
}
assert_eq!(
budget.active_segments.lock().all.len(),
MAX_CONCURRENT_SEGMENTS
);
let b = Arc::clone(&budget);
let parked = tokio::spawn(async move {
b.reserve_with_priority(1, ti(0), &["overflow".to_owned()])
.await;
});
for _ in 0..16 {
tokio::task::yield_now().await;
}
assert!(!parked.is_finished());
budget.publish_segment_finalized("seg0");
parked.await.unwrap();
}
#[tokio::test]
async fn test_reserve_spanning_more_than_cap_segments_does_not_deadlock() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
let segments: Vec<String> = (0..(MAX_CONCURRENT_SEGMENTS + 4))
.map(|i| format!("seg{i}"))
.collect();
let b = Arc::clone(&budget);
let segs = segments.clone();
let reserve = tokio::spawn(async move { b.reserve_with_priority(1, ti(0), &segs).await });
for _ in 0..16 {
tokio::task::yield_now().await;
}
assert!(
reserve.is_finished(),
"reservation spanning >MAX_CONCURRENT_SEGMENTS segments must bypass the gate, not park"
);
let reserved = reserve.await.unwrap();
assert_eq!(reserved, 1);
let gate = budget.active_segments.lock();
assert_eq!(gate.all.len(), segments.len());
assert_eq!(gate.bypass.len(), segments.len());
assert_eq!(gate.effective_len(), 0);
}
#[tokio::test]
async fn test_stall_breaker_fires_after_threshold() {
let budget = Arc::new(PipelineBudget::with_exact_budget(100));
budget.set_multiplier(1.0);
budget.reserve_with_priority(100, ti(0), &[]).await;
for _ in 0..(STALL_EMPTY_EMIT_THRESHOLD - 1) {
budget.notify_empty_emit();
}
assert!(
!budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
budget.notify_empty_emit();
assert!(
budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
let extra = budget.reserve_with_priority(50, ti(0), &[]).await;
assert_eq!(extra, 50);
}
#[tokio::test]
async fn test_stall_breaker_bypasses_segment_count_gate() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
for i in 0..MAX_CONCURRENT_SEGMENTS {
budget
.reserve_with_priority(1, ti(0), &[format!("seg{i}")])
.await;
}
assert_eq!(
budget.active_segments.lock().effective_len(),
MAX_CONCURRENT_SEGMENTS
);
let b = Arc::clone(&budget);
let blocked = tokio::spawn(async move {
b.reserve_with_priority(1, ti(0), &["new-cap-blocked".to_owned()])
.await
});
for _ in 0..16 {
tokio::task::yield_now().await;
}
assert!(
!blocked.is_finished(),
"control: cap-blocked reserve must park without the bypass"
);
budget
.force_overcommit
.store(true, std::sync::atomic::Ordering::Release);
budget.wake_next();
let reserved = blocked.await.unwrap();
assert_eq!(reserved, 1);
let segments = budget.active_segments.lock();
assert!(segments.all.contains("new-cap-blocked"));
assert!(segments.bypass.contains("new-cap-blocked"));
}
#[tokio::test]
async fn test_stall_breaker_clears_on_progress() {
let budget = Arc::new(PipelineBudget::with_exact_budget(100));
budget.set_multiplier(1.0);
budget.reserve_with_priority(100, ti(0), &[]).await;
for _ in 0..STALL_EMPTY_EMIT_THRESHOLD {
budget.notify_empty_emit();
}
assert!(
budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
budget.notify_row_emitted();
assert!(
!budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
assert_eq!(
budget
.empty_emit_count
.load(std::sync::atomic::Ordering::Acquire),
0
);
for _ in 0..STALL_EMPTY_EMIT_THRESHOLD {
budget.notify_empty_emit();
}
assert!(
budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
budget.release(50);
assert!(
!budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
}
#[tokio::test]
async fn test_stall_breaker_ignores_unsaturated_budget() {
let budget = Arc::new(PipelineBudget::with_exact_budget(100));
budget.set_multiplier(1.0);
for _ in 0..(STALL_EMPTY_EMIT_THRESHOLD * 3) {
budget.notify_empty_emit();
}
assert!(
!budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
}
#[tokio::test]
async fn test_stall_breaker_does_not_carry_unsaturated_count_into_saturation() {
let budget = Arc::new(PipelineBudget::with_exact_budget(100));
budget.set_multiplier(1.0);
for _ in 0..(STALL_EMPTY_EMIT_THRESHOLD * 5) {
budget.notify_empty_emit();
}
assert_eq!(
budget
.empty_emit_count
.load(std::sync::atomic::Ordering::Acquire),
0,
"unsaturated empty emits must reset the counter, not accumulate"
);
budget.reserve_with_priority(100, ti(0), &[]).await;
budget.notify_empty_emit();
assert!(
!budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire),
"first saturated empty emit must not trip the breaker"
);
assert_eq!(
budget
.empty_emit_count
.load(std::sync::atomic::Ordering::Acquire),
1,
);
for _ in 0..(STALL_EMPTY_EMIT_THRESHOLD - 1) {
budget.notify_empty_emit();
}
assert!(
budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
}
#[tokio::test]
async fn test_segment_cap_self_heals_after_bypass() {
let budget = Arc::new(PipelineBudget::with_exact_budget(1 << 30));
budget.set_multiplier(1.0);
for i in 0..MAX_CONCURRENT_SEGMENTS {
budget
.reserve_with_priority(1, ti(0), &[format!("seg{i}")])
.await;
}
assert_eq!(
budget.active_segments.lock().effective_len(),
MAX_CONCURRENT_SEGMENTS
);
budget
.force_overcommit
.store(true, std::sync::atomic::Ordering::Release);
budget
.reserve_with_priority(1, ti(0), &["bypass-a".to_owned()])
.await;
{
let segments = budget.active_segments.lock();
assert_eq!(segments.all.len(), MAX_CONCURRENT_SEGMENTS + 1);
assert_eq!(segments.bypass.len(), 1);
assert_eq!(segments.effective_len(), MAX_CONCURRENT_SEGMENTS);
}
budget
.force_overcommit
.store(false, std::sync::atomic::Ordering::Release);
let b = Arc::clone(&budget);
let parked = tokio::spawn(async move {
b.reserve_with_priority(1, ti(0), &["after-bypass".to_owned()])
.await;
});
for _ in 0..16 {
tokio::task::yield_now().await;
}
assert!(!parked.is_finished());
budget.publish_segment_finalized("seg0");
parked.await.unwrap();
let segments = budget.active_segments.lock();
assert_eq!(segments.all.len(), MAX_CONCURRENT_SEGMENTS + 1);
assert_eq!(segments.bypass.len(), 1);
assert_eq!(segments.effective_len(), MAX_CONCURRENT_SEGMENTS);
}
#[tokio::test]
async fn test_wake_next_skips_cancelled_orphans() {
use std::sync::atomic::Ordering::Acquire;
let budget = PipelineBudget::with_exact_budget(100);
let orphan_notify = Arc::new(Notify::new());
let orphan_cancelled = Arc::new(AtomicBool::new(true));
let real_notify = Arc::new(Notify::new());
let real_cancelled = Arc::new(AtomicBool::new(false));
{
let mut queue = budget.wait_queue.lock();
queue.push(Reverse(PriorityWaiter {
task_time_min: ti(1),
seq: 0,
notify: Arc::clone(&orphan_notify),
cancelled: Arc::clone(&orphan_cancelled),
reserved_bytes: 0,
segment_ids: Vec::new(),
}));
queue.push(Reverse(PriorityWaiter {
task_time_min: ti(1_000),
seq: 1,
notify: Arc::clone(&real_notify),
cancelled: Arc::clone(&real_cancelled),
reserved_bytes: 1,
segment_ids: Vec::new(),
}));
}
budget.wake_next();
assert_eq!(budget.wait_queue.lock().len(), 0);
let real_received = tokio::time::timeout(
std::time::Duration::from_millis(100),
real_notify.notified(),
)
.await;
assert!(
real_received.is_ok(),
"real waiter must receive its wake despite the lower-priority orphan ahead of it"
);
assert!(orphan_cancelled.load(Acquire));
}
#[tokio::test]
async fn test_adjust_reservation_shrink_resets_stall_detector() {
let budget = Arc::new(PipelineBudget::with_exact_budget(100));
budget.set_multiplier(1.0);
let reserved = budget.reserve_with_priority(100, ti(0), &[]).await;
for _ in 0..STALL_EMPTY_EMIT_THRESHOLD {
budget.notify_empty_emit();
}
assert!(
budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire)
);
budget.adjust_reservation(reserved, reserved, 30);
assert!(
!budget
.force_overcommit
.load(std::sync::atomic::Ordering::Acquire),
"shrink that frees budget must clear `force_overcommit`"
);
assert_eq!(
budget
.empty_emit_count
.load(std::sync::atomic::Ordering::Acquire),
0,
"shrink that frees budget must reset the empty-emit counter"
);
}
#[tokio::test]
async fn test_release_wakes_admittable_not_segment_blocked_waiter() {
let budget = Arc::new(PipelineBudget::with_exact_budget(100));
budget.set_multiplier(1.0);
budget
.reserve_with_priority(60, ti(5), &["seg0".to_owned()])
.await;
budget
.reserve_with_priority(20, ti(5), &["seg1".to_owned()])
.await;
budget
.reserve_with_priority(20, ti(5), &["seg2".to_owned()])
.await;
let b1 = Arc::clone(&budget);
let w1 = tokio::spawn(async move {
b1.reserve_with_priority(10, ti(1), &["seg-new".to_owned()])
.await
});
let b2 = Arc::clone(&budget);
let w2 = tokio::spawn(async move {
b2.reserve_with_priority(40, ti(9), &["seg0".to_owned()])
.await
});
for _ in 0..32 {
tokio::task::yield_now().await;
}
assert!(
!w1.is_finished(),
"control: W1 must park on the segment gate"
);
assert!(!w2.is_finished(), "control: W2 must park on bytes");
budget.release(40);
tokio::time::timeout(std::time::Duration::from_secs(1), w2)
.await
.expect(
"freed bytes must wake the admittable byte-only waiter, not the segment-blocked one",
)
.expect("W2 panicked");
assert!(
!w1.is_finished(),
"segment-gate-blocked W1 must stay parked despite higher priority",
);
budget.publish_segment_finalized("seg1");
budget.release(60);
tokio::time::timeout(std::time::Duration::from_secs(1), w1)
.await
.expect("W1 should finish once its segment slot and bytes free")
.expect("W1 panicked");
}