1use 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
21const EMITTER_INTERVAL: Duration = Duration::from_secs(1);
23
24#[non_exhaustive]
29#[derive(Clone, Copy, Debug)]
30pub struct Stats {
31 pub in_use: u32,
33 pub queued: u32,
35 pub bulk_in_use: u32,
37}
38
39#[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 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 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 pub fn with_timeout(mut self, timeout: Duration) -> Self {
99 self.timeout = timeout;
100 self
101 }
102
103 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 pub async fn acquire(&self) -> Result<ConcurrencyPermit> {
125 if self.tasks_total == 0 {
126 return Err(Error::AtCapacity);
127 }
128
129 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 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 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 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 pub fn available_permits(&self) -> u32 {
213 u32::try_from(self.tasks.available_permits()).unwrap_or(self.tasks_total)
214 }
215
216 pub fn used_permits(&self) -> u32 {
218 self.tasks_total - self.available_permits()
219 }
220
221 pub fn total_permits(&self) -> u32 {
223 self.tasks_total
224 }
225
226 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 pub fn total_queue(&self) -> u32 {
234 self.queue_total
235 }
236
237 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 pub fn total_bulk(&self) -> u32 {
245 self.bulk_total
246 }
247
248 #[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 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 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
286pub 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
313pub 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 drop(p1);
485 assert!(futures::poll!(&mut wait).is_pending());
486
487 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 #[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 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 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 #[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 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 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 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}