Skip to main content

shuttle_std/sync/
mpsc.rs

1//! Multi-producer, single-consumer FIFO queue communication primitives.
2
3use crate::sync::{ResourceSignature, ResourceType};
4use shuttle_engine::runtime::execution::ExecutionState;
5use shuttle_engine::runtime::task::clock::VectorClock;
6use shuttle_engine::runtime::task::{TaskId, DEFAULT_INLINE_TASKS};
7use shuttle_engine::runtime::thread;
8use smallvec::SmallVec;
9use std::cell::RefCell;
10use std::fmt::Debug;
11use std::rc::Rc;
12use std::result::Result;
13pub use std::sync::mpsc::{RecvError, RecvTimeoutError, SendError, TryRecvError, TrySendError};
14use std::sync::Arc;
15use std::time::Duration;
16use tracing::trace;
17
18const MAX_INLINE_MESSAGES: usize = 32;
19
20/// Create an unbounded channel
21#[track_caller]
22pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
23    let channel = Arc::new(Channel::new(None));
24    let sender = Sender {
25        inner: Arc::clone(&channel),
26    };
27    let receiver = Receiver {
28        inner: Arc::clone(&channel),
29    };
30    (sender, receiver)
31}
32
33/// Create a bounded channel
34#[track_caller]
35pub fn sync_channel<T>(bound: usize) -> (SyncSender<T>, Receiver<T>) {
36    let channel = Arc::new(Channel::new(Some(bound)));
37    let sender = SyncSender {
38        inner: Arc::clone(&channel),
39    };
40    let receiver = Receiver {
41        inner: Arc::clone(&channel),
42    };
43    (sender, receiver)
44}
45
46#[derive(Debug)]
47struct Channel<T> {
48    bound: Option<usize>, // None for an unbounded channel, Some(k) for a bounded channel of size k
49    state: Rc<RefCell<ChannelState<T>>>,
50    #[allow(unused)]
51    signature: ResourceSignature,
52}
53
54// For tracking causality on channels, we timestamp each message with the clock of the sender.
55// When the receiver gets the message, it updates its clock with the the associated timestamp.
56// For unbounded channels, that's all the work we need to do.
57//
58// For bounded and rendezvous channels, things get a bit more interesting.
59// Consider a bounded channel of depth K.  As soon as the sender successfully sends its K+1'th
60// message, it knows that the receiver has received at least 1 message.  At this point, the
61// first receive event causally precedes the (K+1)'th send.  By the rule for vector clocks,
62//  (clock of the first receive)  <  (clock of the K+1'th send)
63// In order to ensure this ordering, we add a return queue of depth K to bounded channels.
64// Initially, this queue contains K empty vector clocks.  On each receive, we push the
65// receiver's clock at the time of the receive to the end of this queue.  Whenever the sender
66// successfully sends a message, it pops the clock at the front of the queue, and updates its
67// own clock with this value.  Thus, on the (K+1)'th send, the sender's clock will be updated
68// with the clock at the first receive, as needed.
69//
70// The story is similar for rendezvous channels, except we have to handle things a bit more
71// specially because K=0.
72
73struct TimestampedValue<T> {
74    value: T,
75    clock: VectorClock,
76}
77
78impl<T> TimestampedValue<T> {
79    fn new(value: T, clock: VectorClock) -> Self {
80        Self { value, clock }
81    }
82}
83
84// Note: The channels in std::sync::mpsc only support a single Receiver (which cannot be
85// cloned).  The state below admits a more general use case, where multiple Senders
86// and Receivers can share a single channel.
87struct ChannelState<T> {
88    messages: SmallVec<[TimestampedValue<T>; MAX_INLINE_MESSAGES]>, // messages in the channel
89    receiver_clock: Option<SmallVec<[VectorClock; MAX_INLINE_MESSAGES]>>, // receiver vector clocks for bounded case
90    known_senders: usize,                                           // number of senders referencing this channel
91    known_receivers: usize,                                         // number or receivers referencing this channel
92    waiting_senders: SmallVec<[TaskId; DEFAULT_INLINE_TASKS]>,      // list of currently blocked senders
93    waiting_receivers: SmallVec<[TaskId; DEFAULT_INLINE_TASKS]>,    // list of currently blocked receivers
94}
95
96impl<T> Debug for ChannelState<T> {
97    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
98        write!(f, "Channel {{ ")?;
99        write!(f, "num_messages: {} ", self.messages.len())?;
100        write!(
101            f,
102            "known_senders {} known_receivers {} ",
103            self.known_senders, self.known_receivers
104        )?;
105        write!(f, "waiting_senders: [{:?}] ", self.waiting_senders)?;
106        write!(f, "waiting_receivers: [{:?}] ", self.waiting_receivers)?;
107        write!(f, "}}")
108    }
109}
110
111/// Which list of waiters a task blocked on a channel is on.
112enum Waiter {
113    Sender,
114    Receiver,
115}
116
117/// Takes a task blocked on a channel off its list of waiters, if execution teardown unwinds the task
118/// from the switch where it blocked (see `ExecutionState::tear_down`). If the task was first in line,
119/// and could have gone on, the next one in line can go on instead.
120struct StopWaitingOnUnwind<'a, T>(&'a Channel<T>, TaskId, Waiter);
121
122impl<T> Drop for StopWaitingOnUnwind<'_, T> {
123    fn drop(&mut self) {
124        let mut state = self.0.state.borrow_mut();
125        let state = &mut *state;
126        let waiters = match self.2 {
127            Waiter::Sender => &mut state.waiting_senders,
128            Waiter::Receiver => &mut state.waiting_receivers,
129        };
130        let Some(position) = waiters.iter().position(|task| *task == self.1) else {
131            return;
132        };
133        waiters.remove(position);
134        let Some(&next) = waiters.first().filter(|_| position == 0) else {
135            return;
136        };
137        // As `send_internal` and `recv_internal` decide whether the next waiter can go on.
138        let can_go_on = match self.2 {
139            Waiter::Sender => match self.0.bound.expect("only a bounded channel has waiting senders") {
140                0 => !state.waiting_receivers.is_empty(),
141                bound => state.messages.len() < bound,
142            },
143            Waiter::Receiver => !state.messages.is_empty(),
144        };
145        if can_go_on {
146            ExecutionState::with(|s| s.get_mut(next).unblock());
147        }
148    }
149}
150
151impl<T> Channel<T> {
152    #[track_caller]
153    fn new(bound: Option<usize>) -> Self {
154        let receiver_clock = if let Some(bound) = bound {
155            let mut s = SmallVec::with_capacity(bound);
156            for _ in 0..bound {
157                s.push(VectorClock::new());
158            }
159            Some(s)
160        } else {
161            None
162        };
163        Self {
164            bound,
165            state: Rc::new(RefCell::new(ChannelState {
166                messages: SmallVec::new(),
167                receiver_clock,
168                known_senders: 1,
169                known_receivers: 1,
170                waiting_senders: SmallVec::new(),
171                waiting_receivers: SmallVec::new(),
172            })),
173            signature: ExecutionState::new_resource_signature(ResourceType::MpscChannel),
174        }
175    }
176
177    fn try_send(&self, message: T) -> Result<(), TrySendError<T>> {
178        self.send_internal(message, false)
179    }
180
181    fn send(&self, message: T) -> Result<(), SendError<T>> {
182        self.send_internal(message, true).map_err(|e| match e {
183            TrySendError::Full(_) => unreachable!(),
184            TrySendError::Disconnected(m) => SendError(m),
185        })
186    }
187
188    fn is_rendezvous(&self) -> bool {
189        self.bound == Some(0)
190    }
191
192    fn sender_must_block(&self, state: &ChannelState<T>) -> bool {
193        let (is_rendezvous, is_full) = if let Some(bound) = self.bound {
194            // For a rendezvous channel (bound = 0), "is_full" holds when there is a message in the channel.
195            // For a non-rendezvous channel (bound > 0), "is_full" holds when the capacity is reached.
196            // We cover both these cases at once using max(bound, 1) below.
197            (bound == 0, state.messages.len() >= std::cmp::max(bound, 1))
198        } else {
199            (false, false)
200        };
201
202        // The sender should block in any of the following situations:
203        //    the channel is full (as defined above)
204        //    there are already waiting senders
205        //    this is a rendezvous channel and there are no waiting receivers
206        is_full || !state.waiting_senders.is_empty() || (is_rendezvous && state.waiting_receivers.is_empty())
207    }
208
209    fn send_internal(&self, message: T, can_block: bool) -> Result<(), TrySendError<T>> {
210        // Because channels are always fair wrt. waiting senders (waiting senders is an *ordered* list),
211        // blocking sends do not commute thus must always provide a switch before blocking
212        thread::switch();
213
214        let me = ExecutionState::me();
215        let mut state = self.state.borrow_mut();
216        let should_block = self.sender_must_block(&state);
217
218        trace!(
219            state = ?state,
220            "sender {:?} starting send on channel {:p}",
221            me,
222            self,
223        );
224        if state.known_receivers == 0 {
225            // No receivers are left, so the channel is disconnected.  Stop and return failure.
226            return Err(TrySendError::Disconnected(message));
227        }
228
229        if should_block {
230            if !can_block {
231                return Err(TrySendError::Full(message));
232            }
233
234            state.waiting_senders.push(me);
235            trace!(
236                state = ?state,
237                "blocking sender {:?} on channel {:p}",
238                me,
239                self,
240            );
241            ExecutionState::with(|s| s.current_mut().block(false));
242            drop(state);
243
244            // If execution teardown unwinds the task from the switch below (see
245            // `ExecutionState::tear_down`), the task is no longer waiting. Destructors that run later
246            // use this channel as if it did.
247            let stop_waiting_on_unwind = StopWaitingOnUnwind(self, me, Waiter::Sender);
248            thread::switch();
249            std::mem::forget(stop_waiting_on_unwind);
250
251            state = self.state.borrow_mut();
252            trace!(
253                state = ?state,
254                "unblocked sender {:?} on channel {:p}",
255                me,
256                self,
257            );
258
259            // Check again that we still have a receiver; if not, return with error.
260            // We repeat this check because the receivers may have disconnected while the sender was blocked.
261            if state.known_receivers == 0 {
262                state.waiting_senders.retain(|t| *t != me);
263                // No receivers are left, so the channel is disconnected.  Stop and return failure.
264                return Err(TrySendError::Disconnected(message));
265            }
266
267            let head = state.waiting_senders.remove(0);
268            assert_eq!(head, me);
269        }
270
271        ExecutionState::with(|s| {
272            let clock = s.increment_clock();
273            state.messages.push(TimestampedValue::new(message, clock.clone()));
274        });
275
276        // The sender has just added a message to the channel, so unblock the first waiting receiver if any
277        if let Some(&tid) = state.waiting_receivers.first() {
278            ExecutionState::with(|s| {
279                s.get_mut(tid).unblock();
280
281                // When a sender successfully sends on a rendezvous channel, it knows that the receiver will perform
282                // the matching receive, so we need to update the sender's clock with the receiver's.
283                if self.is_rendezvous() {
284                    let recv_clock = s.get_clock(tid).clone();
285                    s.update_clock(&recv_clock);
286                }
287            });
288        }
289        // Check and unblock the next the waiting sender, if eligible
290        if let Some(&tid) = state.waiting_senders.first() {
291            let bound = self.bound.expect("can't have waiting senders on an unbounded channel");
292            if state.messages.len() < bound {
293                ExecutionState::with(|s| s.get_mut(tid).unblock());
294            }
295        }
296
297        if !self.is_rendezvous() {
298            if let Some(receiver_clock) = &mut state.receiver_clock {
299                let recv_clock = receiver_clock.remove(0);
300                ExecutionState::with(|s| s.update_clock(&recv_clock));
301            }
302        }
303
304        Ok(())
305    }
306
307    fn recv(&self) -> Result<T, RecvError> {
308        self.recv_internal(true).map_err(|e| match e {
309            TryRecvError::Disconnected => RecvError,
310            TryRecvError::Empty => unreachable!(),
311        })
312    }
313
314    fn try_recv(&self) -> Result<T, TryRecvError> {
315        self.recv_internal(false)
316    }
317
318    fn receiver_must_block(&self, state: &ChannelState<T>) -> bool {
319        // The receiver should block in any of the following situations:
320        //    the channel is empty
321        //    there are waiting receivers
322        state.messages.is_empty() || !state.waiting_receivers.is_empty()
323    }
324
325    fn recv_internal(&self, can_block: bool) -> Result<T, TryRecvError> {
326        // Because channels are always fair wrt. waiting receivers (waiting receivers is an *ordered* list),
327        // blocking receives do not commute thus must always provide a switch before blocking
328        thread::switch();
329
330        let me = ExecutionState::me();
331        let mut state = self.state.borrow_mut();
332        let should_block = self.receiver_must_block(&state);
333
334        trace!(
335            state = ?state,
336            "starting recv on channel {:p}",
337            self,
338        );
339        // Check if there are any senders left; if not, and the channel is empty, fail with error
340        // (If there are no senders, but the channel is nonempty, the receiver can successfully consume that message.)
341        if state.messages.is_empty() && state.known_senders == 0 {
342            return Err(TryRecvError::Disconnected);
343        }
344
345        // If this is a rendezvous channel, and the channel is empty, and there are waiting senders,
346        // notify the first waiting sender
347        if self.is_rendezvous() && state.messages.is_empty() {
348            if let Some(&tid) = state.waiting_senders.first() {
349                // Note: another receiver may have unblocked the sender already
350                ExecutionState::with(|s| s.get_mut(tid).unblock());
351            } else if !can_block {
352                // Nobody to rendezvous with
353                return Err(TryRecvError::Empty);
354            }
355        }
356
357        // Handle the try_recv case, accounting for the number of msgs available and already waiting receivers.
358        if !self.is_rendezvous() && !can_block && state.waiting_receivers.len() >= state.messages.len() {
359            return Err(TryRecvError::Empty);
360        }
361
362        // Pre-increment the receiver's clock before continuing
363        //
364        // Note: The reason for pre-incrementing the receiver's clock is to deal properly with rendezvous channels.
365        // Here's the scenario we have to handle:
366        //   1. the receiver arrives at a rendezvous channel and blocks
367        //   2. the sender arrives, sees the receiver is waiting and does not block
368        //   3. the sender drops the message in the channel and updates its clock with the receiver's clock and continues
369        //   4. later, the receiver unblocks and picks up the message and updates its clock with the sender's
370        // Without the pre-increment, in step 3, the sender would update its clock with the receiver's clock before
371        // it is incremented.  (The increment records the fact that the receiver arrived at the synchronization point.)
372        ExecutionState::with(|s| {
373            let _ = s.increment_clock();
374        });
375
376        if should_block {
377            state.waiting_receivers.push(me);
378            trace!(
379                state = ?state,
380                "blocking receiver {:?} on channel {:p}",
381                me,
382                self,
383            );
384            ExecutionState::with(|s| s.current_mut().block(false));
385            drop(state);
386
387            // As for a blocked sender.
388            let stop_waiting_on_unwind = StopWaitingOnUnwind(self, me, Waiter::Receiver);
389            thread::switch();
390            std::mem::forget(stop_waiting_on_unwind);
391
392            state = self.state.borrow_mut();
393            trace!(
394                state = ?state,
395                "unblocked receiver {:?} on channel {:p}",
396                me,
397                self,
398            );
399
400            // Check again if there are any senders left; if not, and the channel is empty, fail with error
401            // (If there are no senders, but the channel is nonempty, the receiver can successfully consume that message.)
402            // We repeat this check because the senders may have disconnected while the receiver was blocked.
403            if state.messages.is_empty() && state.known_senders == 0 {
404                state.waiting_receivers.retain(|t| *t != me);
405                return Err(TryRecvError::Disconnected);
406            }
407
408            let head = state.waiting_receivers.remove(0);
409            assert_eq!(head, me);
410        }
411
412        let item = state.messages.remove(0);
413        // The receiver has just removed an element from the channel.  Check if any waiting senders
414        // need to be notified.
415        if let Some(&tid) = state.waiting_senders.first() {
416            let bound = self.bound.expect("can't have waiting senders on an unbounded channel");
417            // Unblock the first waiting sender provided one of the following conditions hold:
418            // - this is a non-rendezvous bounded channel (bound > 0)
419            // - this is a rendezvous channel and we have additional waiting receivers
420            if bound > 0 || !state.waiting_receivers.is_empty() {
421                ExecutionState::with(|s| s.get_mut(tid).unblock());
422            }
423        }
424        // Check and unblock the next the waiting receiver, if eligible
425        // Note: this is a no-op for mpsc channels, since there can only be one receiver
426        if let Some(&tid) = state.waiting_receivers.first() {
427            if !state.messages.is_empty() {
428                ExecutionState::with(|s| s.get_mut(tid).unblock());
429            }
430        }
431
432        // Update receiver clock from the clock attached to the message received
433        let TimestampedValue { value, clock } = item;
434        ExecutionState::with(|s| {
435            // Since we already incremented the receiver's clock above, just update it here
436            s.get_clock_mut(me).update(&clock);
437
438            // If this is a (non-rendezvous) bounded channel, propagate causality backwards to sender
439            if let Some(receiver_clock) = &mut state.receiver_clock {
440                let bound = self.bound.expect("unexpected internal error"); // must be defined for bounded channels
441                if bound > 0 {
442                    // non-rendezvous
443                    assert!(receiver_clock.len() < bound);
444                    receiver_clock.push(s.get_clock(me).clone());
445                }
446            }
447        });
448        Ok(value)
449    }
450}
451
452// Safety: A Channel is never actually passed across true threads, only across continuations. The
453// Rc<RefCell<_>> type therefore can't be preempted mid-bookkeeping-operation.
454// TODO We use this workaround in several places in Shuttle.  Maybe there's a cleaner solution.
455unsafe impl<T: Send> Send for Channel<T> {}
456unsafe impl<T: Send> Sync for Channel<T> {}
457
458/// The receiving half of Rust's [`channel`] (or [`sync_channel`]) type.
459/// This half can only be owned by one thread.
460#[derive(Debug)]
461pub struct Receiver<T> {
462    inner: Arc<Channel<T>>,
463}
464
465impl<T> Receiver<T> {
466    /// Attempts to wait for a value on this receiver, returning an error if the
467    /// corresponding channel has hung up.
468    pub fn recv(&self) -> Result<T, RecvError> {
469        self.inner.recv()
470    }
471
472    /// Attempts to wait for a value on this receiver, returning an error if the
473    /// corresponding channel has hung up.
474    pub fn try_recv(&self) -> Result<T, TryRecvError> {
475        self.inner.try_recv()
476    }
477
478    /// Attempts to wait for a value on this receiver, returning an error if the
479    /// corresponding channel has hung up, or if it waits more than timeout.
480    pub fn recv_timeout(&self, _timeout: Duration) -> Result<T, RecvTimeoutError> {
481        // TODO support the timeout case -- this method never times out
482        self.inner.recv().map_err(|_| RecvTimeoutError::Disconnected)
483    }
484
485    /// Returns an iterator that will block waiting for messages, but never
486    /// [`panic!`]. It will return [`None`] when the channel has hung up.
487    pub fn iter(&self) -> Iter<'_, T> {
488        Iter { rx: self }
489    }
490
491    /// Returns an iterator that will attempt to yield all pending values.
492    /// It will return `None` if there are no more pending values or if the
493    /// channel has hung up. The iterator will never [`panic!`] or block the
494    /// user by waiting for values.
495    pub fn try_iter(&self) -> TryIter<'_, T> {
496        TryIter { rx: self }
497    }
498}
499
500impl<T> Drop for Receiver<T> {
501    fn drop(&mut self) {
502        if ExecutionState::should_stop() {
503            return;
504        }
505        let mut state = self.inner.state.borrow_mut();
506        assert!(state.known_receivers > 0);
507        state.known_receivers -= 1;
508        if state.known_receivers == 0 {
509            // Last receiver was dropped; wake up all senders
510            for &tid in state.waiting_senders.iter() {
511                ExecutionState::with(|s| s.get_mut(tid).unblock());
512            }
513        }
514    }
515}
516
517/// An iterator over messages on a [`Receiver`], created by [`iter`].
518///
519/// This iterator will block whenever [`next`] is called,
520/// waiting for a new message, and [`None`] will be returned
521/// when the corresponding channel has hung up.
522///
523/// [`iter`]: Receiver::iter
524/// [`next`]: Iterator::next
525#[derive(Debug)]
526pub struct Iter<'a, T: 'a> {
527    rx: &'a Receiver<T>,
528}
529
530/// An iterator that attempts to yield all pending values for a [`Receiver`],
531/// created by [`try_iter`].
532///
533/// [`None`] will be returned when there are no pending values remaining or
534/// if the corresponding channel has hung up.
535///
536/// This iterator will never block the caller in order to wait for data to
537/// become available. Instead, it will return [`None`].
538///
539/// [`try_iter`]: Receiver::try_iter
540#[derive(Debug)]
541pub struct TryIter<'a, T: 'a> {
542    rx: &'a Receiver<T>,
543}
544
545/// An owning iterator over messages on a [`Receiver`],
546/// created by [`into_iter`].
547///
548/// This iterator will block whenever [`next`]
549/// is called, waiting for a new message, and [`None`] will be
550/// returned if the corresponding channel has hung up.
551///
552/// [`into_iter`]: Receiver::into_iter
553/// [`next`]: Iterator::next
554#[derive(Debug)]
555pub struct IntoIter<T> {
556    rx: Receiver<T>,
557}
558
559impl<T> Iterator for Iter<'_, T> {
560    type Item = T;
561
562    fn next(&mut self) -> Option<T> {
563        self.rx.recv().ok()
564    }
565}
566
567impl<T> Iterator for TryIter<'_, T> {
568    type Item = T;
569
570    fn next(&mut self) -> Option<T> {
571        self.rx.try_recv().ok()
572    }
573}
574
575impl<'a, T> IntoIterator for &'a Receiver<T> {
576    type Item = T;
577    type IntoIter = Iter<'a, T>;
578
579    fn into_iter(self) -> Iter<'a, T> {
580        self.iter()
581    }
582}
583
584impl<T> Iterator for IntoIter<T> {
585    type Item = T;
586    fn next(&mut self) -> Option<T> {
587        self.rx.recv().ok()
588    }
589}
590
591impl<T> IntoIterator for Receiver<T> {
592    type Item = T;
593    type IntoIter = IntoIter<T>;
594
595    fn into_iter(self) -> IntoIter<T> {
596        IntoIter { rx: self }
597    }
598}
599
600/// The sending-half of Rust's asynchronous [`channel`] type. This half can only be
601/// owned by one thread, but it can be cloned to send to other threads.
602#[derive(Debug)]
603pub struct Sender<T> {
604    inner: Arc<Channel<T>>,
605}
606
607impl<T> Sender<T> {
608    /// Attempts to send a value on this channel, returning it back if it could
609    /// not be sent.
610    pub fn send(&self, t: T) -> Result<(), SendError<T>> {
611        self.inner.send(t)
612    }
613}
614
615impl<T> Clone for Sender<T> {
616    fn clone(&self) -> Self {
617        let mut state = self.inner.state.borrow_mut();
618        state.known_senders += 1;
619        drop(state);
620        Self {
621            inner: self.inner.clone(),
622        }
623    }
624}
625
626impl<T> Drop for Sender<T> {
627    fn drop(&mut self) {
628        if ExecutionState::should_stop() {
629            return;
630        }
631        let mut state = self.inner.state.borrow_mut();
632        assert!(state.known_senders > 0);
633        state.known_senders -= 1;
634        if state.known_senders == 0 {
635            // Last sender was dropped; wake up all receivers
636            for &tid in state.waiting_receivers.iter() {
637                ExecutionState::with(|s| s.get_mut(tid).unblock());
638            }
639        }
640    }
641}
642
643/// The sending-half of Rust's synchronous [`sync_channel`] type.
644///
645/// Messages can be sent through this channel with [`SyncSender::send`] or \[`try_send`\] (TODO)
646///
647/// [`SyncSender::send`] will block if there is no space in the internal buffer.
648#[derive(Debug)]
649pub struct SyncSender<T> {
650    inner: Arc<Channel<T>>,
651}
652
653impl<T> SyncSender<T> {
654    /// Sends a value on this synchronous channel.
655    ///
656    /// This function will *block* until space in the internal buffer becomes
657    /// available or a receiver is available to hand off the message to.
658    pub fn send(&self, t: T) -> Result<(), SendError<T>> {
659        self.inner.send(t)
660    }
661
662    /// Attempts to send a value on this channel without blocking.
663    ///
664    /// This method differs from [`send`] by returning immediately if the
665    /// channel's buffer is full or no receiver is waiting to acquire some
666    /// data. Compared with [`send`], this function has two failure cases
667    /// instead of one (one for disconnection, one for a full buffer).
668    ///
669    /// [`send`]: Self::send
670    pub fn try_send(&self, t: T) -> Result<(), TrySendError<T>> {
671        self.inner.try_send(t)
672    }
673}
674
675impl<T> Clone for SyncSender<T> {
676    fn clone(&self) -> Self {
677        let mut state = self.inner.state.borrow_mut();
678        state.known_senders += 1;
679        drop(state);
680        Self {
681            inner: self.inner.clone(),
682        }
683    }
684}
685
686impl<T> Drop for SyncSender<T> {
687    fn drop(&mut self) {
688        if ExecutionState::should_stop() {
689            return;
690        }
691        let mut state = self.inner.state.borrow_mut();
692        assert!(state.known_senders > 0);
693        state.known_senders -= 1;
694        if state.known_senders == 0 {
695            // Last sender was dropped; wake up any receivers
696            for &tid in state.waiting_receivers.iter() {
697                ExecutionState::with(|s| s.get_mut(tid).unblock());
698            }
699        }
700    }
701}
702
703#[cfg(test)]
704mod tests {
705    use super::*;
706
707    #[test]
708    fn unique_resource_signature_mpsc() {
709        shuttle_schedulers::check_random(
710            || {
711                let (sender1, _) = channel::<i32>();
712                let (sender2, _) = channel::<i32>();
713                assert_ne!(sender1.inner.signature, sender2.inner.signature);
714            },
715            1,
716        );
717    }
718}