Skip to main content

moq_e2ee/
track.rs

1//! Grouped-frame and datagram writers and readers for one physical track.
2
3use 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
15/// Exclusive producer for one physical track.
16///
17/// Owns the grouped-frame and datagram keys and the shared sequence namespace.
18/// Sequences are allocated monotonically; reuse and exhaustion are refused before encryption.
19pub 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	/// The physical track name.
37	pub fn name(&self) -> &str {
38		self.inner.name()
39	}
40
41	/// Allocate the next sequence and start a group there.
42	///
43	/// # Errors
44	///
45	/// [`Error::Identity`] if the next sequence exceeds `2^53-1`, or a net error.
46	pub fn append_group(&mut self) -> Result<group::Producer> {
47		let sequence = self.next;
48		self.create_group(sequence)
49	}
50
51	/// Start a group at an explicit sequence, which must be at or above the next unallocated one.
52	///
53	/// # Errors
54	///
55	/// [`Error::Reuse`] if `sequence` was already allocated, [`Error::Identity`] if
56	/// it exceeds `2^53-1`, or a net error.
57	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	/// Encrypt `plaintext` at the next sequence and insert the datagram, returning that sequence.
64	///
65	/// # Errors
66	///
67	/// [`Error::Identity`], [`Error::Exhausted`], [`Error::Oversize`], or a net write error.
68	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	/// Encrypt `plaintext` at an explicit sequence and insert the datagram.
75	///
76	/// Plaintext is capped at [`MAX_DATAGRAM_PLAINTEXT`](crate::MAX_DATAGRAM_PLAINTEXT).
77	/// A predictable failure is refused before the sequence is allocated; once
78	/// encrypted, the identity is spent even if the net write fails.
79	///
80	/// # Errors
81	///
82	/// [`Error::Reuse`] if `sequence` was already allocated, [`Error::Identity`],
83	/// [`Error::Exhausted`], [`Error::Oversize`], or a net write error.
84	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	/// Finish the track after the last allocated sequence.
95	///
96	/// # Errors
97	///
98	/// A net error if the track is already closed.
99	pub fn finish(self) -> Result<()> {
100		self.inner.finish()?;
101		Ok(())
102	}
103
104	/// Abort the track, consuming the handle.
105	///
106	/// # Errors
107	///
108	/// A net error if the track is already closed.
109	pub fn abort(self) -> Result<()> {
110		self.inner.abort(moq_net::Error::Cancel)?;
111		Ok(())
112	}
113
114	/// Refuse a sequence this track already allocated or cannot represent.
115	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
143/// Subscriber for one physical track.
144///
145/// Grouped authentication failure is sticky and ends the track. A bad datagram
146/// is dropped with [`Event::Authentication`] and the track continues.
147pub 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	/// The physical track name.
167	pub fn name(&self) -> &str {
168		self.inner.name()
169	}
170
171	/// Poll for the next protected group in arrival order.
172	///
173	/// # Errors
174	///
175	/// [`Error::Authentication`] once a grouped frame has failed to open, or a net error.
176	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	/// Receive the next protected group in arrival order.
191	///
192	/// # Errors
193	///
194	/// Same as [`Self::poll_recv_group`].
195	pub async fn recv_group(&mut self) -> Result<Option<group::Consumer>> {
196		kio::wait(|waiter| self.poll_recv_group(waiter)).await
197	}
198
199	/// Poll for the next datagram event.
200	///
201	/// A bad tag and a duplicate are events, not track-ending errors.
202	///
203	/// # Errors
204	///
205	/// [`Error::Authentication`] only if a grouped frame already failed,
206	/// [`Error::Oversize`] or [`Error::Exhausted`] from opening, or a net error.
207	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	/// Receive the next datagram event.
236	///
237	/// # Errors
238	///
239	/// Same as [`Self::poll_recv_datagram`].
240	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}