Skip to main content

objectstore_service/
concurrency.rs

1//! Concurrency limiter for backend operations.
2//!
3//! [`ConcurrencyLimiter`] caps the number of in-flight backend operations
4//! using a tokio semaphore. Each acquired [`ConcurrencyPermit`] notifies
5//! waiters on drop, allowing [`ConcurrencyLimiter::wait_all`] to resolve once
6//! all permits have been returned.
7//!
8//! [`spawn_metered`] spawns an arbitrary future as an isolated task with panic
9//! recovery and `service.task.*` metric emission.
10
11use std::future::Future;
12use std::sync::Arc;
13use std::time::Duration;
14
15use futures_util::FutureExt;
16use sentry::{Hub, SentryFutureExt, TransactionContext};
17use tokio::sync::{AcquireError, Notify, OwnedSemaphorePermit, Semaphore};
18
19use crate::error::{Error, Result};
20
21/// Interval for the periodic metrics emitter.
22const EMITTER_INTERVAL: Duration = Duration::from_secs(1);
23
24/// Snapshot of concurrency limiter state.
25///
26/// Passed to the callback registered via
27/// [`ConcurrencyLimiter::run_emitter`].
28#[non_exhaustive]
29#[derive(Clone, Copy, Debug)]
30pub struct Stats {
31    /// Number of execution permits currently held.
32    pub in_use: u32,
33    /// Number of callers waiting in the queue for a permit.
34    pub queued: u32,
35    /// Number of bulk operations currently in flight.
36    pub bulk_in_use: u32,
37}
38
39/// Limits concurrent backend operations and tracks in-flight count.
40///
41/// Permits are acquired with [`acquire`](Self::acquire) or
42/// [`acquire_bulk`](Self::acquire_bulk) and automatically returned when
43/// the [`ConcurrencyPermit`] is dropped.
44///
45/// Bulk operations use a separate budget semaphore that limits how many
46/// execution slots they may occupy. This is intended as a safe operating
47/// point — below this level there should be little-to-no performance
48/// degradation, leaving room for more tasks to be admitted via the queue
49/// before rejection is necessary.
50#[derive(Clone, Debug)]
51pub struct ConcurrencyLimiter {
52    tasks: Arc<Semaphore>,
53    queue: Arc<Semaphore>,
54    bulk: Arc<Semaphore>,
55    tasks_total: u32,
56    queue_total: u32,
57    bulk_total: u32,
58    timeout: Duration,
59    released: Arc<Notify>,
60}
61
62impl ConcurrencyLimiter {
63    /// Creates a new limiter with the given maximum number of permits.
64    ///
65    /// By default the queue depth is zero, preserving the original try-or-reject behavior. Use
66    /// [`with_queue`](Self::with_queue) to enable bounded waiting.
67    ///
68    /// The bulk budget defaults to 100% of `max` (no restriction); use
69    /// [`with_bulk`](Self::with_bulk) to set a safe operating point for bulk traffic.
70    pub fn new(max: u32) -> Self {
71        Self {
72            tasks: Arc::new(Semaphore::new(max as usize)),
73            queue: Arc::new(Semaphore::new(0)),
74            bulk: Arc::new(Semaphore::new(max as usize)),
75            tasks_total: max,
76            queue_total: 0,
77            bulk_total: max,
78            timeout: Duration::from_secs(1),
79            released: Arc::new(Notify::new()),
80        }
81    }
82
83    /// Enables bounded waiting when all execution permits are held.
84    ///
85    /// Up to `size` additional callers may park in [`acquire`](Self::acquire)
86    /// waiting for a permit. Callers beyond that are rejected immediately.
87    pub fn with_queue(mut self, size: u32) -> Self {
88        self.queue_total = size;
89        self.queue = Arc::new(Semaphore::new(size as usize));
90        self
91    }
92
93    /// Sets the maximum time a caller may wait for a permit.
94    ///
95    /// Applies to both [`acquire`](Self::acquire) (when parked in the
96    /// queue) and [`acquire_bulk`](Self::acquire_bulk) (waiting for the
97    /// bulk and execution semaphores). Defaults to 1 second.
98    pub fn with_timeout(mut self, timeout: Duration) -> Self {
99        self.timeout = timeout;
100        self
101    }
102
103    /// Sets the bulk concurrency budget as a percentage of `max`.
104    ///
105    /// `percent` is clamped to `1..=100`. At `100` (the default), bulk
106    /// operations can use all execution slots. Lower values set a safe
107    /// operating point below which there is little-to-no performance
108    /// degradation, allowing more tasks to queue before rejection is
109    /// necessary — e.g. `60` means bulk operations can hold at most 60%
110    /// of permits.
111    pub fn with_bulk(mut self, percent: u32) -> Self {
112        let clamped = percent.min(100);
113        self.bulk_total = (self.tasks_total * clamped).div_ceil(100).max(1);
114        self.bulk = Arc::new(Semaphore::new(self.bulk_total as usize));
115        self
116    }
117
118    /// Acquires a single concurrency permit, waiting if necessary.
119    ///
120    /// If a permit is free, returns immediately without touching the
121    /// queue. Otherwise, acquires a queue ticket (bounded by the queue
122    /// depth) and waits up to the configured timeout. Returns
123    /// [`Error::AtCapacity`] if the queue is full or on timeout.
124    pub async fn acquire(&self) -> Result<ConcurrencyPermit> {
125        if self.tasks_total == 0 {
126            return Err(Error::AtCapacity);
127        }
128
129        // Fast path: Instantly grab a free permit without parking.
130        if let Ok(task_permit) = self.tasks.clone().try_acquire_owned() {
131            return Ok(ConcurrencyPermit {
132                task_permit: Some(task_permit),
133                bulk_permit: None,
134                released: Arc::clone(&self.released),
135            });
136        }
137
138        // Slow path: acquire a temporary queue ticket to bound concurrent
139        // waiters. Released in this scope when the permit is constructed.
140        let _ticket = self
141            .queue
142            .clone()
143            .try_acquire_owned()
144            .map_err(|_| Error::AtCapacity)?;
145
146        let acquire = self.tasks.clone().acquire_owned();
147        let task_permit = tokio::time::timeout(self.timeout, acquire)
148            .await
149            .map_err(|_| Error::AtCapacity)?
150            .map_err(|_| Error::AtCapacity)?;
151
152        Ok(ConcurrencyPermit {
153            task_permit: Some(task_permit),
154            bulk_permit: None,
155            released: Arc::clone(&self.released),
156        })
157    }
158
159    /// Tries to acquire a single permit without waiting.
160    ///
161    /// Returns [`Error::AtCapacity`] when no permits are available.
162    pub fn try_acquire(&self) -> Result<ConcurrencyPermit> {
163        let task_permit = self
164            .tasks
165            .clone()
166            .try_acquire_owned()
167            .map_err(|_| Error::AtCapacity)?;
168
169        Ok(ConcurrencyPermit {
170            task_permit: Some(task_permit),
171            bulk_permit: None,
172            released: Arc::clone(&self.released),
173        })
174    }
175
176    /// Acquires a single permit for a bulk operation, waiting if necessary.
177    ///
178    /// Bulk operations are bounded by the bulk budget — a safe operating
179    /// point below which there is little-to-no performance degradation.
180    /// Both the bulk semaphore and the inner execution semaphore are
181    /// acquired under a single timeout deadline configured via
182    /// [`with_timeout`](Self::with_timeout).
183    ///
184    /// Returns [`Error::AtCapacity`] on timeout or when `max` is zero.
185    pub async fn acquire_bulk(&self) -> Result<ConcurrencyPermit> {
186        if self.tasks_total == 0 {
187            return Err(Error::AtCapacity);
188        }
189
190        let bulk_sem = self.bulk.clone();
191        let tasks_sem = self.tasks.clone();
192
193        let acquire = async move {
194            let bulk_permit = bulk_sem.acquire_owned().await?;
195            let task_permit = tasks_sem.acquire_owned().await?;
196            Ok((task_permit, bulk_permit))
197        };
198
199        let (task_permit, bulk_permit) = tokio::time::timeout(self.timeout, acquire)
200            .await
201            .map_err(|_| Error::AtCapacity)?
202            .map_err(|_: AcquireError| Error::AtCapacity)?;
203
204        Ok(ConcurrencyPermit {
205            task_permit: Some(task_permit),
206            bulk_permit: Some(bulk_permit),
207            released: Arc::clone(&self.released),
208        })
209    }
210
211    /// Returns the number of permits currently available.
212    pub fn available_permits(&self) -> u32 {
213        u32::try_from(self.tasks.available_permits()).unwrap_or(self.tasks_total)
214    }
215
216    /// Returns the number of permits currently held.
217    pub fn used_permits(&self) -> u32 {
218        self.tasks_total - self.available_permits()
219    }
220
221    /// Returns the total number of permits.
222    pub fn total_permits(&self) -> u32 {
223        self.tasks_total
224    }
225
226    /// Returns the number of callers currently waiting in the queue.
227    pub fn queued_permits(&self) -> u32 {
228        let available = u32::try_from(self.queue.available_permits()).unwrap_or(self.queue_total);
229        self.queue_total - available
230    }
231
232    /// Returns the configured queue capacity.
233    pub fn total_queue(&self) -> u32 {
234        self.queue_total
235    }
236
237    /// Returns the number of bulk permits currently held.
238    pub fn used_bulk_permits(&self) -> u32 {
239        let available = u32::try_from(self.bulk.available_permits()).unwrap_or(self.bulk_total);
240        self.bulk_total - available
241    }
242
243    /// Returns the bulk concurrency budget.
244    pub fn total_bulk(&self) -> u32 {
245        self.bulk_total
246    }
247
248    /// Waits until all permits have been returned.
249    #[allow(dead_code)]
250    pub async fn wait_all(&self) {
251        loop {
252            let notified = self.released.notified();
253            if self.used_permits() == 0 {
254                return;
255            }
256            notified.await;
257        }
258    }
259
260    /// Returns a snapshot of the current counts.
261    pub fn stats(&self) -> Stats {
262        Stats {
263            in_use: self.used_permits(),
264            queued: self.queued_permits(),
265            bulk_in_use: self.used_bulk_permits(),
266        }
267    }
268
269    /// Periodically calls `emit` with the current in-use and queued counts.
270    ///
271    /// This future runs forever and is intended to be spawned as a background
272    /// task alongside the service.
273    pub async fn run_emitter<F, Fut>(&self, mut emit: F)
274    where
275        F: FnMut(Stats) -> Fut,
276        Fut: Future<Output = ()>,
277    {
278        let mut ticker = tokio::time::interval(EMITTER_INTERVAL);
279        loop {
280            ticker.tick().await;
281            emit(self.stats()).await;
282        }
283    }
284}
285
286/// RAII guard for a concurrency permit.
287///
288/// Dropping this permit releases it back to the [`ConcurrencyLimiter`] and
289/// notifies any task waiting in [`ConcurrencyLimiter::wait_all`].
290///
291/// The execution permit drops before the bulk permit so that a waiter
292/// blocked on the task semaphore sees the freed slot immediately.
293pub struct ConcurrencyPermit {
294    task_permit: Option<OwnedSemaphorePermit>,
295    bulk_permit: Option<OwnedSemaphorePermit>,
296    released: Arc<Notify>,
297}
298
299impl std::fmt::Debug for ConcurrencyPermit {
300    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
301        f.debug_struct("ConcurrencyPermit").finish_non_exhaustive()
302    }
303}
304
305impl Drop for ConcurrencyPermit {
306    fn drop(&mut self) {
307        drop(self.task_permit.take());
308        drop(self.bulk_permit.take());
309        self.released.notify_waiters();
310    }
311}
312
313/// Spawns a future on a dedicated task with panic isolation and timing metrics.
314///
315/// The `guard` is moved into the spawned task and dropped after the future
316/// completes, ensuring any resource it represents (e.g. a concurrency permit)
317/// outlives the operation.
318///
319/// Emits `service.task.start` (counter) before spawning and
320/// `service.task.duration` (distribution) when the task completes, both tagged
321/// with the given `operation` name. The duration tag includes an `outcome` of
322/// `"success"` or `"error"`.
323pub async fn spawn_metered<T, G, F>(operation: &'static str, guard: G, f: F) -> Result<T>
324where
325    T: Send + 'static,
326    G: Send + 'static,
327    F: Future<Output = Result<T>> + Send + 'static,
328{
329    objectstore_metrics::count!("service.task.start", operation = operation);
330
331    let hub = Hub::current();
332    let span = hub.configure_scope(|scope| scope.get_span());
333
334    let new_hub = Hub::new_from_top(hub);
335    let transaction = new_hub.start_transaction(TransactionContext::continue_from_span(
336        operation,
337        "tokio.task",
338        span,
339    ));
340
341    let scope_guard = new_hub.push_scope();
342    new_hub.configure_scope(|scope| scope.set_span(Some(transaction.clone().into())));
343
344    let (tx, rx) = tokio::sync::oneshot::channel();
345    tokio::spawn(
346        async move {
347            let start = tokio::time::Instant::now();
348            let result = std::panic::AssertUnwindSafe(f)
349                .catch_unwind()
350                .await
351                .unwrap_or_else(|payload| Err(Error::panic(payload)));
352
353            if let Err(ref e) = result {
354                let error = e as &dyn std::error::Error;
355                objectstore_log::event_dyn!(e.level(), error, operation, "Task failed");
356            }
357
358            objectstore_metrics::record!(
359                "service.task.duration" = start.elapsed(),
360                operation = operation,
361                outcome = if result.is_ok() { "success" } else { "error" },
362            );
363
364            let _ = tx.send(result);
365            drop(guard);
366            transaction.finish();
367            drop(scope_guard);
368        }
369        .bind_hub(new_hub),
370    );
371
372    rx.await.map_err(|_| {
373        objectstore_log::error!(!!&Error::Dropped, operation, "Task failed");
374        Error::Dropped
375    })?
376}
377
378#[cfg(test)]
379mod tests {
380    use std::sync::atomic::{AtomicU32, Ordering};
381
382    use super::*;
383    use crate::error::Error;
384
385    #[test]
386    fn available_permits_tracks_held() {
387        let limiter = ConcurrencyLimiter::new(5);
388        assert_eq!(limiter.available_permits(), 5);
389
390        let p1 = limiter.try_acquire().unwrap();
391        assert_eq!(limiter.available_permits(), 4);
392
393        let p2 = limiter.try_acquire().unwrap();
394        assert_eq!(limiter.available_permits(), 3);
395
396        drop(p1);
397        assert_eq!(limiter.available_permits(), 4);
398
399        drop(p2);
400        assert_eq!(limiter.available_permits(), 5);
401    }
402
403    #[test]
404    fn total_permits_returns_configured_max() {
405        let limiter = ConcurrencyLimiter::new(42);
406        assert_eq!(limiter.total_permits(), 42);
407    }
408
409    #[test]
410    fn acquire_and_release() {
411        let limiter = ConcurrencyLimiter::new(2);
412        assert_eq!(limiter.used_permits(), 0);
413
414        let p1 = limiter.try_acquire().unwrap();
415        assert_eq!(limiter.used_permits(), 1);
416
417        let p2 = limiter.try_acquire().unwrap();
418        assert_eq!(limiter.used_permits(), 2);
419
420        drop(p1);
421        assert_eq!(limiter.used_permits(), 1);
422
423        drop(p2);
424        assert_eq!(limiter.used_permits(), 0);
425    }
426
427    #[test]
428    fn at_capacity_rejects() {
429        let limiter = ConcurrencyLimiter::new(1);
430        let _permit = limiter.try_acquire().unwrap();
431
432        let result = limiter.try_acquire();
433        assert!(matches!(result, Err(Error::AtCapacity)));
434    }
435
436    #[test]
437    fn permit_recovery_after_drop() {
438        let limiter = ConcurrencyLimiter::new(1);
439
440        let permit = limiter.try_acquire().unwrap();
441        assert!(limiter.try_acquire().is_err());
442
443        drop(permit);
444        assert!(limiter.try_acquire().is_ok());
445    }
446
447    #[tokio::test(start_paused = true)]
448    async fn emitter_calls_callback() {
449        let limiter = ConcurrencyLimiter::new(5);
450        let _permit = limiter.try_acquire().unwrap();
451
452        let emitted_in_use = Arc::new(AtomicU32::new(0));
453        let emitted_queued = Arc::new(AtomicU32::new(0));
454        let in_use_clone = Arc::clone(&emitted_in_use);
455        let queued_clone = Arc::clone(&emitted_queued);
456
457        let emitter = limiter.run_emitter(move |stats| {
458            let in_use_ref = Arc::clone(&in_use_clone);
459            let queued_ref = Arc::clone(&queued_clone);
460            async move {
461                in_use_ref.store(stats.in_use, Ordering::Relaxed);
462                queued_ref.store(stats.queued, Ordering::Relaxed);
463            }
464        });
465
466        tokio::select! {
467            _ = emitter => unreachable!("emitter runs forever"),
468            _ = tokio::time::sleep(EMITTER_INTERVAL) => {}
469        }
470
471        assert_eq!(emitted_in_use.load(Ordering::Relaxed), 1);
472        assert_eq!(emitted_queued.load(Ordering::Relaxed), 0);
473    }
474
475    #[tokio::test]
476    async fn wait_all_resolves_when_permits_returned() {
477        let limiter = ConcurrencyLimiter::new(2);
478        let p1 = limiter.try_acquire().unwrap();
479        let p2 = limiter.try_acquire().unwrap();
480
481        let mut wait = Box::pin(limiter.wait_all());
482
483        // Dropping one permit is not enough.
484        drop(p1);
485        assert!(futures::poll!(&mut wait).is_pending());
486
487        // Dropping the last permit should resolve it.
488        drop(p2);
489        assert!(futures::poll!(&mut wait).is_ready());
490    }
491
492    #[tokio::test]
493    async fn wait_all_returns_immediately_when_empty() {
494        let limiter = ConcurrencyLimiter::new(5);
495        let wait = Box::pin(limiter.wait_all());
496        assert!(futures::poll!(wait).is_ready());
497    }
498
499    // --- Queue tests ---
500
501    #[tokio::test(start_paused = true)]
502    async fn queue_zero_rejects_immediately() {
503        let limiter = ConcurrencyLimiter::new(2);
504        assert_eq!(limiter.total_queue(), 0);
505
506        let p1 = limiter.try_acquire().unwrap();
507        let p2 = limiter.try_acquire().unwrap();
508        assert!(matches!(limiter.try_acquire(), Err(Error::AtCapacity)));
509
510        drop(p1);
511        assert!(limiter.try_acquire().is_ok());
512        drop(p2);
513
514        // Bulk holding all permits also rejects instantly (no timeout wait).
515        let mut bulk_permits = Vec::new();
516        for _ in 0..2 {
517            let permit = limiter.acquire_bulk().await.unwrap();
518            bulk_permits.push(permit);
519        }
520
521        let start = tokio::time::Instant::now();
522        let result = limiter.acquire().await;
523        assert!(matches!(result, Err(Error::AtCapacity)));
524        assert_eq!(start.elapsed(), Duration::ZERO);
525        drop(bulk_permits);
526    }
527
528    #[tokio::test(start_paused = true)]
529    async fn acquire_succeeds_immediately_when_available() {
530        let limiter = ConcurrencyLimiter::new(2).with_queue(3);
531
532        let permit = limiter.acquire().await.unwrap();
533        assert_eq!(limiter.used_permits(), 1);
534        assert_eq!(limiter.queued_permits(), 0);
535        drop(permit);
536    }
537
538    #[tokio::test(start_paused = true)]
539    async fn acquire_waits_and_succeeds_after_release() {
540        let limiter = ConcurrencyLimiter::new(1).with_queue(2);
541
542        let held = limiter.acquire().await.unwrap();
543        assert_eq!(limiter.used_permits(), 1);
544
545        let limiter2 = limiter.clone();
546        let waiter = tokio::spawn(async move { limiter2.acquire().await });
547
548        tokio::task::yield_now().await;
549        assert_eq!(limiter.queued_permits(), 1);
550
551        drop(held);
552
553        let permit = waiter.await.unwrap().unwrap();
554        assert_eq!(limiter.used_permits(), 1);
555        assert_eq!(limiter.queued_permits(), 0);
556        drop(permit);
557    }
558
559    #[tokio::test(start_paused = true)]
560    async fn acquire_times_out() {
561        let limiter = ConcurrencyLimiter::new(1).with_queue(2);
562
563        let _held = limiter.acquire().await.unwrap();
564
565        let limiter2 = limiter.clone();
566        let waiter = tokio::spawn(async move { limiter2.acquire().await });
567
568        tokio::time::sleep(Duration::from_secs(2)).await;
569
570        let result = waiter.await.unwrap();
571        assert!(matches!(result, Err(Error::AtCapacity)));
572        assert_eq!(limiter.queued_permits(), 0);
573    }
574
575    #[tokio::test(start_paused = true)]
576    async fn acquire_rejects_over_max_plus_queue() {
577        let limiter = ConcurrencyLimiter::new(1).with_queue(1);
578
579        let _held = limiter.acquire().await.unwrap();
580
581        let limiter2 = limiter.clone();
582        let _waiter = tokio::spawn(async move { limiter2.acquire().await });
583        tokio::task::yield_now().await;
584
585        assert_eq!(limiter.queued_permits(), 1);
586        let result = limiter.acquire().await;
587        assert!(matches!(result, Err(Error::AtCapacity)));
588    }
589
590    #[tokio::test(start_paused = true)]
591    async fn dropping_parked_acquire_releases_queue_slot() {
592        let limiter = ConcurrencyLimiter::new(1).with_queue(1);
593
594        let _held = limiter.acquire().await.unwrap();
595
596        let limiter2 = limiter.clone();
597        let waiter = tokio::spawn(async move { limiter2.acquire().await });
598        tokio::task::yield_now().await;
599        assert_eq!(limiter.queued_permits(), 1);
600
601        waiter.abort();
602        let _ = waiter.await;
603        tokio::task::yield_now().await;
604
605        assert_eq!(limiter.queued_permits(), 0);
606
607        let limiter3 = limiter.clone();
608        let replacement = tokio::spawn(async move { limiter3.acquire().await });
609        tokio::task::yield_now().await;
610        assert_eq!(limiter.queued_permits(), 1);
611        drop(replacement);
612    }
613
614    #[tokio::test(start_paused = true)]
615    async fn queued_permits_reflects_state() {
616        let limiter = ConcurrencyLimiter::new(2).with_queue(3);
617
618        assert_eq!(limiter.queued_permits(), 0);
619        let _p1 = limiter.try_acquire().unwrap();
620        assert_eq!(limiter.queued_permits(), 0);
621        let _p2 = limiter.try_acquire().unwrap();
622        assert_eq!(limiter.queued_permits(), 0);
623        drop(_p1);
624        drop(_p2);
625
626        // Exact count when bulk holds permits.
627        let _bulk = limiter.acquire_bulk().await.unwrap();
628        let _bulk2 = limiter.acquire_bulk().await.unwrap();
629        assert_eq!(limiter.queued_permits(), 0);
630
631        let limiter2 = limiter.clone();
632        let _waiter = tokio::spawn(async move { limiter2.acquire().await });
633        tokio::task::yield_now().await;
634        assert_eq!(limiter.queued_permits(), 1);
635    }
636
637    #[tokio::test(start_paused = true)]
638    async fn acquire_rejects_immediately_when_max_is_zero() {
639        let limiter = ConcurrencyLimiter::new(0).with_queue(5);
640
641        let start = tokio::time::Instant::now();
642        let result = limiter.acquire().await;
643        assert!(matches!(result, Err(Error::AtCapacity)));
644        assert_eq!(start.elapsed(), Duration::ZERO);
645    }
646
647    #[tokio::test(start_paused = true)]
648    async fn emitter_reports_queued_count() {
649        let limiter = ConcurrencyLimiter::new(1).with_queue(2);
650        let _held = limiter.acquire().await.unwrap();
651
652        let limiter2 = limiter.clone();
653        let _waiter = tokio::spawn(async move { limiter2.acquire().await });
654        tokio::task::yield_now().await;
655
656        let emitted_in_use = Arc::new(AtomicU32::new(0));
657        let emitted_queued = Arc::new(AtomicU32::new(0));
658        let in_use_clone = Arc::clone(&emitted_in_use);
659        let queued_clone = Arc::clone(&emitted_queued);
660
661        let emitter = limiter.run_emitter(move |stats| {
662            let in_use_ref = Arc::clone(&in_use_clone);
663            let queued_ref = Arc::clone(&queued_clone);
664            async move {
665                in_use_ref.store(stats.in_use, Ordering::Relaxed);
666                queued_ref.store(stats.queued, Ordering::Relaxed);
667            }
668        });
669
670        tokio::select! {
671            _ = emitter => unreachable!("emitter runs forever"),
672            _ = tokio::time::sleep(EMITTER_INTERVAL) => {}
673        }
674
675        assert_eq!(emitted_in_use.load(Ordering::Relaxed), 1);
676        assert_eq!(emitted_queued.load(Ordering::Relaxed), 1);
677    }
678
679    // --- Bulk tests ---
680
681    #[test]
682    fn bulk_defaults_to_full_capacity() {
683        let limiter = ConcurrencyLimiter::new(100);
684        assert_eq!(limiter.total_bulk(), 100);
685    }
686
687    #[test]
688    fn bulk_percent_computes_correctly() {
689        let limiter = ConcurrencyLimiter::new(100).with_bulk(60);
690        assert_eq!(limiter.total_bulk(), 60);
691
692        let limiter = ConcurrencyLimiter::new(10).with_bulk(90);
693        assert_eq!(limiter.total_bulk(), 9);
694
695        let limiter = ConcurrencyLimiter::new(100).with_bulk(150);
696        assert_eq!(limiter.total_bulk(), 100);
697
698        // Allow at least one bulk permit even when the percentage is zero.
699        let limiter = ConcurrencyLimiter::new(1).with_bulk(60);
700        assert_eq!(limiter.total_bulk(), 1);
701    }
702
703    #[tokio::test(start_paused = true)]
704    async fn bulk_caps_at_budget() {
705        let limiter = ConcurrencyLimiter::new(10).with_queue(5).with_bulk(90);
706        let bulk_budget = limiter.total_bulk();
707        assert_eq!(bulk_budget, 9);
708
709        let mut permits = Vec::new();
710        for _ in 0..bulk_budget {
711            let permit = limiter.acquire_bulk().await.unwrap();
712            permits.push(permit);
713        }
714
715        assert_eq!(limiter.used_bulk_permits(), bulk_budget);
716        assert_eq!(limiter.used_permits(), bulk_budget);
717
718        // Normal request can still acquire the remaining permit.
719        let normal = limiter.acquire().await.unwrap();
720        assert_eq!(limiter.used_permits(), 10);
721        drop(normal);
722        drop(permits);
723    }
724
725    #[tokio::test(start_paused = true)]
726    async fn normal_uses_all_permits_when_bulk_idle() {
727        let limiter = ConcurrencyLimiter::new(10).with_queue(5);
728
729        let mut permits = Vec::new();
730        for _ in 0..10 {
731            let permit = limiter.acquire().await.unwrap();
732            permits.push(permit);
733        }
734
735        assert_eq!(limiter.used_permits(), 10);
736        assert_eq!(limiter.used_bulk_permits(), 0);
737        drop(permits);
738    }
739
740    #[tokio::test(start_paused = true)]
741    async fn bulk_waits_for_inner_permit() {
742        let limiter = ConcurrencyLimiter::new(1).with_queue(2);
743
744        let held = limiter.acquire().await.unwrap();
745        assert_eq!(limiter.used_permits(), 1);
746
747        let limiter2 = limiter.clone();
748        let waiter = tokio::spawn(async move { limiter2.acquire_bulk().await });
749        tokio::task::yield_now().await;
750
751        drop(held);
752
753        let permit = waiter.await.unwrap().unwrap();
754        assert_eq!(limiter.used_permits(), 1);
755        assert_eq!(limiter.used_bulk_permits(), 1);
756        drop(permit);
757    }
758
759    #[tokio::test(start_paused = true)]
760    async fn bulk_timeout_spans_both_waits() {
761        let limiter = ConcurrencyLimiter::new(1);
762
763        let _held = limiter.acquire().await.unwrap();
764
765        let limiter2 = limiter.clone();
766        let waiter = tokio::spawn(async move { limiter2.acquire_bulk().await });
767
768        tokio::time::sleep(Duration::from_secs(2)).await;
769
770        let result = waiter.await.unwrap();
771        assert!(matches!(result, Err(Error::AtCapacity)));
772        assert_eq!(limiter.used_bulk_permits(), 0);
773    }
774
775    #[tokio::test(start_paused = true)]
776    async fn bulk_cancellation_leaks_nothing() {
777        let limiter = ConcurrencyLimiter::new(2).with_queue(2);
778
779        let _held1 = limiter.acquire().await.unwrap();
780        let _held2 = limiter.acquire().await.unwrap();
781
782        let limiter2 = limiter.clone();
783        let waiter = tokio::spawn(async move { limiter2.acquire_bulk().await });
784        tokio::task::yield_now().await;
785
786        waiter.abort();
787        let _ = waiter.await;
788        tokio::task::yield_now().await;
789
790        assert_eq!(limiter.used_bulk_permits(), 0);
791        assert_eq!(limiter.used_permits(), 2);
792    }
793
794    #[tokio::test(start_paused = true)]
795    async fn bulk_rejects_immediately_when_max_is_zero() {
796        let limiter = ConcurrencyLimiter::new(0).with_queue(5);
797
798        let start = tokio::time::Instant::now();
799        let result = limiter.acquire_bulk().await;
800        assert!(matches!(result, Err(Error::AtCapacity)));
801        assert_eq!(start.elapsed(), Duration::ZERO);
802    }
803
804    #[tokio::test(start_paused = true)]
805    async fn bulk_with_zero_percent_allows_one() {
806        let limiter = ConcurrencyLimiter::new(10).with_bulk(0);
807
808        assert_eq!(limiter.total_bulk(), 1);
809
810        let permit = limiter.acquire_bulk().await.unwrap();
811        assert_eq!(limiter.used_bulk_permits(), 1);
812
813        let limiter2 = limiter.clone();
814        let waiter = tokio::spawn(async move { limiter2.acquire_bulk().await });
815        tokio::time::sleep(Duration::from_secs(2)).await;
816
817        let result = waiter.await.unwrap();
818        assert!(matches!(result, Err(Error::AtCapacity)));
819        drop(permit);
820    }
821
822    #[tokio::test(start_paused = true)]
823    async fn queue_bounded_under_bulk_load() {
824        let limiter = ConcurrencyLimiter::new(3).with_queue(2);
825
826        let mut bulk_permits = Vec::new();
827        for _ in 0..3 {
828            let permit = limiter.acquire_bulk().await.unwrap();
829            bulk_permits.push(permit);
830        }
831
832        let limiter2 = limiter.clone();
833        let _w1 = tokio::spawn(async move { limiter2.acquire().await });
834        tokio::task::yield_now().await;
835        assert_eq!(limiter.queued_permits(), 1);
836
837        let limiter3 = limiter.clone();
838        let _w2 = tokio::spawn(async move { limiter3.acquire().await });
839        tokio::task::yield_now().await;
840        assert_eq!(limiter.queued_permits(), 2);
841
842        // Third waiter exceeds queue depth — rejected instantly.
843        let start = tokio::time::Instant::now();
844        let result = limiter.acquire().await;
845        assert!(matches!(result, Err(Error::AtCapacity)));
846        assert_eq!(start.elapsed(), Duration::ZERO);
847
848        drop(bulk_permits);
849    }
850}