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}