1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
use super::{Inner, COMPLETE, NUMBER_BITS, NUMBER_MASK};
use crate::sync::spsc::{SpscInner, SpscInnerErr};
use alloc::sync::Arc;
use core::{
pin::Pin,
ptr,
sync::atomic::Ordering,
task::{Context, Poll},
};
use futures::stream::Stream;
const IS_TX_HALF: bool = false;
#[must_use = "futures do nothing unless you `.await` or poll them"]
pub struct Receiver<T, E> {
inner: Arc<Inner<T, E>>,
}
impl<T, E> Receiver<T, E> {
pub(super) fn new(inner: Arc<Inner<T, E>>) -> Self {
Self { inner }
}
#[inline]
pub fn close(&mut self) {
self.inner.close_half(IS_TX_HALF)
}
#[inline]
pub fn try_next(&mut self) -> Result<Option<T>, E> {
self.inner.try_next()
}
}
impl<T, E> Stream for Receiver<T, E> {
type Item = Result<T, E>;
#[inline]
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
self.inner.poll_half_with_transaction(
cx,
IS_TX_HALF,
Ordering::Acquire,
Ordering::AcqRel,
Inner::take_index_try,
Inner::take_index_finalize,
)
}
}
impl<T, E> Drop for Receiver<T, E> {
#[inline]
fn drop(&mut self) {
self.inner.close_half(IS_TX_HALF);
}
}
impl<T, E> Inner<T, E> {
pub(super) fn take_index(&self, state: &mut usize, length: usize) -> usize {
let cursor = *state >> NUMBER_BITS & NUMBER_MASK;
*state >>= NUMBER_BITS << 1;
*state <<= NUMBER_BITS;
*state |= cursor.wrapping_add(1).wrapping_rem(self.buffer.capacity());
*state <<= NUMBER_BITS;
*state |= length.wrapping_sub(1);
cursor
}
pub(super) fn get_length(state: usize) -> usize {
state & NUMBER_MASK
}
fn try_next(&self) -> Result<Option<T>, E> {
let state = self.state_load(Ordering::Acquire);
self.transaction(state, Ordering::AcqRel, Ordering::Acquire, |state| {
match self.take_index_try(state) {
Some(value) => value.map_err(Ok),
None => Err(Err(())),
}
})
.map(|index| unsafe { Some(self.take_value(index)) })
.or_else(|value| value.map_or_else(|()| Ok(None), |()| self.take_err().transpose()))
}
fn take_index_try(&self, state: &mut usize) -> Option<Result<usize, ()>> {
let length = Self::get_length(*state);
if length != 0 {
Some(Ok(self.take_index(state, length)))
} else if *state & COMPLETE == 0 {
None
} else {
Some(Err(()))
}
}
fn take_index_finalize(&self, value: Result<usize, ()>) -> Option<Result<T, E>> {
match value {
Ok(index) => unsafe { Some(Ok(self.take_value(index))) },
Err(()) => self.take_err(),
}
}
unsafe fn take_value(&self, index: usize) -> T {
unsafe { ptr::read(self.buffer.ptr().add(index)) }
}
}