1#![cfg_attr(not(feature = "std"), no_std)]
37
38#[cfg(not(feature = "std"))]
39extern crate alloc;
40
41#[cfg(not(feature = "std"))]
42use core::future::poll_fn;
43
44#[cfg(feature = "std")]
45use std::future::poll_fn;
46
47#[cfg(not(feature = "std"))]
48use core::
49{
50 fmt,
51 mem,
52 ops::{Deref, DerefMut},
53 ptr::{self, NonNull},
54 sync::atomic::{AtomicU64, Ordering}
55};
56
57#[cfg(not(feature = "std"))]
58use core::{marker::PhantomData, task::Poll, time::Duration};
59
60#[cfg(not(feature = "std"))]
61use alloc::{boxed::Box, collections::vec_deque::VecDeque, vec::Vec};
62
63#[cfg(not(feature = "std"))]
64use alloc::vec;
65
66
67
68use crossbeam_utils::Backoff;
69
70#[cfg(feature = "std")]
71use std::
72{
73 collections::VecDeque,
74 ops::{Deref, DerefMut},
75 ptr::{self, NonNull},
76 sync::atomic::{AtomicU64, Ordering},
77 fmt,
78 mem
79};
80
81#[cfg(feature = "std")]
82use std::{marker::PhantomData, task::Poll};
83
84#[cfg(all(feature = "std", not(feature = "clone_wait_indef")))]
85use std::time::Duration;
86
87extern crate crossbeam_utils;
88
89pub trait TryClone: Sized
91{
92 type Error;
93
94 fn try_clone(&self) -> Result<Self, Self::Error>;
97}
98
99#[cfg(feature = "enable_async")]
100pub mod sbr_async
101{
102 pub trait LocalAsyncDrop: Send + Sync + 'static
104 {
105 fn async_drop(&mut self) -> impl Future<Output = ()>;
107 }
108
109 pub trait LocalAsyncClone: Send + Sync + 'static
111 {
112 fn async_clone(&self) -> impl Future<Output = Self>;
114 }
115
116 pub async
119 fn async_drop<LAD: LocalAsyncDrop + Send + Sync>(mut lad: LAD)
120 {
121 lad.async_drop().await;
122
123 drop(lad);
124 }
125}
126
127#[cfg(feature = "enable_async")]
128pub use self::sbr_async::{async_drop, LocalAsyncClone, LocalAsyncDrop};
129
130#[derive(Debug, Clone, Copy, PartialEq, Eq)]
131pub enum RwBufferError
132{
133 TooManyRead,
135
136 TooManyBase,
138
139 ReadTryAgianLater,
141
142 WriteTryAgianLater,
144
145 BaseTryAgainLater,
147
148 OutOfBuffers,
150
151 DowngradeFailed,
153
154 InvalidArguments,
156
157 Busy,
159}
160
161impl fmt::Display for RwBufferError
162{
163 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result
164 {
165 match self
166 {
167 Self::TooManyRead =>
168 write!(f, "TooManyRead: read soft limit reached"),
169 Self::TooManyBase =>
170 write!(f, "TooManyBase: base soft limit reached"),
171 Self::ReadTryAgianLater =>
172 write!(f, "ReadTryAgianLater: shared access not available, try again later"),
173 Self::WriteTryAgianLater =>
174 write!(f, "WriteTryAgianLater: exclusive access not available, try again later"),
175 Self::BaseTryAgainLater =>
176 write!(f, "BaseTryAgainLater: failed to obtain a clone in reasonable time"),
177 Self::OutOfBuffers =>
178 write!(f, "OutOfBuffers: no more free bufers are left"),
179 Self::DowngradeFailed =>
180 write!(f, "DowngradeFailed: can not downgrade exclusive to shared, race condition"),
181 Self::InvalidArguments =>
182 write!(f, "InvalidArguments: arguments are not valid"),
183 Self::Busy =>
184 write!(f, "RwBuffer is busy and cannot be acquired"),
185 }
186 }
187}
188
189pub type RwBufferRes<T> = Result<T, RwBufferError>;
190
191#[derive(Debug)]
194pub struct RBuffer
195{
196 inner: NonNull<RwBufferInner>,
198
199 a_dropped: bool,
201}
202
203unsafe impl Send for RBuffer {}
204unsafe impl Sync for RBuffer {}
205
206impl RwBufType for RBuffer {}
207
208impl Eq for RBuffer {}
209
210impl PartialEq for RBuffer
211{
212 fn eq(&self, other: &Self) -> bool
213 {
214 return self.inner == other.inner;
215 }
216}
217
218impl RBuffer
219{
220 #[inline]
221 fn new(inner: NonNull<RwBufferInner>) -> Self
222 {
223 return Self{ inner, a_dropped: false };
224 }
225
226 #[cfg(test)]
227 fn get_flags(&self) -> RwBufferFlags<Self>
228 {
229 use core::sync::atomic::Ordering;
230
231 let inner = unsafe{ self.inner.as_ref() };
232
233 let flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
234
235 return flags;
236 }
237
238 pub
240 fn as_slice(&self) -> &[u8]
241 {
242 let inner = unsafe { self.inner.as_ref() };
243
244 return inner.buf.as_ref().unwrap().as_slice();
245 }
246
247 pub
265 fn try_inner(mut self) -> Result<Vec<u8>, Self>
266 {
267 let inner = unsafe { self.inner.as_ref() };
268
269 let current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::SeqCst).into();
270
271 if current_flags.try_inner_check() == true
274 {
275 let inner = unsafe { self.inner.as_mut() };
279
280 let buf = inner.buf.take().unwrap();
281
282 drop(self);
285
286 return Ok(buf);
287 }
288
289 return Err(self);
290 }
291
292 fn inner(&self) -> &RwBufferInner
293 {
294 return unsafe { self.inner.as_ref() };
295 }
296}
297
298impl Deref for RBuffer
299{
300 type Target = Vec<u8>;
301
302 fn deref(&self) -> &Vec<u8>
303 {
304 let inner = self.inner();
305
306 return inner.buf.as_ref().unwrap();
307 }
308}
309
310#[cfg(feature = "enable_async")]
311impl sbr_async::LocalAsyncClone for RBuffer
312{
313 fn async_clone(&self) -> impl Future<Output = Self>
314 {
315 return
316 poll_fn(
317 |_cx|
318 {
319 let inner = self.inner();
320
321 let current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
322 let mut new_flags = current_flags.clone();
323
324 new_flags.read().unwrap();
325
326 let res =
327 inner
328 .flags
329 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
330
331 if let Ok(_) = res
332 {
333 return Poll::Ready( Self{ inner: self.inner, a_dropped: self.a_dropped } );
334 }
335
336 return Poll::Pending;
337 }
338 );
339 }
340}
341
342impl Clone for RBuffer
343{
344 fn clone(&self) -> Self
360 {
361 let inner = self.inner();
362
363 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
364 let mut new_flags = current_flags.clone();
365
366 new_flags.read().unwrap();
367
368 let backoff = Backoff::new();
369
370 #[cfg(all(feature = "std", not(feature = "clone_wait_indef")))]
371 let mut parked = false;
372
373 loop
374 {
375 let res =
377 inner
378 .flags
379 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
380
381 if let Ok(_) = res
382 {
383 return Self{ inner: self.inner, a_dropped: self.a_dropped };
384 }
385
386 current_flags = res.err().unwrap().into();
387 new_flags = current_flags.clone();
388
389 new_flags.read().unwrap();
390
391 if backoff.is_completed() == false
392 {
393 backoff.snooze();
394 }
395 else
396 {
397 #[cfg(all(feature = "std", not(feature = "clone_wait_indef")))]
398 {
399 if parked == false
400 {
401 std::thread::park_timeout(Duration::from_millis(1));
403
404 parked = true;
405 }
406 else
407 {
408 panic!("can not obtain a clone of RBuffer in reasonable time!");
409 }
410 }
411
412 #[cfg(all(not(feature = "std"), not(feature = "clone_wait_indef")))]
413 {
414 panic!("can not obtain a clone of RBuffer in reasonable time!");
415 }
416
417 #[cfg(feature = "clone_wait_indef")]
418 {
419 backoff.reset();
420 }
421
422 }
423 }
424 }
425}
426
427impl TryClone for RBuffer
428{
429 type Error = RwBufferError;
430
431 fn try_clone(&self) -> Result<Self, Self::Error>
443 {
444 let inner = self.inner();
445
446 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::SeqCst).into();
447 let mut new_flags = current_flags.clone();
448
449 new_flags.read()?;
450
451 let backoff = Backoff::new();
452
453 loop
454 {
455 let res =
456 inner
457 .flags
458 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
459
460 if let Ok(_) = res
461 {
462 return Ok(Self{ inner: self.inner, a_dropped: self.a_dropped });
463 }
464
465 current_flags = res.err().unwrap().into();
466 new_flags = current_flags.clone();
467
468 new_flags.read()?;
469
470 if backoff.is_completed() == false
471 {
472 backoff.snooze();
473 }
474 else
475 {
476 break;
477 }
478 }
479
480 return Err(RwBufferError::ReadTryAgianLater);
481 }
482}
483
484#[cfg(feature = "enable_async")]
485impl sbr_async::LocalAsyncDrop for RBuffer
486{
487 fn async_drop(&mut self) -> impl Future<Output = ()>
488 {
489 self.a_dropped = true;
490
491 return
492 poll_fn(
493 |cx|
494 {
495 let inner = self.inner();
496
497 let current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::SeqCst).into();
498 let mut new_flags = current_flags.clone();
499
500 new_flags.unread();
501
502 let res =
503 inner
504 .flags
505 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
506
507 if let Ok(flags) = res.map(|v| <u64 as Into<RwBufferFlags<Self>>>::into(v))
508 {
509 if flags.is_drop_inplace() == true
510 {
511 unsafe { ptr::drop_in_place(self.inner.as_ptr()) };
513 }
514
515 return Poll::Ready(());
516 }
517
518 cx.waker().wake_by_ref();
519
520 return Poll::Pending;
521 }
522 );
523 }
524}
525
526impl Drop for RBuffer
527{
528 fn drop(&mut self)
536 {
537 if self.a_dropped == true
538 {
539 return;
540 }
541
542 let inner = self.inner();
543
544 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
545 let mut new_flags = current_flags.clone();
546
547 new_flags.unread();
548
549 let backoff = Backoff::new();
550
551 loop
552 {
553 let res =
554 inner
555 .flags
556 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
557
558 if let Ok(flags) = res.map(|v| <u64 as Into<RwBufferFlags<Self>>>::into(v))
559 {
560 if flags.is_drop_inplace() == true
561 {
562 unsafe { ptr::drop_in_place(self.inner.as_ptr()) };
564 }
565
566 return;
567 }
568
569 if backoff.is_completed() == true
570 {
571 panic!("assertion trap: RBuffer::drop can not drop RBuffer in reasonable time!");
573 }
574
575 current_flags = res.err().unwrap().into();
576 new_flags = current_flags.clone();
577
578 new_flags.unread();
579
580 backoff.snooze();
581 }
582
583
584 }
585}
586
587#[derive(Debug, PartialEq, Eq)]
591pub struct WBuffer
592{
593 buf: NonNull<RwBufferInner>,
595
596 downgraded: bool,
598}
599
600unsafe impl Send for WBuffer{}
601unsafe impl Sync for WBuffer{}
602
603impl RwBufType for WBuffer{}
604
605impl WBuffer
606{
607 #[inline]
608 fn new(inner: NonNull<RwBufferInner>) -> Self
609 {
610 return Self{ buf: inner, downgraded: false };
611 }
612
613 pub
624 fn downgrade(mut self) -> Result<RBuffer, Self>
625 {
626 let inner = unsafe { self.buf.as_ref() };
627
628 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
629 let mut new_flags = current_flags.clone();
630
631 new_flags.downgrade();
632
633 let backoff = Backoff::new();
634
635 while backoff.is_completed() == false
636 {
637 let res =
638 inner
639 .flags
640 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
641
642 if let Ok(_) = res
643 {
644 self.downgraded = true;
645
646 return Ok(RBuffer::new(self.buf.clone()));
647 }
648
649 current_flags = res.err().unwrap().into();
650 new_flags = current_flags.clone();
651
652 new_flags.downgrade();
653
654 backoff.snooze();
655 }
656
657 return Err(self);
658 }
659
660 pub
661 fn as_slice(&self) -> &[u8]
662 {
663 let inner = unsafe { self.buf.as_ref() };
664
665 return inner.buf.as_ref().unwrap()
666 }
667}
668
669impl Deref for WBuffer
670{
671 type Target = Vec<u8>;
672
673 fn deref(&self) -> &Vec<u8>
674 {
675 let inner = unsafe { self.buf.as_ref() };
676
677 return inner.buf.as_ref().unwrap();
678 }
679}
680
681impl DerefMut for WBuffer
682{
683 fn deref_mut(&mut self) -> &mut Vec<u8>
684 {
685 let inner = unsafe { self.buf.as_mut() };
686
687 return inner.buf.as_mut().unwrap();
688 }
689}
690
691impl Drop for WBuffer
692{
693 fn drop(&mut self)
695 {
696 if self.downgraded == true
697 {
698 return;
699 }
700
701 let inner = unsafe { self.buf.as_ref() };
702
703 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
704 let mut new_flags = current_flags.clone();
705
706 new_flags.unwrite();
707
708 let backoff = Backoff::new();
709
710 for _ in 0..1000
711 {
712 let res =
713 inner
714 .flags
715 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
716
717 if let Ok(flags) = res.map(|v| <u64 as Into<RwBufferFlags<Self>>>::into(v))
718 {
719 if flags.is_drop_inplace() == true
720 {
721 unsafe { ptr::drop_in_place(self.buf.as_ptr()) };
723 }
724
725 return;
726 }
727
728 current_flags = res.err().unwrap().into();
729 new_flags = current_flags.clone();
730
731 new_flags.unwrite();
732
733 backoff.snooze();
734 }
735
736 panic!("assertion trap: WBuffer::drop can not drop RBuffer in reasonable time!");
739 }
740}
741
742#[cfg(feature = "enable_async")]
743impl sbr_async::LocalAsyncDrop for WBuffer
744{
745 fn async_drop(&mut self) -> impl Future<Output = ()>
746 {
747 self.downgraded = true;
748
749 return
750 poll_fn(
751 move |cx|
752 {
753 let inner = unsafe { self.buf.as_ref() };
754
755 let current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::SeqCst).into();
756 let mut new_flags = current_flags.clone();
757
758 new_flags.unwrite();
759
760 let res =
761 inner
762 .flags
763 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
764
765 if let Ok(flags) = res.map(|v| <u64 as Into<RwBufferFlags<Self>>>::into(v))
766 {
767 if flags.is_drop_inplace() == true
768 {
769 unsafe { ptr::drop_in_place(self.buf.as_ptr()) };
771 }
772
773 return Poll::Ready(());
774 }
775
776 cx.waker().wake_by_ref();
777
778 return Poll::Pending;
779 }
780 );
781 }
782}
783
784trait RwBufType {}
785
786#[repr(align(8))]
789#[derive(Debug, PartialEq, Eq)]
790struct RwBufferFlags<TP>
791{
792 read: u32, write: bool, base: u16, unused0: u8, _p: PhantomData<TP>,
807}
808
809
810impl<TP: RwBufType> From<u64> for RwBufferFlags<TP>
811{
812 fn from(value: u64) -> Self
813 {
814 return unsafe { mem::transmute(value) };
815 }
816}
817
818impl<TP: RwBufType> From<RwBufferFlags<TP>> for u64
819{
820 fn from(value: RwBufferFlags<TP>) -> Self
821 {
822 return unsafe { mem::transmute(value) };
823 }
824}
825
826impl<TP: RwBufType> Default for RwBufferFlags<TP>
827{
828 fn default() -> RwBufferFlags<TP>
829 {
830 return
831 Self
832 {
833 read: 0,
834 write: false,
835 base: 1,
836 unused0: 0,
837 _p: PhantomData
838 };
839 }
840}
841
842impl Copy for RwBufferFlags<WBuffer>{}
843
844impl Clone for RwBufferFlags<WBuffer>
845{
846 fn clone(&self) -> Self
847 {
848 return
849 Self
850 {
851 read: self.read.clone(),
852 write: self.write.clone(),
853 base: self.base.clone(),
854 unused0: self.unused0.clone(),
855 _p: PhantomData
856 }
857 }
858}
859
860impl Copy for RwBufferFlags<RBuffer>{}
861
862impl Clone for RwBufferFlags<RBuffer>
863{
864 fn clone(&self) -> Self
865 {
866 return
867 Self
868 {
869 read: self.read.clone(),
870 write: self.write.clone(),
871 base: self.base.clone(),
872 unused0: self.unused0.clone(),
873 _p: PhantomData
874 }
875 }
876}
877
878impl Copy for RwBufferFlags<RwBuffer>{}
879
880impl Clone for RwBufferFlags<RwBuffer>
881{
882 fn clone(&self) -> Self
883 {
884 return
885 Self
886 {
887 read: self.read.clone(),
888 write: self.write.clone(),
889 base: self.base.clone(),
890 unused0: self.unused0.clone(),
891 _p: PhantomData
892 }
893 }
894}
895
896impl RwBufferFlags<WBuffer>
897{
898 #[inline]
899 fn write(&mut self) -> RwBufferRes<()>
900 {
901 if self.read == 0
902 {
903 self.write = true;
904
905 return Ok(());
906 }
907 else
908 {
909 return Err(RwBufferError::WriteTryAgianLater);
910 }
911 }
912
913 #[inline]
914 fn downgrade(&mut self)
915 {
916 self.write = false;
917 self.read += 1;
918 }
919
920 #[inline]
921 fn unwrite(&mut self)
922 {
923 self.write = false;
924 }
925}
926
927impl RwBufferFlags<RBuffer>
928{
929 #[inline]
930 fn try_inner_check(&self) -> bool
931 {
932 return self.read == 1 && self.write == false && self.base == 0;
933 }
934
935 #[inline]
936 fn unread(&mut self)
937 {
938 self.read -= 1;
939 }
940
941 #[inline]
942 fn read(&mut self) -> RwBufferRes<()>
943 {
944 if self.write == false
945 {
946 self.read += 1;
947
948 if self.read <= Self::MAX_READ_REFS
949 {
950 return Ok(());
951 }
952
953 return Err(RwBufferError::TooManyRead);
954 }
955
956 return Err(RwBufferError::ReadTryAgianLater);
957 }
958}
959
960impl RwBufferFlags<RwBuffer>
961{
962 #[inline]
963 fn make_pre_unused() -> Self
964 {
965 return Self{ read: 0, write: false, base: 1, unused0: 0, _p: PhantomData };
966 }
967
968 #[inline]
969 fn read(&mut self) -> RwBufferRes<()>
970 {
971 if self.write == false
972 {
973 self.read += 1;
974
975 if self.read <= Self::MAX_READ_REFS
976 {
977 return Ok(());
978 }
979
980 return Err(RwBufferError::TooManyRead);
981 }
982
983 return Err(RwBufferError::ReadTryAgianLater);
984 }
985
986 #[inline]
987 fn write(&mut self) -> RwBufferRes<()>
988 {
989 if self.read == 0
990 {
991 self.write = true;
992
993 return Ok(());
994 }
995 else
996 {
997 return Err(RwBufferError::WriteTryAgianLater);
998 }
999 }
1000
1001 #[inline]
1002 fn base(&mut self) -> RwBufferRes<()>
1003 {
1004 self.base += 1;
1005
1006 if self.base <= Self::MAX_BASE_REFS
1007 {
1008 return Ok(());
1009 }
1010
1011 return Err(RwBufferError::TooManyBase);
1012 }
1013
1014 #[inline]
1015 fn unbase(&mut self) -> bool
1016 {
1017 self.base -= 1;
1018
1019 return self.base != 0;
1020 }
1021}
1022
1023impl<TP: RwBufType> RwBufferFlags<TP>
1024{
1025 pub const MAX_READ_REFS: u32 = u32::MAX - 2;
1027
1028 pub const MAX_BASE_REFS: u16 = u16::MAX - 2;
1030
1031
1032 #[inline]
1033 fn is_free(&self) -> bool
1034 {
1035 return self.write == false && self.read == 0 && self.base == 1;
1036 }
1037
1038 #[inline]
1039 fn is_drop_inplace(&self) -> bool
1040 {
1041 return self.read == 0 && self.write == false && self.base == 0;
1042 }
1043}
1044
1045#[derive(Debug)]
1046pub struct RwBufferInner
1047{
1048 flags: AtomicU64,
1050
1051 buf: Option<Vec<u8>>,
1053}
1054
1055impl RwBufferInner
1056{
1057 fn new(buf_size: usize) -> Self
1058 {
1059 return
1060 Self
1061 {
1062 flags:
1063 AtomicU64::new(RwBufferFlags::<RwBuffer>::default().into()),
1064 buf:
1065 Some(vec![0_u8; buf_size])
1066 };
1067 }
1068}
1069
1070#[derive(Debug, PartialEq, Eq)]
1077pub struct RwBuffer(NonNull<RwBufferInner>);
1078
1079unsafe impl Send for RwBuffer {}
1080unsafe impl Sync for RwBuffer {}
1081
1082impl RwBufType for RwBuffer {}
1083
1084impl RwBuffer
1085{
1086 #[inline]
1087 fn new(buf_size: usize) -> Self
1088 {
1089 let status = Box::new(RwBufferInner::new(buf_size));
1090
1091 return Self(Box::leak(status).into());
1092 }
1093
1094 #[inline]
1095 fn inner(&self) -> &RwBufferInner
1096 {
1097 return unsafe { self.0.as_ref() };
1098 }
1099
1100 #[inline]
1116 pub
1117 fn is_free(&self) -> bool
1118 {
1119 let inner = self.inner();
1120
1121 let flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
1122
1123 return flags.is_free();
1124 }
1125
1126 #[inline]
1143 pub(crate)
1144 fn acqiure_if_free(&self) -> RwBufferRes<Self>
1145 {
1146 let inner = self.inner();
1147
1148 let current_flags: RwBufferFlags<Self> = RwBufferFlags::make_pre_unused();
1149 let mut new_flags = current_flags.clone();
1150
1151 new_flags.base()?;
1152
1153 let res =
1154 inner
1155 .flags
1156 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1157
1158 if let Ok(_) = res
1159 {
1160 return Ok(Self(self.0.clone()));
1161 }
1162
1163 return Err(RwBufferError::Busy);
1164 }
1165
1166 pub
1179 fn write(&self) -> RwBufferRes<WBuffer>
1180 {
1181 let inner = self.inner();
1182
1183 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
1184 let mut new_flags = current_flags.clone();
1185
1186 new_flags.write()?;
1187
1188 let backoff = Backoff::new();
1189
1190 while backoff.is_completed() == false
1191 {
1192 let res =
1193 inner
1194 .flags
1195 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1196
1197 if let Ok(_) = res
1198 {
1199 return Ok(WBuffer::new(self.0.clone()));
1200 }
1201
1202 current_flags = res.err().unwrap().into();
1203 new_flags = current_flags.clone();
1204
1205 new_flags.write()?;
1206
1207 backoff.snooze();
1208 }
1209
1210 return Err(RwBufferError::WriteTryAgianLater);
1211 }
1212
1213 #[cfg(feature = "enable_async")]
1219 pub async
1220 fn write_async(&self) -> RwBufferRes<WBuffer>
1221 {
1222 return
1223 poll_fn(
1224 |cx|
1225 {
1226 let inner = self.inner();
1227
1228 let current_flags: RwBufferFlags<WBuffer> = inner.flags.load(Ordering::SeqCst).into();
1229 let mut new_flags = current_flags.clone();
1230
1231 if let Err(e) = new_flags.write()
1232 {
1233 return Poll::Ready(Err(e));
1234 }
1235
1236 let res =
1237 inner
1238 .flags
1239 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1240
1241 if let Ok(_) = res
1242 {
1243 return Poll::Ready( Ok( WBuffer::new(self.0.clone()) ) );
1244 }
1245
1246 cx.waker().wake_by_ref();
1247
1248 return Poll::Pending;
1249 }
1250 )
1251 .await;
1252 }
1253
1254
1255 pub
1273 fn read(&self) -> RwBufferRes<RBuffer>
1274 {
1275
1276 let inner = self.inner();
1277
1278 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
1279 let mut new_flags = current_flags.clone();
1280
1281 new_flags.read()?;
1282
1283 let backoff = Backoff::new();
1284
1285 while backoff.is_completed() == false
1286 {
1287 let res =
1288 inner
1289 .flags
1290 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1291
1292 if let Ok(_) = res
1293 {
1294 return Ok(RBuffer::new(self.0.clone()));
1295 }
1296
1297 current_flags = res.err().unwrap().into();
1298 new_flags = current_flags.clone();
1299
1300 new_flags.read()?;
1301
1302 backoff.snooze();
1303 }
1304
1305 return Err(RwBufferError::ReadTryAgianLater);
1306 }
1307
1308 #[cfg(feature = "enable_async")]
1314 pub async
1315 fn read_async(&self) -> RwBufferRes<RBuffer>
1316 {
1317 return
1318 poll_fn(
1319 |cx|
1320 {
1321 let inner = self.inner();
1322
1323 let current_flags: RwBufferFlags<RBuffer> = inner.flags.load(Ordering::SeqCst).into();
1324 let mut new_flags = current_flags.clone();
1325
1326 match new_flags.read()
1327 {
1328 Ok(_) => {},
1329 Err(RwBufferError::TooManyRead) =>
1330 return Poll::Ready(Err(RwBufferError::TooManyRead)),
1331 Err(RwBufferError::ReadTryAgianLater) =>
1332 {
1333 cx.waker().wake_by_ref();
1334
1335 return Poll::Pending;
1336 },
1337 Err(e) =>
1338 panic!("assertion trap: unknown error {} in Future for AsyncRBuffer", e)
1339 }
1340
1341 let res =
1342 inner
1343 .flags
1344 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1345
1346 if let Ok(_) = res
1347 {
1348 return Poll::Ready( Ok( RBuffer::new(self.0.clone()) ) );
1349 }
1350
1351 cx.waker().wake_by_ref();
1352
1353 return Poll::Pending;
1354 }
1355 )
1356 .await;
1357 }
1358
1359 #[cfg(test)]
1360 fn get_flags(&self) -> RwBufferFlags<Self>
1361 {
1362 let inner = self.inner();
1363
1364 let flags: RwBufferFlags<Self> = inner.flags.load(Ordering::SeqCst).into();
1365
1366 return flags;
1367 }
1368
1369 fn clone_single(&self) -> RwBufferRes<Self>
1376 {
1377 let inner = self.inner();
1378
1379 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
1380
1381 current_flags.base()?;
1382
1383 inner.flags.store(current_flags.into(), Ordering::Relaxed);
1384
1385 return Ok(Self(self.0));
1386 }
1387}
1388
1389impl Clone for RwBuffer
1390{
1391 fn clone(&self) -> Self
1395 {
1396 let inner = self.inner();
1397
1398 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
1399 let mut new_flags = current_flags.clone();
1400
1401 new_flags.base().unwrap();
1402
1403 let backoff = Backoff::new();
1404
1405 #[cfg(all(feature = "std", not(feature = "clone_wait_indef")))]
1406 let mut parked = false;
1407
1408 loop
1409 {
1410 let res =
1411 inner
1412 .flags
1413 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1414
1415 if let Ok(_) = res
1416 {
1417 return Self(self.0);
1418 }
1419
1420 current_flags = res.err().unwrap().into();
1421 new_flags = current_flags.clone();
1422
1423 new_flags.base().unwrap();
1424
1425 if backoff.is_completed() == false
1426 {
1427 backoff.snooze();
1428 }
1429 else
1430 {
1431 #[cfg(all(feature = "std", not(feature = "clone_wait_indef")))]
1432 {
1433 if parked == false
1434 {
1435 std::thread::park_timeout(Duration::from_millis(1));
1437
1438 parked = true;
1439 }
1440 else
1441 {
1442 panic!("can not obtain a clone of RBuffer in reasonable time!");
1443 }
1444 }
1445
1446 #[cfg(all(not(feature = "std"), not(feature = "clone_wait_indef")))]
1447 {
1448 panic!("can not obtain a clone of RBuffer in reasonable time!");
1449 }
1450
1451 #[cfg(feature = "clone_wait_indef")]
1452 {
1453 backoff.reset();
1454 }
1455 }
1456 }
1457 }
1458}
1459
1460impl TryClone for RwBuffer
1461{
1462 type Error = RwBufferError;
1463
1464 fn try_clone(&self) -> Result<Self, Self::Error>
1476 {
1477 let inner = self.inner();
1478
1479 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::SeqCst).into();
1480 let mut new_flags = current_flags.clone();
1481
1482 new_flags.base()?;
1483
1484 let backoff = Backoff::new();
1485
1486 while backoff.is_completed() == false
1487 {
1488 let res =
1489 inner
1490 .flags
1491 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1492
1493 if let Ok(_) = res
1494 {
1495 return Ok(Self(self.0));
1496 }
1497
1498 current_flags = res.err().unwrap().into();
1499 new_flags = current_flags.clone();
1500
1501 new_flags.base()?;
1502
1503 backoff.snooze();
1504 }
1505
1506 return Err(RwBufferError::BaseTryAgainLater);
1507 }
1508}
1509
1510impl Drop for RwBuffer
1511{
1512 fn drop(&mut self)
1517 {
1518 let inner = self.inner();
1519
1520 let mut current_flags: RwBufferFlags<Self> = inner.flags.load(Ordering::Relaxed).into();
1521 let mut new_flags = current_flags.clone();
1522
1523 new_flags.unbase();
1524
1525 let backoff = Backoff::new();
1526
1527 for _ in 0..1000
1528 {
1529 let res =
1530 inner
1531 .flags
1532 .compare_exchange_weak(current_flags.into(), new_flags.into(), Ordering::Acquire, Ordering::Relaxed);
1533
1534 if let Ok(flags) = res.map(|v| <u64 as Into<RwBufferFlags<Self>>>::into(v))
1535 {
1536 if flags.is_drop_inplace() == true
1537 {
1538 unsafe { ptr::drop_in_place(self.0.as_ptr()) };
1540 }
1541
1542 return;
1543 }
1544
1545 current_flags = res.err().unwrap().into();
1546 new_flags = current_flags.clone();
1547
1548 new_flags.unbase();
1549
1550 backoff.snooze();
1551 }
1552
1553 panic!("assertion trap: RwBuffer::drop can not drop RwBuffer in reasonable time!");
1555 }
1556}
1557
1558#[derive(Debug)]
1562pub struct RwBuffers
1563{
1564 buf_len: usize,
1566
1567 bufs_cnt_lim: usize,
1569
1570 buffs: VecDeque<RwBuffer>
1572
1573}
1574
1575impl RwBuffers
1576{
1577 pub
1601 fn new(buf_len: usize, pre_init_cnt: usize, bufs_cnt_lim: usize) -> RwBufferRes<Self>
1602 {
1603 if pre_init_cnt > bufs_cnt_lim
1604 {
1605 return Err(RwBufferError::InvalidArguments);
1606 }
1607 else if buf_len == 0
1608 {
1609 return Err(RwBufferError::InvalidArguments);
1610 }
1611
1612 let buffs: VecDeque<RwBuffer> =
1613 if pre_init_cnt > 0
1614 {
1615 let mut buffs = VecDeque::with_capacity(bufs_cnt_lim);
1616
1617 for _ in 0..pre_init_cnt
1618 {
1619 buffs.push_back(RwBuffer::new(buf_len));
1620 }
1621
1622 buffs
1623 }
1624 else
1625 {
1626 VecDeque::with_capacity(bufs_cnt_lim)
1627 };
1628
1629 return Ok(
1630 Self
1631 {
1632 buf_len: buf_len,
1633 bufs_cnt_lim: bufs_cnt_lim,
1634 buffs: buffs,
1635 }
1636 )
1637 }
1638
1639 pub
1652 fn new_unbounded(buf_len: usize, pre_init_cnt: usize) -> Self
1653 {
1654 let mut buffs = VecDeque::with_capacity(pre_init_cnt);
1655
1656 for _ in 0..pre_init_cnt
1657 {
1658 buffs.push_back(RwBuffer::new(buf_len));
1659 }
1660
1661 return
1662 Self
1663 {
1664 buf_len: buf_len,
1665 bufs_cnt_lim: 0,
1666 buffs: buffs,
1667 };
1668 }
1669
1670 pub
1688 fn allocate(&mut self) -> RwBufferRes<RwBuffer>
1689 {
1690 for buf in self.buffs.iter()
1692 {
1693 if let Ok(rwbuf) = buf.acqiure_if_free()
1694 {
1695 return Ok(rwbuf);
1696 }
1697 }
1698
1699 if self.bufs_cnt_lim == 0 || self.buffs.len() < self.bufs_cnt_lim
1700 {
1701 let buf = RwBuffer::new(self.buf_len);
1702 let c_buf = buf.clone_single()?;
1703
1704 self.buffs.push_back(buf);
1705
1706 return Ok(c_buf);
1707 }
1708
1709 return Err(RwBufferError::OutOfBuffers);
1710 }
1711
1712 pub
1720 fn allocate_in_place(&mut self) -> RwBuffer
1721 {
1722 let mut idx = Option::None;
1723
1724 for (i, _item) in self.buffs.iter().enumerate()
1725 {
1726 if let Ok(_) = self.buffs[i].acqiure_if_free()
1727 {
1728 idx = Some(i);
1729
1730 break;
1731 }
1732 }
1733
1734 return
1735 idx
1736 .map_or(
1737 RwBuffer::new(self.buf_len),
1738 |f| self.buffs.remove(f).unwrap()
1739 );
1740
1741 }
1742
1743 pub
1756 fn compact(&mut self, mut cnt: usize) -> usize
1757 {
1758 let p_cnt = cnt;
1759
1760 self
1761 .buffs
1762 .retain(
1763 |buf|
1764 {
1765 if buf.is_free() == true
1766 {
1767 cnt -= 1;
1768
1769 return false;
1770 }
1771
1772 return true;
1773 }
1774 );
1775
1776 return p_cnt - cnt;
1777 }
1778
1779 #[cfg(test)]
1780 fn get_flags_by_index(&self, index: usize) -> Option<RwBufferFlags<RwBuffer>>
1781 {
1782 return Some(self.buffs.get(index)?.get_flags());
1783 }
1784}
1785
1786
1787#[cfg(test)]
1788mod tests;