Skip to main content

moirai_core/channel/spsc/
shared.rs

1//! The `Arc`-backed half pair: `spsc(capacity)`.
2//!
3//! Each half owns a reference-counted handle to the ring, so both are
4//! `'static` and can be moved into freshly spawned threads. The cost is one
5//! allocation for the ring plus an atomic refcount decrement per half at
6//! drop; [`borrowed`](super::borrowed) trades that for a scope.
7
8use super::ring::{SpscChannel, blocking};
9use crate::channel::error::{Channel, Result};
10use crate::channel::roles::{Consumer, Producer};
11use std::cell::Cell;
12use std::sync::Arc;
13use std::sync::atomic::Ordering;
14
15impl<T> SpscChannel<T> {
16    /// Create a channel pair (sender, receiver) for ergonomic usage
17    pub fn channel(capacity: usize) -> (SpscSender<T>, SpscReceiver<T>) {
18        let channel = Arc::new(Self::new(capacity));
19        (
20            SpscSender {
21                channel: channel.clone(),
22                cached_tail: Cell::new(0),
23            },
24            SpscReceiver {
25                channel,
26                cached_head: Cell::new(0),
27            },
28        )
29    }
30}
31
32/// Sender half of SPSC channel.
33///
34/// The single-producer half of the pair: not `Clone`, so only one exists, and
35/// not `Sync`, so it cannot be shared while it exists.
36pub struct SpscSender<T> {
37    pub(super) channel: Arc<SpscChannel<T>>,
38    /// Last known consumer index, and the reason the sender is `!Sync`.
39    ///
40    /// Reading the real `tail` on every send touches the consumer's cache line
41    /// each time, which dominates the cost of a queue that is otherwise two
42    /// loads and a store. Consulting this first means a producer that is not
43    /// hitting a full queue never reads the consumer's line at all.
44    ///
45    /// It doubles as the marker keeping the sender unshareable: `send` takes
46    /// `&self`, so a `Sync` sender could be driven by two threads at once, both
47    /// claiming the same slot. `Cell` is `Send` but not `Sync`, which permits
48    /// moving the sender to another thread while forbidding sharing it.
49    /// Replacing it with something `Sync` reopens that race, and
50    /// `halves_are_not_sync` fails if it happens.
51    cached_tail: Cell<usize>,
52}
53
54impl<T: Send> SpscSender<T> {
55    /// Send a value through the channel, blocking until there is room.
56    pub fn send(&self, value: T) -> Result<()> {
57        self.channel.send_cached(value, &self.cached_tail)
58    }
59
60    /// Try to send a value without blocking.
61    pub fn try_send(&self, value: T) -> Result<()> {
62        self.channel.try_send_cached(value, &self.cached_tail)
63    }
64}
65
66/// Receiver half of SPSC channel.
67///
68/// The single-consumer half of the pair, `!Sync` for the same reason as
69/// [`SpscSender`] — two shared receivers would `assume_init_read` one slot
70/// twice, moving out of it and then dropping it twice.
71pub struct SpscReceiver<T> {
72    pub(super) channel: Arc<SpscChannel<T>>,
73    /// Last known producer index, and the reason the receiver is `!Sync`.
74    /// Mirrors `SpscSender::cached_tail` in both roles.
75    cached_head: Cell<usize>,
76}
77
78impl<T: Send> SpscReceiver<T> {
79    /// Receive a value from the channel, blocking until one arrives.
80    pub fn recv(&self) -> Result<T> {
81        blocking(|| self.channel.try_recv_cached(&self.cached_head))
82    }
83
84    /// Try to receive a value without blocking.
85    pub fn try_recv(&self) -> Result<T> {
86        self.channel.try_recv_cached(&self.cached_head)
87    }
88}
89
90impl<T: Send> Producer<T> for SpscSender<T> {
91    #[inline]
92    fn send(&self, value: T) -> Result<()> {
93        SpscSender::send(self, value)
94    }
95
96    #[inline]
97    fn try_send(&self, value: T) -> Result<()> {
98        SpscSender::try_send(self, value)
99    }
100
101    #[inline]
102    fn is_full(&self) -> bool {
103        Channel::is_full(&*self.channel)
104    }
105
106    #[inline]
107    fn capacity(&self) -> Option<usize> {
108        Channel::capacity(&*self.channel)
109    }
110}
111
112impl<T: Send> Consumer<T> for SpscReceiver<T> {
113    #[inline]
114    fn recv(&self) -> Result<T> {
115        SpscReceiver::recv(self)
116    }
117
118    #[inline]
119    fn try_recv(&self) -> Result<T> {
120        SpscReceiver::try_recv(self)
121    }
122
123    #[inline]
124    fn is_empty(&self) -> bool {
125        Channel::is_empty(&*self.channel)
126    }
127}
128
129impl<T> Drop for SpscSender<T> {
130    fn drop(&mut self) {
131        self.channel.closed.store(true, Ordering::Release);
132    }
133}
134
135impl<T> Drop for SpscReceiver<T> {
136    fn drop(&mut self) {
137        self.channel.closed.store(true, Ordering::Release);
138    }
139}
140
141#[cfg(test)]
142mod auto_traits {
143    use super::{SpscReceiver, SpscSender};
144    use static_assertions::{assert_impl_all, assert_not_impl_any};
145
146    // Each half may be moved to the thread that owns its role.
147    assert_impl_all!(SpscSender<u64>: Send);
148    assert_impl_all!(SpscReceiver<u64>: Send);
149
150    /// Neither half may be *shared*, which is what keeps the channel SPSC.
151    ///
152    /// Both `send` and `recv` take `&self`, so a `Sync` half could be driven by
153    /// two threads at once: two producers would write the same slot, and two
154    /// consumers would read one slot twice. The `Cell` cache on each half is
155    /// what prevents it, so changing that field to a `Sync` type would reopen
156    /// the race — this assertion is what catches it.
157    #[allow(dead_code)]
158    fn halves_are_not_sync() {
159        assert_not_impl_any!(SpscSender<u64>: Sync);
160        assert_not_impl_any!(SpscReceiver<u64>: Sync);
161    }
162}