noq_proto/connection/
datagrams.rs1use std::collections::VecDeque;
2
3use bytes::Bytes;
4use thiserror::Error;
5use tracing::{debug, trace};
6
7use super::Connection;
8use crate::{
9 FrameStats, TransportError,
10 connection::PacketBuilder,
11 frame::{Datagram, FrameStruct},
12};
13
14pub struct Datagrams<'a> {
16 pub(super) conn: &'a mut Connection,
17}
18
19impl Datagrams<'_> {
20 pub fn send(&mut self, data: Bytes, drop: bool) -> Result<(), SendDatagramError> {
30 if self.conn.config.datagram_receive_buffer_size.is_none() {
31 return Err(SendDatagramError::Disabled);
32 }
33 let max = self
34 .max_size()
35 .ok_or(SendDatagramError::UnsupportedByPeer)?;
36 let send_buffer_size = self.conn.config.datagram_send_buffer_size;
37 if data.len() > Ord::min(max, send_buffer_size) {
38 return Err(SendDatagramError::TooLarge);
39 }
40 if drop {
41 self.conn
42 .datagrams
43 .make_space_for(data.len(), send_buffer_size);
44 } else if !self
45 .conn
46 .datagrams
47 .has_send_buffer_space(data.len(), send_buffer_size)
48 {
49 self.conn.datagrams.send_blocked = true;
50 return Err(SendDatagramError::Blocked(data));
51 }
52 self.conn.datagrams.outgoing_total += data.len();
53 self.conn.datagrams.outgoing.push_back(Datagram { data });
54 Ok(())
55 }
56
57 pub fn send_many(
76 &mut self,
77 datagrams: &[Bytes],
78 drop: bool,
79 ) -> Result<usize, SendDatagramError> {
80 if self.conn.config.datagram_receive_buffer_size.is_none() {
81 return Err(SendDatagramError::Disabled);
82 }
83 let max = self
84 .max_size()
85 .ok_or(SendDatagramError::UnsupportedByPeer)?;
86 let send_buffer_size = self.conn.config.datagram_send_buffer_size;
87 if datagrams
88 .iter()
89 .any(|data| data.len() > Ord::min(max, send_buffer_size))
90 {
91 return Err(SendDatagramError::TooLarge);
92 }
93
94 let mut queued = 0usize;
95 for data in datagrams {
96 if drop {
97 self.conn
98 .datagrams
99 .make_space_for(data.len(), send_buffer_size);
100 } else if !self
101 .conn
102 .datagrams
103 .has_send_buffer_space(data.len(), send_buffer_size)
104 {
105 self.conn.datagrams.send_blocked = true;
106 break;
107 }
108 self.conn.datagrams.outgoing_total += data.len();
109 self.conn
110 .datagrams
111 .outgoing
112 .push_back(Datagram { data: data.clone() });
113 queued += 1;
114 }
115
116 Ok(queued)
117 }
118
119 pub fn max_size(&self) -> Option<usize> {
132 let max_size = self.conn.current_mtu() as usize
136 - self.conn.predict_1rtt_overhead_no_pn()
137 - Datagram::SIZE_BOUND;
138 let limit = self
139 .conn
140 .peer_params
141 .max_datagram_frame_size?
142 .into_inner()
143 .saturating_sub(Datagram::SIZE_BOUND as u64);
144 Some(limit.min(max_size as u64) as usize)
145 }
146
147 pub fn recv(&mut self) -> Option<Bytes> {
149 self.conn.datagrams.recv()
150 }
151
152 pub fn recv_many(&mut self, out: &mut [Bytes]) -> usize {
160 self.conn.datagrams.recv_many(out)
161 }
162
163 pub fn send_buffer_space(&self) -> usize {
168 self.conn
169 .config
170 .datagram_send_buffer_size
171 .saturating_sub(self.conn.datagrams.outgoing_total)
172 }
173}
174
175#[derive(Default)]
176pub(super) struct DatagramState {
177 pub(super) recv_buffered: usize,
180 pub(super) incoming: VecDeque<Datagram>,
181 pub(super) outgoing: VecDeque<Datagram>,
182 pub(super) outgoing_total: usize,
183 pub(super) send_blocked: bool,
184}
185
186impl DatagramState {
187 pub(super) fn received(
188 &mut self,
189 datagram: Datagram,
190 window: &Option<usize>,
191 ) -> Result<bool, TransportError> {
192 let window = match window {
193 None => {
194 return Err(TransportError::PROTOCOL_VIOLATION(
195 "unexpected DATAGRAM frame",
196 ));
197 }
198 Some(x) => *x,
199 };
200
201 if datagram.data.len() > window {
202 return Err(TransportError::PROTOCOL_VIOLATION("oversized datagram"));
203 }
204
205 let was_empty = self.recv_buffered == 0;
206 while datagram.data.len() + self.recv_buffered > window {
207 debug!("dropping stale datagram");
208 self.recv();
209 }
210
211 self.recv_buffered += datagram.data.len();
212 self.incoming.push_back(datagram);
213 Ok(was_empty)
214 }
215
216 fn make_space_for(&mut self, datagram_len: usize, send_buffer_size: usize) {
217 while !self.has_send_buffer_space(datagram_len, send_buffer_size) {
218 let Some(prev) = self.outgoing.pop_front() else {
219 break;
220 };
221 trace!(len = prev.data.len(), "dropping outgoing datagram");
222 self.outgoing_total -= prev.data.len();
223 }
224 }
225
226 fn has_send_buffer_space(&self, datagram_len: usize, send_buffer_size: usize) -> bool {
227 let Some(total) = self.outgoing_total.checked_add(datagram_len) else {
228 return false;
229 };
230
231 total <= send_buffer_size
232 }
233
234 pub(super) fn drop_oversized(&mut self, max_payload: usize) -> bool {
241 let mut dropped_any = false;
242 self.outgoing.retain(|datagram| {
243 let result = datagram.data.len() < max_payload;
244 if !result {
245 trace!(
246 "dropping {} byte datagram violating {} byte limit",
247 datagram.data.len(),
248 max_payload
249 );
250 self.outgoing_total -= datagram.data.len();
251 dropped_any = true;
252 }
253 result
254 });
255 dropped_any
256 }
257
258 pub(super) fn write<'a, 'b>(
263 &mut self,
264 buf: &mut PacketBuilder<'a, 'b>,
265 stat: &mut FrameStats,
266 ) -> bool {
267 let Some(datagram) = self.outgoing.pop_front() else {
268 return false;
269 };
270
271 if buf.frame_space_remaining() < datagram.size(true) {
272 self.outgoing.push_front(datagram);
275 return false;
276 }
277
278 self.outgoing_total -= datagram.data.len();
279 buf.write_frame(datagram, stat);
280 true
281 }
282
283 pub(super) fn recv(&mut self) -> Option<Bytes> {
284 let x = self.incoming.pop_front()?.data;
285 self.recv_buffered -= x.len();
286 Some(x)
287 }
288
289 pub(super) fn recv_many(&mut self, out: &mut [Bytes]) -> usize {
294 let n = out.len().min(self.incoming.len());
295 let mut received_bytes = 0;
296 for (i, d) in self.incoming.drain(..n).enumerate() {
297 received_bytes += d.data.len();
298 out[i] = d.data;
299 }
300 self.recv_buffered -= received_bytes;
301 n
302 }
303}
304
305#[cfg(test)]
306mod tests {
307 use super::*;
308
309 #[test]
310 fn make_space_for_accounts_for_new_datagram() {
311 let mut state = DatagramState::default();
312 state.outgoing.push_back(Datagram {
313 data: Bytes::from_static(&[0; 7]),
314 });
315 state.outgoing.push_back(Datagram {
316 data: Bytes::from_static(&[0; 2]),
317 });
318 state.outgoing_total = 9;
319
320 state.make_space_for(4, 10);
321
322 assert_eq!(state.outgoing.len(), 1);
323 assert_eq!(state.outgoing[0].data.len(), 2);
324 assert_eq!(state.outgoing_total, 2);
325 }
326
327 #[test]
328 fn make_space_for_handles_overflowing_capacity_check() {
329 let mut state = DatagramState::default();
330 state.outgoing.push_back(Datagram {
331 data: Bytes::from_static(&[0]),
332 });
333 state.outgoing_total = usize::MAX - 1;
334
335 state.make_space_for(2, usize::MAX);
336
337 assert!(state.outgoing.is_empty());
338 assert_eq!(state.outgoing_total, usize::MAX - 2);
339 }
340}
341
342#[derive(Debug, Error, Clone, Eq, PartialEq, Ord, PartialOrd, Hash)]
344pub enum SendDatagramError {
345 #[error("datagrams not supported by peer")]
347 UnsupportedByPeer,
348 #[error("datagram support disabled")]
350 Disabled,
351 #[error("datagram too large")]
356 TooLarge,
357 #[error("datagram send blocked")]
359 Blocked(Bytes),
360}