1use std::fmt;
4use std::sync::atomic::{AtomicBool, Ordering};
5use std::sync::{Arc, Mutex};
6use std::task::{Poll, ready};
7
8use crate::datagram::{Datagram, Event};
9use crate::error::{Error, Result};
10use crate::group;
11use crate::key::TrackKey;
12use crate::limits::{MAX_DATAGRAM_PAYLOAD, check_u53};
13use crate::window::DatagramWindow;
14
15pub struct Producer {
20 inner: moq_net::track::Producer,
21 group_key: Arc<Mutex<TrackKey>>,
22 datagram_key: TrackKey,
23 next: u64,
24}
25
26impl Producer {
27 pub(crate) fn new(inner: moq_net::track::Producer, group_key: TrackKey, datagram_key: TrackKey) -> Self {
28 Self {
29 inner,
30 group_key: Arc::new(Mutex::new(group_key)),
31 datagram_key,
32 next: 0,
33 }
34 }
35
36 pub fn name(&self) -> &str {
38 self.inner.name()
39 }
40
41 pub fn append_group(&mut self) -> Result<group::Producer> {
47 let sequence = self.next;
48 self.create_group(sequence)
49 }
50
51 pub fn create_group(&mut self, sequence: u64) -> Result<group::Producer> {
58 self.allocate(sequence)?;
59 let inner = self.inner.create_group(moq_net::group::Info { sequence })?;
60 Ok(group::Producer::new(inner, self.group_key.clone()))
61 }
62
63 pub fn append_datagram(&mut self, timestamp: moq_net::Timestamp, plaintext: &[u8]) -> Result<u64> {
69 let sequence = self.next;
70 self.insert_datagram(sequence, timestamp, plaintext)?;
71 Ok(sequence)
72 }
73
74 pub fn insert_datagram(&mut self, sequence: u64, timestamp: moq_net::Timestamp, plaintext: &[u8]) -> Result<()> {
85 self.reserve(sequence)?;
86 let payload = self
87 .datagram_key
88 .protect(sequence, 0, plaintext, MAX_DATAGRAM_PAYLOAD)?;
89 self.next = sequence + 1;
90 self.inner.insert_datagram(sequence, timestamp, payload)?;
91 Ok(())
92 }
93
94 pub fn finish(self) -> Result<()> {
100 self.inner.finish()?;
101 Ok(())
102 }
103
104 pub fn abort(self) -> Result<()> {
110 self.inner.abort(moq_net::Error::Cancel)?;
111 Ok(())
112 }
113
114 fn reserve(&self, sequence: u64) -> Result<()> {
116 if sequence < self.next {
117 return Err(Error::Reuse);
118 }
119 check_u53(sequence)
120 }
121
122 fn allocate(&mut self, sequence: u64) -> Result<()> {
123 self.reserve(sequence)?;
124 self.next = sequence + 1;
125 Ok(())
126 }
127
128 #[cfg(test)]
129 pub(crate) fn datagram_invocations(&self) -> u64 {
130 self.datagram_key.invocations()
131 }
132}
133
134impl fmt::Debug for Producer {
135 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
136 f.debug_struct("track::Producer")
137 .field("name", &self.name())
138 .field("next", &self.next)
139 .finish()
140 }
141}
142
143pub struct Consumer {
148 inner: moq_net::track::Subscriber,
149 group_key: Arc<Mutex<TrackKey>>,
150 datagram_key: TrackKey,
151 window: DatagramWindow,
152 auth_failed: Arc<AtomicBool>,
153}
154
155impl Consumer {
156 pub(crate) fn new(inner: moq_net::track::Subscriber, group_key: TrackKey, datagram_key: TrackKey) -> Self {
157 Self {
158 inner,
159 group_key: Arc::new(Mutex::new(group_key)),
160 datagram_key,
161 window: DatagramWindow::default(),
162 auth_failed: Arc::new(AtomicBool::new(false)),
163 }
164 }
165
166 pub fn name(&self) -> &str {
168 self.inner.name()
169 }
170
171 pub fn poll_recv_group(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<group::Consumer>>> {
177 if self.auth_failed.load(Ordering::Acquire) {
178 return Poll::Ready(Err(Error::Authentication));
179 }
180 let Some(inner) = ready!(self.inner.poll_recv_group(waiter)?) else {
181 return Poll::Ready(Ok(None));
182 };
183 Poll::Ready(Ok(Some(group::Consumer::new(
184 inner,
185 self.group_key.clone(),
186 self.auth_failed.clone(),
187 ))))
188 }
189
190 pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
196 kio::wait(|waiter| self.poll_recv_group(waiter)).await
197 }
198
199 pub fn poll_recv_datagram(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Event>>> {
208 if self.auth_failed.load(Ordering::Acquire) {
209 return Poll::Ready(Err(Error::Authentication));
210 }
211 let Some(datagram) = ready!(self.inner.poll_recv_datagram(waiter)?) else {
212 return Poll::Ready(Ok(None));
213 };
214 let sequence = datagram.sequence;
215 if self.window.is_duplicate(sequence) {
216 return Poll::Ready(Ok(Some(Event::Duplicate { sequence })));
217 }
218 match self
219 .datagram_key
220 .open(sequence, 0, &datagram.payload, MAX_DATAGRAM_PAYLOAD)
221 {
222 Ok(plaintext) => {
223 self.window.mark(sequence);
224 Poll::Ready(Ok(Some(Event::Datagram(Datagram {
225 sequence,
226 timestamp: datagram.timestamp,
227 plaintext,
228 }))))
229 }
230 Err(Error::Authentication) => Poll::Ready(Ok(Some(Event::Authentication { sequence }))),
231 Err(err) => Poll::Ready(Err(err)),
232 }
233 }
234
235 pub async fn recv_datagram(&mut self) -> Result<Option<Event>> {
241 kio::wait(|waiter| self.poll_recv_datagram(waiter)).await
242 }
243}
244
245impl fmt::Debug for Consumer {
246 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
247 f.debug_struct("track::Consumer")
248 .field("name", &self.name())
249 .field("auth_failed", &self.auth_failed.load(Ordering::Acquire))
250 .finish()
251 }
252}