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}