1mod class;
62mod freelist;
63
64use super::{IoBufMut, page_size};
65use crate::{
66 iobuf::owner::PooledBuffer,
67 telemetry::metrics::{Counter, CounterFamily, EncodeLabelSet, GaugeFamily, Register, raw},
68};
69pub use class::BufferPoolThreadCache;
70use class::SizeClassHandle;
71pub(crate) use class::SizeClassLease;
72use commonware_utils::{NZU32, NZUsize};
73pub(super) use freelist::Freelist;
74use std::{
75 collections::BTreeMap,
76 num::{NonZeroU32, NonZeroUsize},
77 sync::atomic::{AtomicUsize, Ordering},
78};
79use thiserror::Error;
80
81cfg_if::cfg_if! {
82 if #[cfg(feature = "loom")] {
83 use loom::sync::Arc;
84 } else {
85 use std::sync::Arc;
86 }
87}
88
89#[derive(Error, Debug, Clone, Copy, PartialEq, Eq)]
91pub enum PoolError {
92 #[error("requested capacity exceeds maximum buffer size")]
94 Oversized,
95 #[error("pool exhausted for required size class")]
97 Exhausted,
98}
99
100#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub(crate) enum BufferPoolThreadCacheConfig {
103 Enabled(Option<NonZeroUsize>),
114 Disabled,
116}
117
118#[derive(Clone, Copy, Debug, PartialEq, Eq)]
120pub struct BufferPoolClassConfig {
121 pub size: NonZeroUsize,
123 pub max_buffers: NonZeroU32,
128}
129
130impl From<(NonZeroUsize, NonZeroU32)> for BufferPoolClassConfig {
131 fn from((size, max_buffers): (NonZeroUsize, NonZeroU32)) -> Self {
132 Self { size, max_buffers }
133 }
134}
135
136#[derive(Clone, Debug)]
149pub struct BufferPoolConfig {
150 pool_min_size: usize,
156 class_limits: BTreeMap<NonZeroUsize, NonZeroU32>,
161 prefill: bool,
168 alignment: NonZeroUsize,
170 parallelism: NonZeroUsize,
176 pub(crate) thread_cache_config: BufferPoolThreadCacheConfig,
183}
184
185impl BufferPoolConfig {
186 pub fn for_network() -> Self {
192 Self {
193 pool_min_size: 0,
194 class_limits: BTreeMap::new(),
195 prefill: false,
196 alignment: NZUsize!(1),
197 parallelism: NZUsize!(1),
198 thread_cache_config: BufferPoolThreadCacheConfig::Enabled(None),
199 }
200 .with_size_class_range(NZUsize!(1024), NZUsize!(128 * 1024), NZU32!(4096))
201 }
202
203 pub fn for_storage() -> Self {
206 Self {
207 pool_min_size: 0,
208 class_limits: BTreeMap::new(),
209 prefill: false,
210 alignment: NZUsize!(1),
212 parallelism: NZUsize!(1),
213 thread_cache_config: BufferPoolThreadCacheConfig::Enabled(None),
214 }
215 .with_size_class_range(
216 NZUsize!(page_size()),
217 NZUsize!(8 * 1024 * 1024),
218 NZU32!(64),
219 )
220 }
221
222 const fn validate_class_size(size: NonZeroUsize) {
228 assert!(
229 size.get().is_power_of_two(),
230 "class size must be a power of two"
231 );
232 assert!(
233 size.get() <= isize::MAX as usize,
234 "class size must not exceed isize::MAX"
235 );
236 }
237
238 pub const fn with_pool_min_size(mut self, pool_min_size: usize) -> Self {
240 self.pool_min_size = pool_min_size;
241 self
242 }
243
244 pub fn with_size_class_range(
255 self,
256 min: NonZeroUsize,
257 max: NonZeroUsize,
258 max_buffers: NonZeroU32,
259 ) -> Self {
260 Self::validate_class_size(min);
261 Self::validate_class_size(max);
262 assert!(max >= min, "max size must be >= min size");
263
264 self.with_size_classes(
265 (min.get().trailing_zeros()..=max.get().trailing_zeros())
266 .map(|exponent| (NZUsize!(1 << exponent), max_buffers)),
267 )
268 }
269
270 pub fn with_size_classes<I, C>(mut self, classes: I) -> Self
282 where
283 I: IntoIterator<Item = C>,
284 C: Into<BufferPoolClassConfig>,
285 {
286 let mut limits = BTreeMap::new();
287 for class in classes {
288 let class = class.into();
289 Self::validate_class_size(class.size);
290 assert!(
291 limits.insert(class.size, class.max_buffers).is_none(),
292 "duplicate class size {}",
293 class.size
294 );
295 }
296 assert!(
297 !limits.is_empty(),
298 "class layout must enable at least one class"
299 );
300 self.class_limits = limits;
301 self
302 }
303
304 pub fn with_size_class(mut self, size: NonZeroUsize, max_buffers: NonZeroU32) -> Self {
312 Self::validate_class_size(size);
313 self.class_limits.insert(size, max_buffers);
314 self
315 }
316
317 pub fn without_size_class(mut self, size: NonZeroUsize) -> Self {
329 Self::validate_class_size(size);
330 assert!(
331 self.class_limits.remove(&size).is_some(),
332 "cannot remove a class that is not enabled"
333 );
334 assert!(
335 !self.class_limits.is_empty(),
336 "cannot remove the final enabled class"
337 );
338 self
339 }
340
341 pub fn with_max_per_class(mut self, max_buffers: NonZeroU32) -> Self {
343 for limit in self.class_limits.values_mut() {
344 *limit = max_buffers;
345 }
346 self
347 }
348
349 pub fn with_bytes_per_class(mut self, bytes: NonZeroUsize) -> Self {
360 for (size, limit) in self.class_limits.iter_mut() {
361 let count = bytes.get() / size.get();
362 assert!(
363 count <= u32::MAX as usize,
364 "per-class byte weight derives a limit above u32::MAX"
365 );
366 *limit = NonZeroU32::new(count.max(1) as u32).expect("count is at least one");
367 }
368 self
369 }
370
371 pub const fn with_parallelism(mut self, parallelism: NonZeroUsize) -> Self {
379 self.parallelism = parallelism;
380 self
381 }
382
383 pub const fn with_max_thread_cache_capacity(mut self, capacity: NonZeroUsize) -> Self {
403 self.thread_cache_config = BufferPoolThreadCacheConfig::Enabled(Some(capacity));
404 self
405 }
406
407 pub const fn with_thread_cache_disabled(mut self) -> Self {
411 self.thread_cache_config = BufferPoolThreadCacheConfig::Disabled;
412 self
413 }
414
415 pub const fn with_prefill(mut self, prefill: bool) -> Self {
417 self.prefill = prefill;
418 self
419 }
420
421 pub const fn with_alignment(mut self, alignment: NonZeroUsize) -> Self {
423 self.alignment = alignment;
424 self
425 }
426
427 pub fn with_budget_bytes(mut self, budget: NonZeroUsize) -> Self {
449 let budget = budget.get() as u128;
450
451 let minimum: u128 = self
453 .class_limits
454 .keys()
455 .map(|size| size.get() as u128)
456 .sum();
457 assert!(
458 budget >= minimum,
459 "budget must cover at least one buffer from every enabled class"
460 );
461
462 const FRACTION_BITS: u32 = 64;
469
470 let scaled = |limit: NonZeroU32, scale: u128| -> u128 {
475 ((limit.get() as u128).saturating_mul(scale) >> FRACTION_BITS).max(1)
476 };
477 let evaluate = |scale: u128| -> (u128, bool) {
482 self.class_limits
483 .iter()
484 .fold((0u128, true), |(total, fits), (&size, &limit)| {
485 let count = scaled(limit, scale);
486 (
487 total.saturating_add(count.saturating_mul(size.get() as u128)),
488 fits && count <= u32::MAX as u128,
489 )
490 })
491 };
492 let feasible = |scale: u128| -> bool {
493 let (total, fits) = evaluate(scale);
494 total <= budget && fits
495 };
496
497 let mut lo: u128 = 0;
502 let mut hi: u128 = (u32::MAX as u128 + 1) << FRACTION_BITS;
503 assert!(feasible(lo), "scale zero must be feasible");
504 while hi - lo > 1 {
505 let mid = lo + (hi - lo) / 2;
506 if feasible(mid) {
507 lo = mid;
508 } else {
509 hi = mid;
510 }
511 }
512
513 let (next_total, _) = evaluate(lo + 1);
517 assert!(
518 next_total > budget,
519 "budget requires scaling a class limit above u32::MAX"
520 );
521
522 for limit in self.class_limits.values_mut() {
524 let count = u32::try_from(scaled(*limit, lo)).expect("feasible count fits u32");
525 *limit = NonZeroU32::new(count).expect("count is at least one");
526 }
527 self
528 }
529
530 pub fn size_classes(&self) -> impl ExactSizeIterator<Item = BufferPoolClassConfig> + '_ {
532 self.class_limits
533 .iter()
534 .map(|(&size, &max_buffers)| BufferPoolClassConfig { size, max_buffers })
535 }
536
537 pub fn class_for(&self, size: usize) -> Option<BufferPoolClassConfig> {
547 self.size_classes().find(|class| class.size.get() >= size)
548 }
549
550 pub const fn pool_min_size(&self) -> usize {
552 self.pool_min_size
553 }
554
555 pub const fn prefill(&self) -> bool {
557 self.prefill
558 }
559
560 pub const fn alignment(&self) -> NonZeroUsize {
562 self.alignment
563 }
564
565 pub const fn parallelism(&self) -> NonZeroUsize {
567 self.parallelism
568 }
569
570 pub fn min_size(&self) -> NonZeroUsize {
572 *self
573 .class_limits
574 .first_key_value()
575 .expect("class layout must enable at least one class")
576 .0
577 }
578
579 pub fn max_size(&self) -> NonZeroUsize {
581 *self
582 .class_limits
583 .last_key_value()
584 .expect("class layout must enable at least one class")
585 .0
586 }
587
588 pub fn max_tracked_bytes(&self) -> usize {
593 self.class_limits
594 .iter()
595 .map(|(size, limit)| size.get().saturating_mul(limit.get() as usize))
596 .fold(0usize, usize::saturating_add)
597 }
598
599 fn validate(&self) {
611 assert!(
612 self.alignment.is_power_of_two(),
613 "alignment must be a power of two"
614 );
615 let min_size = self.min_size();
616 assert!(
617 min_size >= self.alignment,
618 "smallest class ({}) must be >= alignment ({})",
619 min_size,
620 self.alignment
621 );
622 assert!(
623 self.pool_min_size <= min_size.get(),
624 "pool_min_size ({}) must be <= smallest class ({})",
625 self.pool_min_size,
626 min_size
627 );
628 }
629
630 fn resolve_thread_cache_capacity(&self, class_limit: NonZeroU32) -> usize {
637 let class_limit = class_limit.get() as usize;
638 match self.thread_cache_config {
639 BufferPoolThreadCacheConfig::Enabled(None) => {
640 let effective_threads = self.parallelism.get().min(class_limit);
641 class_limit / effective_threads.saturating_mul(2)
642 }
643 BufferPoolThreadCacheConfig::Enabled(Some(capacity)) => capacity.get().min(class_limit),
644 BufferPoolThreadCacheConfig::Disabled => 0,
645 }
646 }
647}
648
649#[derive(Clone, Debug, Hash, PartialEq, Eq, EncodeLabelSet)]
651struct SizeClassLabel {
652 size_class: u64,
653}
654
655struct PoolMetrics {
657 created: GaugeFamily<SizeClassLabel>,
659 exhausted_total: CounterFamily<SizeClassLabel>,
661 oversized_total: Counter,
663}
664
665impl PoolMetrics {
666 fn new(registry: &mut impl Register) -> Self {
667 Self {
668 created: registry.register(
669 "buffer_pool_created",
670 "Number of tracked buffers created for the pool",
671 raw::Family::default(),
672 ),
673 exhausted_total: registry.register(
676 "buffer_pool_exhausted",
677 "Total number of failed allocations due to pool exhaustion",
678 raw::Family::default(),
679 ),
680 oversized_total: registry.register(
681 "buffer_pool_oversized",
682 "Total number of allocation requests exceeding max buffer size",
683 raw::Counter::default(),
684 ),
685 }
686 }
687}
688
689struct Allocation {
691 buffer: PooledBuffer,
692 is_new: bool,
693}
694
695pub(crate) struct BufferPoolInner {
697 config: BufferPoolConfig,
698 classes: Vec<SizeClassHandle>,
706 min_size: usize,
709 max_size: usize,
711 metrics: PoolMetrics,
712}
713
714impl Drop for BufferPoolInner {
715 fn drop(&mut self) {
716 self.classes.dedup_by(|a, b| a.same_class(b));
726 assert_eq!(self.classes.len(), self.config.size_classes().len());
727 for class in &self.classes {
728 class.drain_global();
729 }
730 }
731}
732
733impl BufferPoolInner {
734 #[inline(always)]
746 fn try_alloc(&self, class_index: usize, zero_on_new: bool) -> Option<Allocation> {
747 let class = &self.classes[class_index];
748
749 if let Some(buffer) = BufferPoolThreadCache::pop(class) {
752 return Some(Allocation {
753 buffer,
754 is_new: false,
755 });
756 }
757
758 self.try_alloc_new(class, zero_on_new)
760 }
761
762 #[inline(never)]
768 fn try_alloc_new(&self, class: &SizeClassHandle, zeroed: bool) -> Option<Allocation> {
769 let label = SizeClassLabel {
770 size_class: class.size() as u64,
771 };
772 let Some(buffer) = class.try_create(zeroed) else {
773 self.metrics.exhausted_total.get_or_create(&label).inc();
774 return None;
775 };
776
777 self.metrics.created.get_or_create(&label).inc();
778 Some(Allocation {
779 buffer,
780 is_new: true,
781 })
782 }
783}
784
785#[derive(Clone)]
808pub struct BufferPool {
809 inner: Arc<BufferPoolInner>,
810}
811
812impl std::fmt::Debug for BufferPool {
813 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
814 f.debug_struct("BufferPool")
815 .field("config", &self.inner.config)
816 .field("num_classes", &self.inner.config.size_classes().len())
817 .finish()
818 }
819}
820
821static NEXT_SIZE_CLASS_ID: AtomicUsize = AtomicUsize::new(0);
837
838impl BufferPool {
839 pub(crate) fn new(config: BufferPoolConfig, registry: &mut impl Register) -> Self {
845 config.validate();
846 let metrics = PoolMetrics::new(registry);
847 let min_size = config.min_size().get();
848 let max_size = config.max_size().get();
849 let min_exponent = min_size.trailing_zeros() as usize;
850 let max_exponent = max_size.trailing_zeros() as usize;
851
852 let mut classes = Vec::with_capacity(max_exponent - min_exponent + 1);
858 for class_config in config.size_classes() {
859 let class_id = NEXT_SIZE_CLASS_ID.fetch_add(1, Ordering::Relaxed);
860 let handle = SizeClassHandle::new(
861 class_id,
862 class_config.size.get(),
863 config.alignment.get(),
864 class_config.max_buffers,
865 config.parallelism,
866 config.resolve_thread_cache_capacity(class_config.max_buffers),
867 config.prefill,
868 );
869
870 if config.prefill {
872 let label = SizeClassLabel {
873 size_class: class_config.size.get() as u64,
874 };
875 metrics
876 .created
877 .get_or_create(&label)
878 .set(class_config.max_buffers.get() as i64);
879 }
880
881 let index = class_config.size.get().trailing_zeros() as usize - min_exponent;
882 classes.resize(index + 1, handle);
883 }
884
885 Self {
886 inner: Arc::new(BufferPoolInner {
887 config,
888 classes,
889 min_size,
890 max_size,
891 metrics,
892 }),
893 }
894 }
895
896 #[inline(always)]
903 fn class_index(&self, size: usize) -> Option<usize> {
904 let min_size = self.inner.min_size;
905 let max_size = self.inner.max_size;
906 if size > max_size {
907 return None;
908 }
909 if size <= min_size {
910 return Some(0);
911 }
912
913 Some(
919 size.next_power_of_two()
920 .trailing_zeros()
921 .wrapping_sub(min_size.trailing_zeros()) as usize,
922 )
923 }
924
925 #[inline]
927 fn class_index_or_record_oversized(&self, capacity: usize) -> Option<usize> {
928 let class_index = self.class_index(capacity);
929 if class_index.is_none() {
930 self.inner.metrics.oversized_total.inc();
931 }
932 class_index
933 }
934
935 #[inline(always)]
957 pub fn try_alloc(&self, capacity: usize) -> Result<IoBufMut, PoolError> {
958 if capacity == 0 {
959 return Ok(IoBufMut::default());
960 }
961 if capacity < self.inner.config.pool_min_size {
962 return Ok(IoBufMut::with_alignment(
963 capacity,
964 self.inner.config.alignment,
965 ));
966 }
967
968 let class_index = self
969 .class_index_or_record_oversized(capacity)
970 .ok_or(PoolError::Oversized)?;
971
972 let buffer = self
973 .inner
974 .try_alloc(class_index, false)
975 .map(|allocation| {
976 unsafe { IoBufMut::from_pooled_parts(allocation.buffer) }
979 })
980 .ok_or(PoolError::Exhausted)?;
981 Ok(buffer)
982 }
983
984 #[inline]
1009 pub fn alloc(&self, capacity: usize) -> IoBufMut {
1010 self.try_alloc(capacity).unwrap_or_else(|_| {
1011 let size = capacity.max(1);
1012 IoBufMut::with_alignment(size, self.inner.config.alignment)
1013 })
1014 }
1015
1016 pub unsafe fn alloc_len(&self, len: usize) -> IoBufMut {
1025 let mut buf = self.alloc(len);
1026 unsafe { buf.set_len(len) };
1028 buf
1029 }
1030
1031 pub fn try_alloc_zeroed(&self, len: usize) -> Result<IoBufMut, PoolError> {
1053 if len == 0 {
1054 return Ok(IoBufMut::default());
1055 }
1056 if len < self.inner.config.pool_min_size {
1057 return Ok(IoBufMut::zeroed_with_alignment(
1058 len,
1059 self.inner.config.alignment,
1060 ));
1061 }
1062
1063 let class_index = self
1064 .class_index_or_record_oversized(len)
1065 .ok_or(PoolError::Oversized)?;
1066 let allocation = self
1067 .inner
1068 .try_alloc(class_index, true)
1069 .ok_or(PoolError::Exhausted)?;
1070 let mut buf = unsafe { IoBufMut::from_pooled_parts(allocation.buffer) };
1073 if allocation.is_new {
1074 unsafe { buf.set_len(len) };
1077 } else {
1078 unsafe {
1081 std::ptr::write_bytes(buf.as_mut_ptr(), 0, len);
1082 buf.set_len(len);
1083 }
1084 }
1085 Ok(buf)
1086 }
1087
1088 pub fn alloc_zeroed(&self, len: usize) -> IoBufMut {
1112 self.try_alloc_zeroed(len).unwrap_or_else(|_| {
1113 let size = len.max(1);
1115 let mut buf = IoBufMut::zeroed_with_alignment(size, self.inner.config.alignment);
1116 buf.truncate(len);
1117 buf
1118 })
1119 }
1120
1121 pub fn config(&self) -> &BufferPoolConfig {
1123 &self.inner.config
1124 }
1125}
1126
1127#[cfg(all(test, not(feature = "loom")))]
1128mod tests {
1129 use super::{
1130 class::tests::{
1131 get_global_created, get_global_len, get_global_num_stripes, get_local_len,
1132 get_thread_cache_capacity,
1133 },
1134 *,
1135 };
1136 use crate::{
1137 iobuf::{IoBuf, cache_line_size},
1138 telemetry::metrics::Registry,
1139 };
1140 use bytes::{Buf, BufMut};
1141 use commonware_utils::NZU32;
1142 use std::{
1143 sync::{Arc, mpsc},
1144 thread,
1145 };
1146
1147 fn test_pool(config: BufferPoolConfig) -> BufferPool {
1148 let mut registry = Registry::default();
1149 BufferPool::new(config, &mut registry)
1150 }
1151
1152 fn test_config(min_size: usize, max_size: usize, max_per_class: u32) -> BufferPoolConfig {
1154 BufferPoolConfig::for_network()
1155 .with_pool_min_size(0)
1156 .with_size_class_range(
1157 NZUsize!(min_size),
1158 NZUsize!(max_size),
1159 NZU32!(max_per_class),
1160 )
1161 .with_alignment(NZUsize!(page_size()))
1162 }
1163
1164 fn sparse_config(classes: impl IntoIterator<Item = (usize, u32)>) -> BufferPoolConfig {
1166 BufferPoolConfig::for_network()
1167 .with_pool_min_size(0)
1168 .with_size_classes(
1169 classes
1170 .into_iter()
1171 .map(|(size, max_buffers)| (NZUsize!(size), NZU32!(max_buffers))),
1172 )
1173 .with_alignment(NZUsize!(page_size()))
1174 }
1175
1176 fn classes_of(config: &BufferPoolConfig) -> Vec<(usize, u32)> {
1178 config
1179 .size_classes()
1180 .map(|class| (class.size.get(), class.max_buffers.get()))
1181 .collect()
1182 }
1183
1184 fn get_allocated(pool: &BufferPool, size: usize) -> usize {
1189 let class_index = pool.class_index(size).unwrap();
1190 let class = &pool.inner.classes[class_index];
1191 get_global_created(class) - get_global_len(class) - get_local_len(class)
1192 }
1193
1194 fn get_available(pool: &BufferPool, size: usize) -> i64 {
1196 let class_index = pool.class_index(size).unwrap();
1197 let class = &pool.inner.classes[class_index];
1198 (get_global_len(class) + get_local_len(class)) as i64
1199 }
1200
1201 #[test]
1202 fn test_page_size() {
1203 let size = page_size();
1204 assert!(size >= 4096);
1205 assert!(size.is_power_of_two());
1206 }
1207
1208 #[test]
1209 fn test_config_validation() {
1210 let page = page_size();
1211 let config = test_config(page, page * 4, 10);
1212 config.validate();
1213 }
1214
1215 #[test]
1216 fn test_explicit_thread_cache_capacity_clamps_to_class_limit() {
1217 let page = page_size();
1218 let config = test_config(page, page * 4, 10).with_max_thread_cache_capacity(NZUsize!(11));
1221 config.validate();
1222 let pool = test_pool(config);
1223 let class_index = pool.class_index(page).unwrap();
1224 assert_eq!(
1225 get_thread_cache_capacity(&pool.inner.classes[class_index]),
1226 10
1227 );
1228
1229 let config = BufferPoolConfig::for_network()
1232 .with_size_classes([(NZUsize!(1024), NZU32!(4)), (NZUsize!(4096), NZU32!(64))])
1233 .with_max_thread_cache_capacity(NZUsize!(16));
1234 let pool = test_pool(config);
1235 let small_index = pool.class_index(1024).unwrap();
1236 let large_index = pool.class_index(4096).unwrap();
1237 assert_eq!(
1238 get_thread_cache_capacity(&pool.inner.classes[small_index]),
1239 4
1240 );
1241 assert_eq!(
1242 get_thread_cache_capacity(&pool.inner.classes[large_index]),
1243 16
1244 );
1245 }
1246
1247 #[test]
1248 #[should_panic(expected = "class size must be a power of two")]
1249 fn test_config_invalid_min_size() {
1250 let _ = BufferPoolConfig::for_network().with_size_class_range(
1251 NZUsize!(3000),
1252 NZUsize!(8192),
1253 NZU32!(10),
1254 );
1255 }
1256
1257 #[test]
1258 #[should_panic(expected = "class size must be a power of two")]
1259 fn test_config_invalid_max_size() {
1260 let _ = BufferPoolConfig::for_network().with_size_class_range(
1261 NZUsize!(4096),
1262 NZUsize!(12000),
1263 NZU32!(10),
1264 );
1265 }
1266
1267 #[test]
1268 #[should_panic(expected = "max size must be >= min size")]
1269 fn test_config_range_rejects_max_below_min() {
1270 let _ = BufferPoolConfig::for_network().with_size_class_range(
1271 NZUsize!(8192),
1272 NZUsize!(1024),
1273 NZU32!(10),
1274 );
1275 }
1276
1277 #[test]
1278 #[should_panic(expected = "class size must not exceed isize::MAX")]
1279 fn test_config_rejects_class_size_above_isize_max() {
1280 let _ = BufferPoolConfig::for_network()
1281 .with_size_class(NZUsize!(1usize << (usize::BITS - 1)), NZU32!(1));
1282 }
1283
1284 #[test]
1285 #[should_panic(expected = "alignment must be a power of two")]
1286 fn test_config_invalid_alignment() {
1287 let page = page_size();
1288 let config = test_config(page, page, 10).with_alignment(NZUsize!(page - 1));
1289 config.validate();
1290 }
1291
1292 #[test]
1293 #[should_panic(expected = "must be >= alignment")]
1294 fn test_config_min_size_below_alignment() {
1295 let page = page_size();
1296 let config = test_config(page, page, 10).with_alignment(NZUsize!(page * 2));
1297 config.validate();
1298 }
1299
1300 #[test]
1301 #[should_panic(expected = "pool_min_size")]
1302 fn test_config_pool_min_size_above_min_size() {
1303 let page = page_size();
1304 let config = test_config(page, page, 10).with_pool_min_size(page + 1);
1305 config.validate();
1306 }
1307
1308 #[test]
1309 fn test_pool_class_index() {
1310 let page = page_size();
1311 let pool = test_pool(test_config(page, page * 8, 10));
1312
1313 assert_eq!(pool.inner.classes.len(), 4);
1315
1316 assert_eq!(pool.class_index(1), Some(0));
1317 assert_eq!(pool.class_index(page), Some(0));
1318 assert_eq!(pool.class_index(page + 1), Some(1));
1319 assert_eq!(pool.class_index(page * 2), Some(1));
1320 assert_eq!(pool.class_index(page * 4 + 1), Some(3));
1321 assert_eq!(pool.class_index(page * 8 - 1), Some(3));
1322 assert_eq!(pool.class_index(page * 8), Some(3));
1323 assert_eq!(pool.class_index(page * 8 + 1), None);
1324 }
1325
1326 #[test]
1327 fn test_size_classes_replacement_normalizes_and_iterates() {
1328 let config = BufferPoolConfig::for_network().with_size_classes([
1330 (NZUsize!(1 << 20), NZU32!(16)),
1331 (NZUsize!(4096), NZU32!(1024)),
1332 (NZUsize!(65536), NZU32!(256)),
1333 ]);
1334 assert_eq!(
1335 classes_of(&config),
1336 vec![(4096, 1024), (65536, 256), (1 << 20, 16)]
1337 );
1338 assert_eq!(config.size_classes().len(), 3);
1339 assert_eq!(config.min_size().get(), 4096);
1340 assert_eq!(config.max_size().get(), 1 << 20);
1341 assert_eq!(
1342 config.max_tracked_bytes(),
1343 4096 * 1024 + 65536 * 256 + (1 << 20) * 16
1344 );
1345
1346 let explicit = BufferPoolConfig::for_network().with_size_classes([BufferPoolClassConfig {
1348 size: NZUsize!(512),
1349 max_buffers: NZU32!(2),
1350 }]);
1351 assert_eq!(classes_of(&explicit), vec![(512, 2)]);
1352 }
1353
1354 #[test]
1355 #[should_panic(expected = "class layout must enable at least one class")]
1356 fn test_size_classes_rejects_empty_input() {
1357 let _ = BufferPoolConfig::for_network()
1358 .with_size_classes(std::iter::empty::<BufferPoolClassConfig>());
1359 }
1360
1361 #[test]
1362 #[should_panic(expected = "duplicate class size 4096")]
1363 fn test_size_classes_rejects_duplicates() {
1364 let _ = BufferPoolConfig::for_network()
1365 .with_size_classes([(NZUsize!(4096), NZU32!(1)), (NZUsize!(4096), NZU32!(2))]);
1366 }
1367
1368 #[test]
1369 #[should_panic(expected = "class size must be a power of two")]
1370 fn test_size_classes_rejects_non_power_of_two() {
1371 let _ = BufferPoolConfig::for_network().with_size_classes([(NZUsize!(3000), NZU32!(1))]);
1372 }
1373
1374 #[test]
1375 fn test_size_class_upsert_and_removal() {
1376 let base = BufferPoolConfig::for_network().with_size_class_range(
1377 NZUsize!(1024),
1378 NZUsize!(8192),
1379 NZU32!(8),
1380 );
1381
1382 let tuned = base.clone().with_size_class(NZUsize!(2048), NZU32!(64));
1384 assert_eq!(
1385 classes_of(&tuned),
1386 vec![(1024, 8), (2048, 64), (4096, 8), (8192, 8)]
1387 );
1388
1389 let extended = base.clone().with_size_class(NZUsize!(32768), NZU32!(2));
1391 assert_eq!(extended.max_size().get(), 32768);
1392 assert_eq!(classes_of(&extended).len(), 5);
1393
1394 let sparse = base.clone().without_size_class(NZUsize!(2048));
1396 assert_eq!(classes_of(&sparse), vec![(1024, 8), (4096, 8), (8192, 8)]);
1397
1398 let narrowed = base.without_size_class(NZUsize!(1024));
1400 assert_eq!(narrowed.min_size().get(), 2048);
1401 let narrowed = narrowed.without_size_class(NZUsize!(8192));
1402 assert_eq!(narrowed.max_size().get(), 4096);
1403
1404 let uniform = sparse.with_max_per_class(NZU32!(3));
1406 assert_eq!(classes_of(&uniform), vec![(1024, 3), (4096, 3), (8192, 3)]);
1407 }
1408
1409 #[test]
1410 #[should_panic(expected = "cannot remove a class that is not enabled")]
1411 fn test_without_size_class_rejects_absent_class() {
1412 let _ = BufferPoolConfig::for_network().without_size_class(NZUsize!(1 << 30));
1413 }
1414
1415 #[test]
1416 #[should_panic(expected = "cannot remove the final enabled class")]
1417 fn test_without_size_class_rejects_final_class() {
1418 let _ = BufferPoolConfig::for_network()
1419 .with_size_classes([(NZUsize!(4096), NZU32!(1))])
1420 .without_size_class(NZUsize!(4096));
1421 }
1422
1423 #[test]
1424 fn test_bytes_per_class_gives_equal_byte_weight() {
1425 let config = BufferPoolConfig::for_network()
1426 .with_size_classes([
1427 (NZUsize!(1024), NZU32!(1)),
1428 (NZUsize!(4096), NZU32!(1)),
1429 (NZUsize!(1 << 20), NZU32!(1)),
1430 ])
1431 .with_bytes_per_class(NZUsize!(64 * 1024));
1432 assert_eq!(
1435 classes_of(&config),
1436 vec![(1024, 64), (4096, 16), (1 << 20, 1)]
1437 );
1438
1439 let exact = BufferPoolConfig::for_network()
1441 .with_size_classes([(NZUsize!(4096), NZU32!(7))])
1442 .with_bytes_per_class(NZUsize!(4096));
1443 assert_eq!(classes_of(&exact), vec![(4096, 1)]);
1444 }
1445
1446 #[test]
1447 fn test_class_for_routes_to_smallest_fitting_class() {
1448 let config = BufferPoolConfig::for_network()
1449 .with_size_classes([(NZUsize!(4096), NZU32!(4)), (NZUsize!(32768), NZU32!(2))]);
1450
1451 assert_eq!(config.class_for(0).unwrap().size.get(), 4096);
1453 assert_eq!(config.class_for(4096).unwrap().size.get(), 4096);
1454
1455 assert_eq!(config.class_for(4097).unwrap().size.get(), 32768);
1458 assert_eq!(config.class_for(16384).unwrap().size.get(), 32768);
1459 assert_eq!(config.class_for(32768).unwrap().size.get(), 32768);
1460
1461 assert_eq!(config.class_for(32769), None);
1463 }
1464
1465 #[test]
1466 fn test_sparse_routing_allocates_next_enabled_class() {
1467 let page = page_size();
1470 let pool = test_pool(sparse_config([(page, 4), (page * 8, 4)]));
1471
1472 let buf = pool.try_alloc(1).unwrap();
1474 assert_eq!(buf.capacity(), page);
1475
1476 let buf = pool.try_alloc(page).unwrap();
1478 assert_eq!(buf.capacity(), page);
1479
1480 let buf = pool.try_alloc(page + 1).unwrap();
1482 assert_eq!(buf.capacity(), page * 8);
1483
1484 let buf = pool.try_alloc(page * 4).unwrap();
1486 assert_eq!(buf.capacity(), page * 8);
1487
1488 let buf = pool.try_alloc(page * 8).unwrap();
1490 assert_eq!(buf.capacity(), page * 8);
1491
1492 assert_eq!(
1494 pool.try_alloc(page * 8 + 1).unwrap_err(),
1495 PoolError::Oversized
1496 );
1497 }
1498
1499 #[test]
1500 fn test_sparse_routing_exhaustion_does_not_cascade() {
1501 let page = page_size();
1502 let pool = test_pool(sparse_config([(page, 1), (page * 8, 1)]));
1503
1504 let _small = pool.try_alloc(page).unwrap();
1507 assert_eq!(pool.try_alloc(page).unwrap_err(), PoolError::Exhausted);
1508
1509 let _large = pool.try_alloc(page * 8).unwrap();
1511
1512 let fallback = pool.alloc(page);
1516 assert!(!fallback.is_pooled());
1517 assert_eq!(fallback.capacity(), page);
1518 let small_fallback = pool.alloc(100);
1519 assert!(!small_fallback.is_pooled());
1520 assert!((100..108).contains(&small_fallback.capacity()));
1521 }
1522
1523 #[test]
1524 fn test_sparse_metrics_use_enabled_class_labels() {
1525 let page = page_size();
1528 let mut registry = Registry::default();
1529 let pool = BufferPool::new(sparse_config([(page, 1), (page * 8, 1)]), &mut registry);
1530
1531 let _held = pool.try_alloc(page * 2).unwrap();
1533 assert!(pool.try_alloc(page * 2).is_err());
1534 assert!(pool.try_alloc(page * 16).is_err());
1536
1537 let encoded = registry.encode();
1538 assert!(
1541 encoded.contains(&format!(
1542 "buffer_pool_created{{size_class=\"{}\"}} 1",
1543 page * 8
1544 )),
1545 "metrics output: {encoded}"
1546 );
1547 assert!(
1548 encoded.contains(&format!(
1549 "buffer_pool_exhausted_total{{size_class=\"{}\"}} 1",
1550 page * 8
1551 )),
1552 "metrics output: {encoded}"
1553 );
1554 assert!(
1555 !encoded.contains(&format!("size_class=\"{}\"", page * 2)),
1556 "metrics output: {encoded}"
1557 );
1558 assert!(
1559 encoded.contains("buffer_pool_oversized_total 1"),
1560 "metrics output: {encoded}"
1561 );
1562 }
1563
1564 #[test]
1565 fn test_sparse_gap_allocations_share_one_class() {
1566 let page = page_size();
1569 let pool = test_pool(sparse_config([(page, 4), (page * 8, 4)]));
1570
1571 let direct = pool.class_index(page * 8).unwrap();
1573 for size in [page + 1, page * 2, page * 4, page * 8] {
1574 let index = pool.class_index(size).unwrap();
1575 assert!(
1576 pool.inner.classes[index].same_class(&pool.inner.classes[direct]),
1577 "size {size} must alias the largest class"
1578 );
1579 }
1580
1581 let mut via_gap = pool.try_alloc(page * 2).unwrap();
1584 let ptr = via_gap.as_mut_ptr();
1585 drop(via_gap);
1586 let mut direct_reuse = pool.try_alloc(page * 8).unwrap();
1587 assert_eq!(direct_reuse.as_mut_ptr(), ptr);
1588 }
1589
1590 #[test]
1591 fn test_sparse_pool_drop_drains_each_unique_class_once() {
1592 let page = page_size();
1596 let pool =
1597 test_pool(sparse_config([(page, 2), (page * 16, 2)]).with_thread_cache_disabled());
1598
1599 let class_index = pool.class_index(page * 16).unwrap();
1600 let class = pool.inner.classes[class_index].clone();
1603
1604 let buf = pool.try_alloc(page * 2).unwrap();
1606 drop(buf);
1607 assert_eq!(get_global_len(&class), 1);
1608
1609 drop(pool);
1610 assert_eq!(get_global_len(&class), 0);
1611 assert_eq!(get_global_created(&class), 1);
1612 }
1613
1614 #[test]
1615 fn test_sparse_pool_debug_reports_unique_classes() {
1616 let page = page_size();
1617 let pool = test_pool(sparse_config([(page, 2), (page * 16, 2)]));
1618 assert_eq!(pool.inner.classes.len(), 5);
1620 let debug = format!("{pool:?}");
1621 assert!(debug.contains("num_classes: 2"), "debug output: {debug}");
1622 }
1623
1624 #[test]
1625 fn test_sparse_prefill_creates_per_class_limits() {
1626 let page = page_size();
1627 let pool = test_pool(sparse_config([(page, 3), (page * 4, 1)]).with_prefill(true));
1628
1629 let small = &pool.inner.classes[pool.class_index(page).unwrap()];
1631 let large = &pool.inner.classes[pool.class_index(page * 4).unwrap()];
1632 assert_eq!(get_global_created(small), 3);
1633 assert_eq!(get_global_len(small), 3);
1634 assert_eq!(get_global_created(large), 1);
1635 assert_eq!(get_global_len(large), 1);
1636
1637 let a = pool.try_alloc(page).unwrap();
1639 let b = pool.try_alloc(page).unwrap();
1640 let c = pool.try_alloc(page).unwrap();
1641 assert!(pool.try_alloc(page).is_err());
1642 drop((a, b, c));
1643 let _gap = pool.try_alloc(page * 2).unwrap();
1644 assert!(pool.try_alloc(page * 4).is_err());
1645 }
1646
1647 #[test]
1648 fn test_pool_alloc_and_return() {
1649 let page = page_size();
1650 let pool = test_pool(test_config(page, page * 4, 2));
1651
1652 let buf = pool.try_alloc(page).unwrap();
1654 assert!(buf.capacity() >= page);
1655 assert_eq!(buf.len(), 0);
1656
1657 drop(buf);
1659
1660 let buf2 = pool.try_alloc(page).unwrap();
1662 assert!(buf2.capacity() >= page);
1663 assert_eq!(buf2.len(), 0);
1664 }
1665
1666 #[test]
1667 fn test_alloc_len_sets_len() {
1668 let page = page_size();
1669 let pool = test_pool(test_config(page, page * 4, 2));
1670
1671 let mut buf = unsafe { pool.alloc_len(100) };
1673 assert_eq!(buf.len(), 100);
1674 buf.as_mut().fill(0xAB);
1675 let frozen = buf.freeze();
1676 assert_eq!(frozen.as_ref(), &[0xAB; 100]);
1677 }
1678
1679 #[test]
1680 fn test_alloc_zeroed_sets_len_and_zeros() {
1681 let page = page_size();
1682 let pool = test_pool(test_config(page, page * 4, 2));
1683
1684 let buf = pool.alloc_zeroed(100);
1685 assert_eq!(buf.len(), 100);
1686 assert!(buf.as_ref().iter().all(|&b| b == 0));
1687 }
1688
1689 #[test]
1690 fn test_try_alloc_zeroed_sets_len_and_zeros() {
1691 let page = page_size();
1692 let pool = test_pool(test_config(page, page * 4, 2));
1693
1694 let buf = pool.try_alloc_zeroed(page).unwrap();
1695 assert!(buf.is_pooled());
1696 assert_eq!(buf.len(), page);
1697 assert!(buf.as_ref().iter().all(|&b| b == 0));
1698 }
1699
1700 #[test]
1701 fn test_alloc_zeroed_fallback_uses_untracked_zeroed_buffer() {
1702 let page = page_size();
1703 let pool = test_pool(test_config(page, page, 1));
1704
1705 let _pooled = pool.try_alloc(page).unwrap();
1707
1708 let buf = pool.alloc_zeroed(100);
1709 assert!(!buf.is_pooled());
1710 assert_eq!(buf.len(), 100);
1711 assert!(buf.as_ref().iter().all(|&b| b == 0));
1712 }
1713
1714 #[test]
1715 fn test_alloc_zeroed_reuses_dirty_pooled_buffer() {
1716 let page = page_size();
1717 let pool = test_pool(test_config(page, page, 1));
1718
1719 let mut first = pool.alloc_zeroed(page);
1720 assert!(first.is_pooled());
1721 assert!(first.as_ref().iter().all(|&b| b == 0));
1722
1723 first.as_mut().fill(0xAB);
1725 drop(first);
1726
1727 let second = pool.alloc_zeroed(page);
1728 assert!(second.is_pooled());
1729 assert_eq!(second.len(), page);
1730 assert!(second.as_ref().iter().all(|&b| b == 0));
1731 }
1732
1733 #[test]
1734 fn test_requests_smaller_than_pool_min_size_bypass_pool() {
1735 let pool = test_pool(
1736 BufferPoolConfig::for_network()
1737 .with_pool_min_size(512)
1738 .with_size_class_range(NZUsize!(512), NZUsize!(1024), NZU32!(2))
1739 .with_alignment(NZUsize!(128)),
1740 );
1741
1742 let buf = pool.try_alloc(200).unwrap();
1743 assert!(!buf.is_pooled());
1744 assert_eq!(buf.capacity(), 200);
1745
1746 let zeroed = pool.try_alloc_zeroed(200).unwrap();
1747 assert!(!zeroed.is_pooled());
1748 assert_eq!(zeroed.len(), 200);
1749 assert!(zeroed.as_ref().iter().all(|&b| b == 0));
1750
1751 let pooled = pool.try_alloc(512).unwrap();
1752 assert!(pooled.is_pooled());
1753 assert_eq!(pooled.capacity(), 512);
1754 }
1755
1756 #[test]
1757 fn test_zero_capacity_requests_bypass_pool() {
1758 let page = page_size();
1761 let pool = test_pool(test_config(page, page, 1));
1762
1763 let empty = pool.try_alloc(0).unwrap();
1764 assert!(!empty.is_pooled());
1765 assert_eq!(empty.capacity(), 0);
1766
1767 let zeroed = pool.try_alloc_zeroed(0).unwrap();
1768 assert!(!zeroed.is_pooled());
1769 assert_eq!(zeroed.len(), 0);
1770 assert_eq!(zeroed.capacity(), 0);
1771
1772 assert_eq!(pool.alloc(0).capacity(), 0);
1773 assert_eq!(pool.alloc_zeroed(0).len(), 0);
1774
1775 let real = pool.try_alloc(page).unwrap();
1776 assert!(real.is_pooled());
1777 assert_eq!(real.capacity(), page);
1778 }
1779
1780 #[test]
1781 fn test_pool_size_classes() {
1782 let page = page_size();
1783 let pool = test_pool(test_config(page, page * 4, 10));
1784
1785 let buf1 = pool.try_alloc(page).unwrap();
1787 assert_eq!(buf1.capacity(), page);
1788
1789 let buf2 = pool.try_alloc(page + 1).unwrap();
1791 assert_eq!(buf2.capacity(), page * 2);
1792
1793 let buf3 = pool.try_alloc(page * 3).unwrap();
1794 assert_eq!(buf3.capacity(), page * 4);
1795 }
1796
1797 #[test]
1798 fn test_prefill() {
1799 let page = NZUsize!(page_size());
1800 let pool = test_pool(
1801 BufferPoolConfig::for_network()
1802 .with_pool_min_size(0)
1803 .with_size_class_range(page, page, NZU32!(5))
1804 .with_alignment(page)
1805 .with_prefill(true),
1806 );
1807
1808 let mut bufs = Vec::new();
1810 for _ in 0..5 {
1811 bufs.push(pool.try_alloc(page.get()).expect("alloc should succeed"));
1812 }
1813
1814 assert!(pool.try_alloc(page.get()).is_err());
1816 }
1817
1818 #[test]
1819 fn test_config_for_network() {
1820 let config = BufferPoolConfig::for_network();
1821 config.validate();
1822 assert_eq!(config.pool_min_size, 0);
1823 assert_eq!(config.min_size().get(), 1024);
1824 assert_eq!(config.max_size().get(), 128 * 1024);
1825 let expected: Vec<(usize, u32)> = (10..=17).map(|e| (1usize << e, 4096)).collect();
1826 assert_eq!(classes_of(&config), expected);
1827 assert_eq!(config.parallelism, NZUsize!(1));
1828 assert_eq!(
1829 config.thread_cache_config,
1830 BufferPoolThreadCacheConfig::Enabled(None)
1831 );
1832 assert!(!config.prefill);
1833 assert_eq!(config.alignment.get(), 1);
1834 }
1835
1836 #[test]
1837 fn test_config_for_storage() {
1838 let config = BufferPoolConfig::for_storage();
1839 config.validate();
1840 assert_eq!(config.pool_min_size, 0);
1841 assert_eq!(config.min_size().get(), page_size());
1842 assert_eq!(config.max_size().get(), 8 * 1024 * 1024);
1843 let min_exponent = page_size().trailing_zeros();
1844 let expected: Vec<(usize, u32)> = (min_exponent..=23).map(|e| (1usize << e, 64)).collect();
1845 assert_eq!(classes_of(&config), expected);
1846 assert_eq!(config.parallelism, NZUsize!(1));
1847 assert_eq!(
1848 config.thread_cache_config,
1849 BufferPoolThreadCacheConfig::Enabled(None)
1850 );
1851 assert!(!config.prefill);
1852 assert_eq!(config.alignment.get(), 1);
1853 }
1854
1855 #[test]
1856 fn test_storage_config_supports_default_allocations() {
1857 let pool = test_pool(BufferPoolConfig::for_storage());
1859
1860 let buf = pool.try_alloc(8 * 1024 * 1024).unwrap();
1861 assert_eq!(buf.capacity(), 8 * 1024 * 1024);
1862 }
1863
1864 #[test]
1865 fn test_config_builders() {
1866 let page = NZUsize!(page_size());
1867 let config = BufferPoolConfig::for_storage()
1868 .with_pool_min_size(1024)
1869 .with_parallelism(NZUsize!(4))
1870 .with_max_thread_cache_capacity(NZUsize!(8))
1871 .with_prefill(true)
1872 .with_size_class_range(page, NZUsize!(128 * 1024), NZU32!(64));
1873
1874 config.validate();
1875 assert_eq!(config.pool_min_size, 1024);
1876 assert_eq!(config.min_size(), page);
1877 assert_eq!(config.max_size().get(), 128 * 1024);
1878 assert!(
1879 config
1880 .size_classes()
1881 .all(|class| class.max_buffers.get() == 64)
1882 );
1883 assert_eq!(config.parallelism, NZUsize!(4));
1884 assert_eq!(
1885 config.thread_cache_config,
1886 BufferPoolThreadCacheConfig::Enabled(Some(NZUsize!(8)))
1887 );
1888 assert!(config.prefill);
1889 assert_eq!(config.alignment.get(), 1);
1890
1891 let aligned = BufferPoolConfig::for_network()
1894 .with_pool_min_size(256)
1895 .with_parallelism(NZUsize!(4))
1896 .with_alignment(NZUsize!(256))
1897 .with_size_class_range(NZUsize!(256), NZUsize!(128 * 1024), NZU32!(4096));
1898 aligned.validate();
1899 assert_eq!(aligned.parallelism, NZUsize!(4));
1900 assert_eq!(
1901 aligned.thread_cache_config,
1902 BufferPoolThreadCacheConfig::Enabled(None)
1903 );
1904 assert_eq!(aligned.alignment.get(), 256);
1905 assert_eq!(aligned.min_size().get(), 256);
1906 }
1907
1908 #[test]
1909 fn test_parallelism_policy_resolves_thread_cache_capacity() {
1910 let page = page_size();
1911
1912 let pool = test_pool(test_config(page, page, 64).with_parallelism(NZUsize!(8)));
1914 let class_index = pool.class_index(page).unwrap();
1915 assert_eq!(
1916 get_thread_cache_capacity(&pool.inner.classes[class_index]),
1917 4
1918 );
1919
1920 let pool = test_pool(test_config(page, page, 4096).with_parallelism(NZUsize!(8)));
1922 let class_index = pool.class_index(page).unwrap();
1923 assert_eq!(
1924 get_thread_cache_capacity(&pool.inner.classes[class_index]),
1925 256
1926 );
1927 }
1928
1929 #[test]
1930 fn test_auto_thread_cache_disables_when_parallelism_exceeds_budget() {
1931 let page = page_size();
1932
1933 let pool = test_pool(test_config(page, page, 2).with_parallelism(NZUsize!(8)));
1938 let class_index = pool.class_index(page).unwrap();
1939 let class = &pool.inner.classes[class_index];
1940 assert_eq!(get_thread_cache_capacity(class), 0);
1941
1942 let first = pool.try_alloc(page).expect("first tracked allocation");
1945 let second = pool.try_alloc(page).expect("second tracked allocation");
1946
1947 let pool_for_thread = pool.clone();
1948 let (returned_tx, returned_rx) = mpsc::channel();
1949 let (release_tx, release_rx) = mpsc::channel();
1950 let handle = thread::spawn(move || {
1951 drop(first);
1955 drop(second);
1956 returned_tx.send(()).expect("signal returned buffers");
1957 release_rx.recv().expect("release worker");
1958 drop(pool_for_thread);
1959 });
1960
1961 returned_rx.recv().expect("wait for returned buffers");
1962
1963 let _first = pool.try_alloc(page).expect("first global reuse");
1968 let _second = pool.try_alloc(page).expect("second global reuse");
1969
1970 release_tx.send(()).expect("release worker");
1971 handle.join().expect("worker should not panic");
1972 }
1973
1974 #[test]
1975 fn test_parallelism_policy_resolves_freelist_stripes() {
1976 let page = page_size();
1977 let pool = test_pool(test_config(page, page, 64).with_parallelism(NZUsize!(16)));
1978
1979 let class_index = pool.class_index(page).unwrap();
1980 assert_eq!(get_global_num_stripes(&pool.inner.classes[class_index]), 16);
1981
1982 let pool = test_pool(test_config(page, page, 12).with_parallelism(NZUsize!(9)));
1985
1986 let class_index = pool.class_index(page).unwrap();
1987 assert_eq!(get_global_num_stripes(&pool.inner.classes[class_index]), 8);
1988
1989 let pool = test_pool(
1991 test_config(page, page, 64)
1992 .with_parallelism(NZUsize!(16))
1993 .with_thread_cache_disabled(),
1994 );
1995
1996 let class_index = pool.class_index(page).unwrap();
1997 assert_eq!(get_global_num_stripes(&pool.inner.classes[class_index]), 16);
1998 }
1999
2000 #[test]
2001 fn test_fixed_thread_cache_capacity_overrides_auto_capacity() {
2002 let page = page_size();
2003 let pool = test_pool(
2004 test_config(page, page, 64)
2005 .with_parallelism(NZUsize!(8))
2006 .with_max_thread_cache_capacity(NZUsize!(7)),
2007 );
2008 let class_index = pool.class_index(page).unwrap();
2009
2010 assert_eq!(
2012 get_thread_cache_capacity(&pool.inner.classes[class_index]),
2013 7
2014 );
2015 assert_eq!(get_global_num_stripes(&pool.inner.classes[class_index]), 8);
2016 }
2017
2018 #[test]
2019 fn test_disabled_thread_cache_does_not_retain_buffers_locally() {
2020 let page = page_size();
2021 let pool = test_pool(test_config(page, page, 2).with_thread_cache_disabled());
2022 let class_index = pool.class_index(page).unwrap();
2023 let class = &pool.inner.classes[class_index];
2024
2025 let tracked = pool.try_alloc(page).expect("tracked allocation");
2026 drop(tracked);
2027
2028 assert_eq!(get_thread_cache_capacity(class), 0);
2031 assert_eq!(get_local_len(class), 0);
2032 assert_eq!(get_global_len(class), 1);
2033 }
2034
2035 #[test]
2036 fn test_config_with_budget_bytes() {
2037 let base = BufferPoolConfig::for_network().with_size_class_range(
2040 NZUsize!(4),
2041 NZUsize!(16),
2042 NZU32!(1),
2043 );
2044 let config = base.clone().with_budget_bytes(NZUsize!(280));
2045 assert_eq!(classes_of(&config), vec![(4, 10), (8, 10), (16, 10)]);
2046 assert_eq!(config.max_tracked_bytes(), 280);
2047
2048 let config = base.clone().with_budget_bytes(NZUsize!(279));
2050 assert_eq!(classes_of(&config), vec![(4, 9), (8, 9), (16, 9)]);
2051
2052 let config = base.clone().with_budget_bytes(NZUsize!(28));
2054 assert_eq!(classes_of(&config), vec![(4, 1), (8, 1), (16, 1)]);
2055
2056 let shaped_base = BufferPoolConfig::for_network()
2060 .with_size_classes([(NZUsize!(4), NZU32!(4)), (NZUsize!(16), NZU32!(1))]);
2061 let shaped = shaped_base.clone().with_budget_bytes(NZUsize!(96));
2062 assert_eq!(classes_of(&shaped), vec![(4, 12), (16, 3)]);
2063
2064 let shrunk = shaped_base.with_budget_bytes(NZUsize!(20));
2066 assert_eq!(classes_of(&shrunk), vec![(4, 1), (16, 1)]);
2067
2068 let uneven = base.with_budget_bytes(NZUsize!(30));
2071 assert_eq!(classes_of(&uneven), vec![(4, 1), (8, 1), (16, 1)]);
2072 }
2073
2074 #[test]
2075 fn test_config_with_budget_bytes_is_one_shot() {
2076 let config = BufferPoolConfig::for_network()
2079 .with_size_class_range(NZUsize!(4), NZUsize!(16), NZU32!(1))
2080 .with_budget_bytes(NZUsize!(280));
2081 assert_eq!(config.max_tracked_bytes(), 280);
2082
2083 let overridden = config.clone().with_max_per_class(NZU32!(100));
2084 assert_eq!(overridden.max_tracked_bytes(), 2800);
2085
2086 let upserted = config.with_size_class(NZUsize!(32), NZU32!(100));
2087 assert_eq!(upserted.max_tracked_bytes(), 280 + 32 * 100);
2088
2089 let base = BufferPoolConfig::for_network()
2093 .with_size_classes([(NZUsize!(1), NZU32!(3)), (NZUsize!(8), NZU32!(2))]);
2094 let once = base.with_budget_bytes(NZUsize!(21));
2095 assert_eq!(classes_of(&once), vec![(1, 4), (8, 2)]);
2096 let twice = once.with_budget_bytes(NZUsize!(21));
2097 assert_eq!(classes_of(&twice), vec![(1, 5), (8, 2)]);
2098 }
2099
2100 #[test]
2101 #[should_panic(expected = "budget must cover at least one buffer from every enabled class")]
2102 fn test_config_with_budget_bytes_below_minimum() {
2103 let _ = BufferPoolConfig::for_network()
2104 .with_size_class_range(NZUsize!(4), NZUsize!(16), NZU32!(1))
2105 .with_budget_bytes(NZUsize!(27));
2106 }
2107
2108 #[test]
2109 #[should_panic(expected = "budget requires scaling a class limit above u32::MAX")]
2110 fn test_config_with_budget_bytes_above_u32() {
2111 let _ = BufferPoolConfig::for_network()
2114 .with_size_classes([(NZUsize!(1), NZU32!(1))])
2115 .with_budget_bytes(NZUsize!(u32::MAX as usize + 2));
2116 }
2117
2118 #[test]
2119 fn test_config_with_budget_bytes_near_u32_breakpoints() {
2120 cfg_if::cfg_if! {
2125 if #[cfg(miri)] {
2126 let budget = 10_000usize;
2127 } else {
2128 let budget = 1_000_000usize;
2129 }
2130 }
2131 let a = u32::MAX;
2132 let b = u32::MAX - 1;
2133 let config = BufferPoolConfig::for_network()
2134 .with_size_classes([
2135 (NZUsize!(1), NonZeroU32::new(a).unwrap()),
2136 (NZUsize!(2), NonZeroU32::new(b).unwrap()),
2137 ])
2138 .with_budget_bytes(NonZeroUsize::new(budget).unwrap());
2139 let expected = brute_force_budget(&[(1, a), (2, b)], budget as u128);
2141 assert_eq!(
2142 classes_of(&config)
2143 .into_iter()
2144 .map(|(_, limit)| limit)
2145 .collect::<Vec<_>>(),
2146 expected
2147 );
2148 }
2149
2150 fn brute_force_budget(shape: &[(usize, u32)], budget: u128) -> Vec<u32> {
2155 let mut best: Option<Vec<u32>> = None;
2158 let mut best_total = 0u128;
2159 let mut candidates: Vec<(u128, u128)> = vec![(0, 1)];
2160 for &(size, limit) in shape {
2161 let max_count = (budget / size as u128).min(u32::MAX as u128);
2162 for k in 1..=max_count {
2163 candidates.push((k, limit as u128));
2164 }
2165 }
2166 for (k, c) in candidates {
2167 let counts: Vec<u128> = shape
2169 .iter()
2170 .map(|&(_, limit)| ((limit as u128 * k) / c).max(1))
2171 .collect();
2172 if counts.iter().any(|&count| count > u32::MAX as u128) {
2173 continue;
2174 }
2175 let total: u128 = counts
2176 .iter()
2177 .zip(shape.iter())
2178 .map(|(&count, &(size, _))| count * size as u128)
2179 .sum();
2180 if total <= budget && total >= best_total {
2183 best_total = total;
2184 best = Some(counts.iter().map(|&count| count as u32).collect());
2185 }
2186 }
2187 best.expect("budget covers one buffer per class")
2188 }
2189
2190 #[test]
2191 fn test_pool_error_display() {
2192 assert_eq!(
2193 PoolError::Oversized.to_string(),
2194 "requested capacity exceeds maximum buffer size"
2195 );
2196 assert_eq!(
2197 PoolError::Exhausted.to_string(),
2198 "pool exhausted for required size class"
2199 );
2200 }
2201
2202 #[test]
2203 fn test_pool_debug_and_config_accessor() {
2204 let page = page_size();
2206 let pool = test_pool(test_config(page, page, 2));
2207
2208 let debug = format!("{pool:?}");
2209 assert!(debug.contains("BufferPool"));
2210 assert!(debug.contains("num_classes"));
2211 assert_eq!(pool.config().min_size().get(), page);
2212 }
2213
2214 #[test]
2215 fn test_pooled_debug_and_empty_freeze_paths() {
2216 let page = page_size();
2219 let pool = test_pool(test_config(page, page, 3));
2220
2221 let pooled_mut = pool.try_alloc(page).expect("pooled allocation");
2222 let pooled_mut_debug = format!("{pooled_mut:?}");
2223 assert!(pooled_mut_debug.contains("IoBufMut"));
2224 assert!(pooled_mut_debug.contains("cap"));
2225 assert!(pooled_mut.is_pooled());
2226
2227 let empty = pool.try_alloc(page).expect("pooled allocation").freeze();
2228 assert!(empty.is_empty());
2229 assert!(!empty.is_pooled());
2230
2231 let mut non_empty = pool.try_alloc(page).expect("pooled allocation");
2232 non_empty.put_slice(b"abc");
2233 let pooled = non_empty.freeze();
2234 let pooled_debug = format!("{pooled:?}");
2235 assert!(pooled_debug.contains("IoBuf"));
2236 assert!(pooled_debug.contains("pooled"));
2237 assert!(pooled.is_pooled());
2238
2239 BufferPoolThreadCache::flush();
2240 }
2241
2242 #[test]
2243 fn test_freeze_returns_buffer_to_pool() {
2244 let page = page_size();
2245 let pool = test_pool(test_config(page, page, 2));
2246
2247 assert_eq!(get_allocated(&pool, page), 0);
2249 assert_eq!(get_available(&pool, page), 0);
2250
2251 let mut buf = pool.try_alloc(page).unwrap();
2254 buf.put_slice(b"x");
2255 assert_eq!(get_allocated(&pool, page), 1);
2256 assert_eq!(get_available(&pool, page), 0);
2257
2258 let iobuf = buf.freeze();
2259 assert_eq!(get_allocated(&pool, page), 1);
2261
2262 drop(iobuf);
2264 assert_eq!(get_allocated(&pool, page), 0);
2265 assert_eq!(get_available(&pool, page), 1);
2266 }
2267
2268 #[test]
2269 fn test_refcount_and_copy_to_bytes_paths() {
2270 let page = page_size();
2271 let pool = test_pool(test_config(page, page, 2));
2272
2273 {
2277 let mut buf = pool.try_alloc(page).unwrap();
2278 buf.put_slice(&[0xAA; 100]);
2279 let iobuf = buf.freeze();
2280 let clone = iobuf.clone();
2281 let slice = iobuf.slice(10..40);
2282 let empty = iobuf.slice(10..10);
2283 assert!(empty.is_empty());
2284 drop(iobuf);
2285 assert_eq!(get_allocated(&pool, page), 1);
2286 drop(slice);
2287 assert_eq!(get_allocated(&pool, page), 1);
2288 drop(clone);
2289 assert_eq!(get_allocated(&pool, page), 0);
2290 }
2291
2292 {
2298 let mut buf = pool.try_alloc(page).unwrap();
2299 buf.put_slice(&[0x42; 100]);
2300 let mut iobuf = buf.freeze();
2301
2302 let zero = iobuf.copy_to_bytes(0);
2303 assert!(zero.is_empty());
2304 assert_eq!(iobuf.remaining(), 100);
2305
2306 let partial = iobuf.copy_to_bytes(30);
2307 assert_eq!(&partial[..], &[0x42; 30]);
2308 assert_eq!(iobuf.remaining(), 70);
2309
2310 let rest = iobuf.copy_to_bytes(70);
2311 assert_eq!(&rest[..], &[0x42; 70]);
2312 assert_eq!(iobuf.remaining(), 0);
2313
2314 let empty = iobuf.copy_to_bytes(0);
2316 assert!(empty.is_empty());
2317
2318 drop(iobuf);
2319 assert_eq!(get_allocated(&pool, page), 1);
2320 drop(zero);
2321 drop(partial);
2322 assert_eq!(get_allocated(&pool, page), 1);
2323 drop(rest);
2324 assert_eq!(get_allocated(&pool, page), 0);
2325 }
2326
2327 {
2329 let buf = pool.try_alloc(page).unwrap();
2330 let mut iobufmut = buf;
2331 iobufmut.put_slice(&[0x7E; 100]);
2332
2333 let zero = iobufmut.copy_to_bytes(0);
2334 assert!(zero.is_empty());
2335 assert_eq!(iobufmut.remaining(), 100);
2336
2337 let partial = iobufmut.copy_to_bytes(30);
2338 assert_eq!(&partial[..], &[0x7E; 30]);
2339 assert_eq!(iobufmut.remaining(), 70);
2340
2341 let rest = iobufmut.copy_to_bytes(70);
2342 assert_eq!(&rest[..], &[0x7E; 70]);
2343 assert_eq!(iobufmut.remaining(), 0);
2344
2345 drop(iobufmut);
2346 assert_eq!(get_allocated(&pool, page), 1);
2347 drop(zero);
2348 drop(partial);
2349 assert_eq!(get_allocated(&pool, page), 1);
2350 drop(rest);
2351 assert_eq!(get_allocated(&pool, page), 0);
2352 }
2353 }
2354
2355 #[test]
2356 fn test_iobuf_to_iobufmut_conversion_reuses_pool_for_non_full_unique_view() {
2357 let page = page_size();
2359 let pool = test_pool(test_config(page, page, 2));
2360
2361 let mut buf = pool.try_alloc(page).unwrap();
2362 buf.put_slice(b"non-full");
2363 assert_eq!(get_allocated(&pool, page), 1);
2364
2365 let iobuf = buf.freeze();
2366 assert_eq!(iobuf.len(), 8);
2367 assert_eq!(get_allocated(&pool, page), 1);
2368
2369 let iobufmut: IoBufMut = iobuf.into();
2370 assert_eq!(iobufmut.as_ref(), b"non-full");
2371
2372 assert_eq!(
2374 get_allocated(&pool, page),
2375 1,
2376 "pooled buffer should remain allocated after zero-copy IoBuf->IoBufMut conversion"
2377 );
2378 assert_eq!(get_available(&pool, page), 0);
2379
2380 drop(iobufmut);
2382 assert_eq!(get_allocated(&pool, page), 0);
2383 assert_eq!(get_available(&pool, page), 1);
2384 }
2385
2386 #[test]
2387 fn test_iobuf_try_into_mut_recycles_full_unique_view() {
2388 let page = page_size();
2391 let pool = test_pool(test_config(page, page, 2));
2392
2393 let mut buf = pool.try_alloc(page).unwrap();
2394 buf.put_slice(&vec![0xAB; page]);
2395 let iobuf = buf.freeze();
2396 assert_eq!(get_allocated(&pool, page), 1);
2397
2398 let recycled = iobuf
2400 .try_into_mut()
2401 .expect("unique full-view pooled buffer should recycle");
2402 assert_eq!(recycled.len(), page);
2403 assert!(recycled.as_ref().iter().all(|&b| b == 0xAB));
2404 assert_eq!(recycled.capacity(), page);
2405 assert_eq!(get_allocated(&pool, page), 1);
2406
2407 drop(recycled);
2408 assert_eq!(get_allocated(&pool, page), 0);
2409 assert_eq!(get_available(&pool, page), 1);
2410 }
2411
2412 #[test]
2413 fn test_iobuf_try_into_mut_succeeds_for_unique_slice_and_fails_for_shared() {
2414 let page = page_size();
2415 let pool = test_pool(test_config(page, page, 2));
2416
2417 let mut buf = pool.try_alloc(page).unwrap();
2419 buf.put_slice(&vec![0xCD; page]);
2420 let iobuf = buf.freeze();
2421 let sliced = iobuf.slice(1..page);
2422 drop(iobuf);
2423 let recycled = sliced
2424 .try_into_mut()
2425 .expect("unique sliced pooled buffer should recycle");
2426 assert_eq!(recycled.len(), page - 1);
2427 assert!(recycled.as_ref().iter().all(|&b| b == 0xCD));
2428 assert_eq!(recycled.capacity(), page - 1);
2429 assert_eq!(get_allocated(&pool, page), 1);
2430 drop(recycled);
2431 assert_eq!(get_allocated(&pool, page), 0);
2432 assert_eq!(get_available(&pool, page), 1);
2433
2434 let mut buf = pool.try_alloc(page).unwrap();
2436 buf.put_slice(&vec![0xEF; page]);
2437 let iobuf = buf.freeze();
2438 let cloned = iobuf.clone();
2439 let iobuf = iobuf
2440 .try_into_mut()
2441 .expect_err("shared pooled buffer must not convert to mutable");
2442
2443 drop(cloned);
2444 drop(iobuf);
2445 assert_eq!(get_allocated(&pool, page), 0);
2446 assert!(get_available(&pool, page) >= 1);
2447 }
2448
2449 #[test]
2450 fn test_multithreaded_alloc_freeze_return() {
2451 let page = page_size();
2452 let pool = Arc::new(test_pool(test_config(page, page, 100)));
2453
2454 let mut handles = vec![];
2455
2456 cfg_if::cfg_if! {
2458 if #[cfg(miri)] {
2459 let iterations = 100;
2460 } else {
2461 let iterations = 1000;
2462 }
2463 }
2464
2465 for _ in 0..10 {
2467 let pool = pool.clone();
2468 let handle = thread::spawn(move || {
2469 for _ in 0..iterations {
2470 let mut buf = pool.try_alloc(page).unwrap();
2471 buf.put_slice(b"x");
2475 let iobuf = buf.freeze();
2476
2477 let clones: Vec<_> = (0..5).map(|_| iobuf.clone()).collect();
2479 drop(iobuf);
2480
2481 for clone in clones {
2483 drop(clone);
2484 }
2485 }
2486 });
2487 handles.push(handle);
2488 }
2489
2490 for handle in handles {
2492 handle.join().unwrap();
2493 }
2494
2495 let _buf = pool
2499 .try_alloc(page)
2500 .expect("pool should remain usable after multithreaded test");
2501 }
2502
2503 #[test]
2504 fn test_cross_thread_buffer_return() {
2505 let page = page_size();
2507 let pool = test_pool(test_config(page, page, 100));
2508
2509 let (tx, rx) = mpsc::channel();
2510
2511 for _ in 0..50 {
2513 let mut buf = pool.try_alloc(page).unwrap();
2514 buf.put_slice(b"x");
2515 let iobuf = buf.freeze();
2516 tx.send(iobuf).unwrap();
2517 }
2518 drop(tx);
2519
2520 let handle = thread::spawn(move || {
2524 while let Ok(iobuf) = rx.recv() {
2525 drop(iobuf);
2526 }
2527
2528 let class_index = pool
2529 .class_index(page)
2530 .expect("class exists for page-sized buffer");
2531 assert_eq!(get_local_len(&pool.inner.classes[class_index]), 50);
2532 assert_eq!(get_global_len(&pool.inner.classes[class_index]), 0);
2533
2534 for _ in 0..50 {
2535 let _buf = pool
2536 .try_alloc(page)
2537 .expect("dropping thread should be able to reuse locally returned buffers");
2538 }
2539 });
2540
2541 handle.join().unwrap();
2542 }
2543
2544 #[test]
2545 fn test_pool_dropped_before_buffer() {
2546 let page = page_size();
2550 let pool = test_pool(test_config(page, page, 2));
2551
2552 let mut buf = pool.try_alloc(page).unwrap();
2553 buf.put_slice(&[0u8; 100]);
2554 let iobuf = buf.freeze();
2555
2556 drop(pool);
2558
2559 assert_eq!(iobuf.len(), 100);
2561
2562 drop(iobuf);
2564 }
2566
2567 #[test]
2568 fn test_pool_exhaustion_and_recovery() {
2569 let page = page_size();
2571 let pool = test_pool(test_config(page, page, 3));
2572
2573 let buf1 = pool.try_alloc(page).expect("first alloc");
2575 let buf2 = pool.try_alloc(page).expect("second alloc");
2576 let buf3 = pool.try_alloc(page).expect("third alloc");
2577 assert!(pool.try_alloc(page).is_err(), "pool should be exhausted");
2578
2579 drop(buf1);
2581
2582 let buf4 = pool.try_alloc(page).expect("alloc after return");
2584 assert!(pool.try_alloc(page).is_err(), "pool exhausted again");
2585
2586 drop(buf2);
2588 drop(buf3);
2589 drop(buf4);
2590
2591 assert_eq!(get_allocated(&pool, page), 0);
2592 assert_eq!(get_available(&pool, page), 3);
2593
2594 let _buf5 = pool.try_alloc(page).expect("reuse from freelist");
2596 assert_eq!(get_available(&pool, page), 2);
2597 }
2598
2599 #[test]
2600 fn test_try_alloc_errors() {
2601 let page = page_size();
2603 let pool = test_pool(test_config(page, page, 2));
2604
2605 let result = pool.try_alloc(page * 10);
2607 assert_eq!(result.unwrap_err(), PoolError::Oversized);
2608
2609 let _buf1 = pool.try_alloc(page).unwrap();
2611 let _buf2 = pool.try_alloc(page).unwrap();
2612 let result = pool.try_alloc(page);
2613 assert_eq!(result.unwrap_err(), PoolError::Exhausted);
2614 }
2615
2616 #[test]
2617 fn test_pool_metrics_track_created_exhausted_oversized() {
2618 let page = page_size();
2619 let mut registry = Registry::default();
2620 let pool = BufferPool::new(test_config(page, page, 1), &mut registry);
2621
2622 let buf = pool.try_alloc(page).unwrap();
2624 assert_eq!(pool.try_alloc(page).unwrap_err(), PoolError::Exhausted);
2625 assert_eq!(pool.try_alloc(page * 2).unwrap_err(), PoolError::Oversized);
2626
2627 let encoded = registry.encode();
2628 assert!(
2629 encoded.contains(&format!("buffer_pool_created{{size_class=\"{page}\"}} 1")),
2630 "created gauge missing: {encoded}"
2631 );
2632 assert!(
2633 encoded.contains(&format!(
2634 "buffer_pool_exhausted_total{{size_class=\"{page}\"}} 1"
2635 )),
2636 "exhausted counter missing: {encoded}"
2637 );
2638 assert!(
2639 encoded.contains("buffer_pool_oversized_total 1"),
2640 "oversized counter missing: {encoded}"
2641 );
2642 drop(buf);
2643 }
2644
2645 #[test]
2646 fn test_try_alloc_zeroed_errors() {
2647 let page = page_size();
2649 let pool = test_pool(test_config(page, page, 2));
2650
2651 let result = pool.try_alloc_zeroed(page * 10);
2653 assert_eq!(result.unwrap_err(), PoolError::Oversized);
2654
2655 let _buf1 = pool.try_alloc_zeroed(page).unwrap();
2657 let _buf2 = pool.try_alloc_zeroed(page).unwrap();
2658 let result = pool.try_alloc_zeroed(page);
2659 assert_eq!(result.unwrap_err(), PoolError::Exhausted);
2660 }
2661
2662 #[test]
2663 fn test_fallback_allocation() {
2664 let page = page_size();
2666 let pool = test_pool(test_config(page, page, 2));
2667
2668 let buf1 = pool.try_alloc(page).unwrap();
2670 let buf2 = pool.try_alloc(page).unwrap();
2671 assert!(buf1.is_pooled());
2672 assert!(buf2.is_pooled());
2673
2674 let mut fallback_exhausted = pool.alloc(page);
2677 assert!(!fallback_exhausted.is_pooled());
2678 assert!((fallback_exhausted.as_mut_ptr() as usize).is_multiple_of(page));
2679 assert_eq!(fallback_exhausted.capacity(), page);
2680
2681 let fallback_small = pool.alloc(100);
2682 assert!(!fallback_small.is_pooled());
2683 assert!((100..108).contains(&fallback_small.capacity()));
2684
2685 let mut fallback_oversized = pool.alloc(page * 10);
2687 assert!(!fallback_oversized.is_pooled());
2688 assert!((fallback_oversized.as_mut_ptr() as usize).is_multiple_of(page));
2689 assert_eq!(fallback_oversized.capacity(), page * 10);
2690
2691 assert_eq!(get_allocated(&pool, page), 2);
2693
2694 drop(fallback_exhausted);
2696 drop(fallback_oversized);
2697 assert_eq!(get_allocated(&pool, page), 2);
2698
2699 drop(buf1);
2701 drop(buf2);
2702 assert_eq!(get_allocated(&pool, page), 0);
2703 }
2704
2705 #[test]
2706 fn test_is_pooled() {
2707 let page = page_size();
2710 let pool = test_pool(test_config(page, page, 10));
2711
2712 let pooled = pool.try_alloc(page).unwrap();
2713 assert!(pooled.is_pooled());
2714
2715 let owned = IoBufMut::with_capacity(100);
2716 assert!(!owned.is_pooled());
2717 }
2718
2719 #[test]
2720 fn test_iobuf_is_pooled() {
2721 let page = page_size();
2722 let pool = test_pool(test_config(page, page, 2));
2723
2724 let mut pooled = pool.try_alloc(page).unwrap();
2725 pooled.put_slice(b"x");
2726 let pooled = pooled.freeze();
2727 assert!(pooled.is_pooled());
2728
2729 let fallback = pool.alloc(page * 10).freeze();
2731 assert!(!fallback.is_pooled());
2732
2733 let bytes = IoBuf::copy_from_slice(b"hello");
2734 assert!(!bytes.is_pooled());
2735 }
2736
2737 #[test]
2738 fn test_buffer_alignment() {
2739 let page = page_size();
2740 let cache_line = cache_line_size();
2741
2742 cfg_if::cfg_if! {
2744 if #[cfg(miri)] {
2745 let storage_config = BufferPoolConfig::for_storage()
2746 .with_alignment(NZUsize!(page))
2747 .with_max_per_class(NZU32!(32));
2748 let network_config = BufferPoolConfig::for_network()
2749 .with_alignment(NZUsize!(cache_line))
2750 .with_max_per_class(NZU32!(32));
2751 } else {
2752 let storage_config =
2753 BufferPoolConfig::for_storage().with_alignment(NZUsize!(page));
2754 let network_config =
2755 BufferPoolConfig::for_network().with_alignment(NZUsize!(cache_line));
2756 }
2757 }
2758
2759 let storage_buffer_pool = test_pool(storage_config);
2761 let mut buf = storage_buffer_pool.try_alloc(100).unwrap();
2762 assert_eq!(
2763 buf.as_mut_ptr() as usize % page,
2764 0,
2765 "storage buffer not page-aligned"
2766 );
2767
2768 let network_buffer_pool = test_pool(network_config);
2770 let mut buf = network_buffer_pool.try_alloc(100).unwrap();
2771 assert_eq!(
2772 buf.as_mut_ptr() as usize % cache_line,
2773 0,
2774 "network buffer not cache-line aligned"
2775 );
2776 }
2777}
2778
2779#[cfg(all(test, feature = "loom"))]
2780mod loom_tests {
2781 use super::*;
2782 use crate::telemetry::metrics::Registry;
2783 use bytes::BufMut;
2784 use loom::thread;
2785
2786 #[test]
2794 fn freeze_clone_cross_thread_drop_then_reuse() {
2795 loom::model(|| {
2796 let mut registry = Registry::default();
2797 let config = BufferPoolConfig::for_network()
2798 .with_size_class_range(NZUsize!(64), NZUsize!(64), NZU32!(2))
2799 .with_thread_cache_disabled();
2800 let pool = BufferPool::new(config, &mut registry);
2801
2802 let mut buf = pool.alloc(64);
2803 assert!(buf.is_pooled());
2804 buf.put_slice(b"payload");
2805 let frozen = buf.freeze();
2806 let clone = frozen.clone();
2807
2808 let t = thread::spawn(move || {
2809 assert_eq!(clone.as_ref(), b"payload");
2810 drop(clone);
2811 });
2812 assert_eq!(frozen.as_ref(), b"payload");
2813 drop(frozen);
2814 t.join().unwrap();
2815
2816 let mut again = pool.alloc(64);
2819 assert!(again.is_pooled());
2820 again.put_slice(b"reuse");
2821 assert_eq!(again.as_ref(), b"reuse");
2822 });
2823 }
2824
2825 #[test]
2837 fn final_drop_races_pool_teardown() {
2838 loom::model(|| {
2839 let mut registry = Registry::default();
2840 let config = BufferPoolConfig::for_network()
2841 .with_size_class_range(NZUsize!(64), NZUsize!(64), NZU32!(1))
2842 .with_thread_cache_disabled();
2843 let pool = BufferPool::new(config, &mut registry);
2844
2845 let mut buf = pool.alloc(64);
2846 assert!(buf.is_pooled());
2847 buf.put_slice(b"x");
2848 let frozen = buf.freeze();
2849
2850 let t = thread::spawn(move || drop(frozen));
2851 drop(pool);
2852 t.join().unwrap();
2853 });
2854 }
2855
2856 #[test]
2864 fn final_drop_races_recheckout() {
2865 loom::model(|| {
2866 let mut registry = Registry::default();
2867 let config = BufferPoolConfig::for_network()
2868 .with_size_class_range(NZUsize!(64), NZUsize!(64), NZU32!(1))
2869 .with_thread_cache_disabled();
2870 let pool = BufferPool::new(config, &mut registry);
2871
2872 let mut buf = pool.alloc(64);
2873 assert!(buf.is_pooled());
2874 buf.put_slice(b"x");
2875 let frozen = buf.freeze();
2876 let clone = frozen.clone();
2877
2878 let t = thread::spawn(move || drop(clone));
2879 drop(frozen);
2880
2881 if let Ok(mut again) = pool.try_alloc(64) {
2885 assert!(again.is_pooled());
2886 again.put_slice(b"y");
2887 assert_eq!(again.as_ref(), b"y");
2888 }
2889 t.join().unwrap();
2890 });
2891 }
2892}