|
1 | 1 | use std::{fmt, sync::Arc}; |
2 | 2 |
|
3 | | -use crate::watch::State; |
| 3 | +use tokio::sync::broadcast; |
4 | 4 |
|
5 | 5 | use super::{ServeError, Track}; |
6 | 6 |
|
| 7 | +const DATAGRAM_CHANNEL_SIZE: usize = 1024; |
| 8 | + |
7 | 9 | pub struct Datagrams { |
8 | 10 | pub track: Arc<Track>, |
9 | 11 | } |
10 | 12 |
|
11 | 13 | impl Datagrams { |
12 | 14 | pub fn produce(self) -> (DatagramsWriter, DatagramsReader) { |
13 | | - let (writer, reader) = State::default().split(); |
| 15 | + let (tx, rx) = broadcast::channel(DATAGRAM_CHANNEL_SIZE); |
14 | 16 |
|
15 | | - let writer = DatagramsWriter::new(writer, self.track.clone()); |
16 | | - let reader = DatagramsReader::new(reader, self.track); |
| 17 | + let writer = DatagramsWriter::new(tx, self.track.clone()); |
| 18 | + let reader = DatagramsReader::new(rx, self.track); |
17 | 19 |
|
18 | 20 | (writer, reader) |
19 | 21 | } |
20 | 22 | } |
21 | 23 |
|
22 | | -struct DatagramsState { |
23 | | - // The latest datagram |
24 | | - latest: Option<Datagram>, |
25 | | - |
26 | | - // Increased each time datagram changes. |
27 | | - epoch: u64, |
28 | | - |
29 | | - // Set when the writer or all readers are dropped. |
30 | | - closed: Result<(), ServeError>, |
31 | | -} |
32 | | - |
33 | | -impl Default for DatagramsState { |
34 | | - fn default() -> Self { |
35 | | - Self { |
36 | | - latest: None, |
37 | | - epoch: 0, |
38 | | - closed: Ok(()), |
39 | | - } |
40 | | - } |
41 | | -} |
42 | | - |
43 | 24 | pub struct DatagramsWriter { |
44 | | - state: State<DatagramsState>, |
| 25 | + tx: broadcast::Sender<Datagram>, |
45 | 26 | pub track: Arc<Track>, |
46 | 27 | } |
47 | 28 |
|
48 | 29 | impl DatagramsWriter { |
49 | | - fn new(state: State<DatagramsState>, track: Arc<Track>) -> Self { |
50 | | - Self { state, track } |
| 30 | + fn new(tx: broadcast::Sender<Datagram>, track: Arc<Track>) -> Self { |
| 31 | + Self { tx, track } |
51 | 32 | } |
52 | 33 |
|
53 | 34 | pub fn write(&mut self, datagram: Datagram) -> Result<(), ServeError> { |
54 | | - let mut state = self.state.lock_mut().ok_or(ServeError::Cancel)?; |
55 | | - |
56 | | - state.latest = Some(datagram); |
57 | | - state.epoch += 1; |
58 | | - |
| 35 | + // Ignore send errors (no receivers) - datagrams are fire-and-forget |
| 36 | + let _ = self.tx.send(datagram); |
59 | 37 | Ok(()) |
60 | 38 | } |
61 | 39 |
|
62 | | - pub fn close(self, err: ServeError) -> Result<(), ServeError> { |
63 | | - let state = self.state.lock(); |
64 | | - state.closed.clone()?; |
65 | | - |
66 | | - let mut state = state.into_mut().ok_or(ServeError::Cancel)?; |
67 | | - state.closed = Err(err); |
68 | | - |
| 40 | + pub fn close(self, _err: ServeError) -> Result<(), ServeError> { |
| 41 | + // Channel closes when tx is dropped |
69 | 42 | Ok(()) |
70 | 43 | } |
71 | 44 | } |
72 | 45 |
|
73 | | -#[derive(Clone)] |
74 | 46 | pub struct DatagramsReader { |
75 | | - state: State<DatagramsState>, |
| 47 | + rx: broadcast::Receiver<Datagram>, |
76 | 48 | pub track: Arc<Track>, |
| 49 | + latest: Option<(u64, u64)>, |
| 50 | +} |
77 | 51 |
|
78 | | - epoch: u64, |
| 52 | +impl Clone for DatagramsReader { |
| 53 | + fn clone(&self) -> Self { |
| 54 | + Self { |
| 55 | + rx: self.rx.resubscribe(), |
| 56 | + track: self.track.clone(), |
| 57 | + latest: self.latest, |
| 58 | + } |
| 59 | + } |
79 | 60 | } |
80 | 61 |
|
81 | 62 | impl DatagramsReader { |
82 | | - fn new(state: State<DatagramsState>, track: Arc<Track>) -> Self { |
| 63 | + fn new(rx: broadcast::Receiver<Datagram>, track: Arc<Track>) -> Self { |
83 | 64 | Self { |
84 | | - state, |
| 65 | + rx, |
85 | 66 | track, |
86 | | - epoch: 0, |
| 67 | + latest: None, |
87 | 68 | } |
88 | 69 | } |
89 | 70 |
|
90 | 71 | pub async fn read(&mut self) -> Result<Option<Datagram>, ServeError> { |
91 | 72 | loop { |
92 | | - { |
93 | | - let state = self.state.lock(); |
94 | | - if self.epoch < state.epoch { |
95 | | - self.epoch = state.epoch; |
96 | | - return Ok(state.latest.clone()); |
| 73 | + match self.rx.recv().await { |
| 74 | + Ok(datagram) => { |
| 75 | + self.latest = Some((datagram.group_id, datagram.object_id)); |
| 76 | + return Ok(Some(datagram)); |
97 | 77 | } |
98 | | - |
99 | | - state.closed.clone()?; |
100 | | - match state.modified() { |
101 | | - Some(notify) => notify, |
102 | | - None => return Ok(None), // No more updates will come |
| 78 | + Err(broadcast::error::RecvError::Lagged(n)) => { |
| 79 | + log::warn!("[DATAGRAMS] reader lagged by {} datagrams", n); |
| 80 | + // Continue reading - we'll get the next available datagram |
| 81 | + } |
| 82 | + Err(broadcast::error::RecvError::Closed) => { |
| 83 | + return Ok(None); // Channel closed |
103 | 84 | } |
104 | 85 | } |
105 | | - .await; |
106 | 86 | } |
107 | 87 | } |
108 | 88 |
|
109 | | - // Returns the largest group/sequence |
110 | 89 | pub fn latest(&self) -> Option<(u64, u64)> { |
111 | | - let state = self.state.lock(); |
112 | | - state |
113 | | - .latest |
114 | | - .as_ref() |
115 | | - .map(|datagram| (datagram.group_id, datagram.object_id)) |
| 90 | + self.latest |
116 | 91 | } |
117 | 92 |
|
118 | 93 | pub fn is_closed(&self) -> bool { |
119 | | - let state = self.state.lock(); |
120 | | - state.closed.is_err() || state.modified().is_none() |
| 94 | + // Check if channel is closed by seeing if there are no more senders |
| 95 | + self.rx.len() == 0 |
121 | 96 | } |
122 | 97 | } |
123 | 98 |
|
|
0 commit comments