1use std::fmt;
8use std::net::SocketAddr;
9use std::num::NonZeroU64;
10use std::pin::Pin;
11use std::sync::Arc;
12use std::sync::atomic::{AtomicU64, Ordering};
13use std::task::{Context, Poll};
14use std::time::{SystemTime, UNIX_EPOCH};
15
16use tokio::io::ReadBuf;
17use tokio::io::{AsyncRead, AsyncWrite};
18use tokio::sync::mpsc;
19
20use super::ServerError;
21use super::qos::StationSocketQos;
22use crate::message::catalog::{ObservationSanitization, observation_sanitization};
23use crate::types::{SignalingQos, StationTransport};
24
25pub trait StationIo: AsyncRead + AsyncWrite + Unpin + Send {}
32
33impl<T> StationIo for T where T: AsyncRead + AsyncWrite + Unpin + Send {}
34
35pub(super) type BoxedStationIo = Box<dyn StationIo>;
36
37#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
39pub struct ObservationConnectionId(NonZeroU64);
40
41impl ObservationConnectionId {
42 pub const fn new(value: u64) -> Option<Self> {
43 match NonZeroU64::new(value) {
44 Some(value) => Some(Self(value)),
45 None => None,
46 }
47 }
48
49 pub const fn get(self) -> u64 {
50 self.0.get()
51 }
52}
53
54#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
56pub enum SignalingDirection {
57 StationToServer,
58 ServerToStation,
59}
60
61#[derive(Clone, Copy, Debug, Eq, PartialEq)]
63pub enum SignalingFidelity {
64 Exact,
65 SecretsRedacted,
66 PayloadSuppressed,
67 IncompletePayloadSuppressed,
68}
69
70#[derive(Clone)]
76pub struct SignalingObservation {
77 pub connection_id: ObservationConnectionId,
78 pub peer: SocketAddr,
79 pub local: SocketAddr,
80 pub transport: StationTransport,
81 pub direction: SignalingDirection,
82 pub protocol_header: Option<u32>,
83 pub message_id: Option<u32>,
84 pub fidelity: SignalingFidelity,
85 pub bytes: Vec<u8>,
86}
87
88impl fmt::Debug for SignalingObservation {
89 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
90 formatter
91 .debug_struct("SignalingObservation")
92 .field("connection_id", &self.connection_id)
93 .field("peer", &self.peer)
94 .field("local", &self.local)
95 .field("transport", &self.transport)
96 .field("direction", &self.direction)
97 .field("protocol_header", &self.protocol_header)
98 .field("message_id", &self.message_id)
99 .field("fidelity", &self.fidelity)
100 .field("byte_count", &self.bytes.len())
101 .finish()
102 }
103}
104
105#[derive(Clone, Debug)]
111pub struct ServerObservation {
112 pub observation_id: u64,
113 pub observed_at_unix_ms: u64,
114 pub dropped_observations: u64,
115 pub kind: ServerObservationKind,
116}
117
118#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
124#[non_exhaustive]
125pub enum StationDisconnectReason {
126 PeerClosure,
128 IoFailure,
130 KeepaliveExpiry,
132 ServerRetirement,
134 StationRequest,
136 RegistrationRejected,
138 ProtocolFailure,
140 ServerFailure,
142}
143
144#[derive(Clone, Debug)]
146#[non_exhaustive]
147pub enum ServerObservationKind {
148 Connected {
149 connection_id: ObservationConnectionId,
150 peer: SocketAddr,
151 local: SocketAddr,
152 transport: StationTransport,
153 },
154 Signaling(SignalingObservation),
155 Identified {
156 connection_id: ObservationConnectionId,
157 device_id: crate::types::DeviceId,
158 session_generation: crate::types::SessionGeneration,
159 },
160 Disconnected {
161 connection_id: ObservationConnectionId,
162 reason: StationDisconnectReason,
163 },
164}
165
166#[derive(Clone, Default)]
167pub(super) struct ObservationSink {
168 sender: Option<mpsc::Sender<ServerObservation>>,
169 next_observation_id: Arc<AtomicU64>,
170 dropped: Arc<AtomicU64>,
171}
172
173impl ObservationSink {
174 pub(super) fn new(sender: mpsc::Sender<ServerObservation>) -> Self {
175 Self {
176 sender: Some(sender),
177 next_observation_id: Arc::new(AtomicU64::new(1)),
178 dropped: Arc::new(AtomicU64::new(0)),
179 }
180 }
181
182 pub(super) fn observe(&self, kind: ServerObservationKind) {
183 let Some(sender) = &self.sender else {
184 return;
185 };
186 let Some(observation_id) = self
187 .next_observation_id
188 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
189 current.checked_add(1)
190 })
191 .ok()
192 else {
193 return;
194 };
195 let dropped_observations = self.dropped.swap(0, Ordering::Relaxed);
196 let observation = ServerObservation {
197 observation_id,
198 observed_at_unix_ms: unix_time_ms(),
199 dropped_observations,
200 kind,
201 };
202 if sender.try_send(observation).is_err() {
203 self.dropped
204 .fetch_add(dropped_observations.saturating_add(1), Ordering::Relaxed);
205 }
206 }
207
208 pub(super) fn is_active(&self) -> bool {
209 self.sender.is_some()
210 }
211}
212
213impl fmt::Debug for ObservationSink {
214 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
215 formatter
216 .debug_tuple("ObservationSink")
217 .field(&self.sender.as_ref().map(|_| "<registered>"))
218 .finish()
219 }
220}
221
222pub(super) struct ObservedStationIo {
223 inner: BoxedStationIo,
224 sink: ObservationSink,
225 connection_id: ObservationConnectionId,
226 peer: SocketAddr,
227 local: SocketAddr,
228 transport: StationTransport,
229 station_to_server: SignalingFrameBuffer,
230 server_to_station: SignalingFrameBuffer,
231}
232
233impl ObservedStationIo {
234 pub(super) fn new(
235 inner: BoxedStationIo,
236 sink: ObservationSink,
237 connection_id: ObservationConnectionId,
238 peer: SocketAddr,
239 local: SocketAddr,
240 transport: StationTransport,
241 ) -> Self {
242 Self {
243 inner,
244 sink,
245 connection_id,
246 peer,
247 local,
248 transport,
249 station_to_server: SignalingFrameBuffer::new(),
250 server_to_station: SignalingFrameBuffer::new(),
251 }
252 }
253
254 fn record(&mut self, direction: SignalingDirection, bytes: &[u8]) {
255 let frames = match direction {
256 SignalingDirection::StationToServer => self.station_to_server.push(bytes),
257 SignalingDirection::ServerToStation => self.server_to_station.push(bytes),
258 };
259 for frame in frames {
260 self.sink
261 .observe(ServerObservationKind::Signaling(SignalingObservation {
262 connection_id: self.connection_id,
263 peer: self.peer,
264 local: self.local,
265 transport: self.transport,
266 direction,
267 protocol_header: frame.protocol_header,
268 message_id: frame.message_id,
269 fidelity: frame.fidelity,
270 bytes: frame.bytes,
271 }));
272 }
273 }
274}
275
276impl AsyncRead for ObservedStationIo {
277 fn poll_read(
278 mut self: Pin<&mut Self>,
279 context: &mut Context<'_>,
280 buffer: &mut ReadBuf<'_>,
281 ) -> Poll<std::io::Result<()>> {
282 let previous = buffer.filled().len();
283 let result = Pin::new(&mut *self.inner).poll_read(context, buffer);
284 if matches!(result, Poll::Ready(Ok(()))) {
285 self.record(
286 SignalingDirection::StationToServer,
287 &buffer.filled()[previous..],
288 );
289 }
290 result
291 }
292}
293
294impl AsyncWrite for ObservedStationIo {
295 fn poll_write(
296 mut self: Pin<&mut Self>,
297 context: &mut Context<'_>,
298 bytes: &[u8],
299 ) -> Poll<Result<usize, std::io::Error>> {
300 let result = Pin::new(&mut *self.inner).poll_write(context, bytes);
301 if let Poll::Ready(Ok(written)) = result {
302 self.record(SignalingDirection::ServerToStation, &bytes[..written]);
303 }
304 result
305 }
306
307 fn poll_flush(
308 mut self: Pin<&mut Self>,
309 context: &mut Context<'_>,
310 ) -> Poll<Result<(), std::io::Error>> {
311 Pin::new(&mut *self.inner).poll_flush(context)
312 }
313
314 fn poll_shutdown(
315 mut self: Pin<&mut Self>,
316 context: &mut Context<'_>,
317 ) -> Poll<Result<(), std::io::Error>> {
318 Pin::new(&mut *self.inner).poll_shutdown(context)
319 }
320}
321
322impl Drop for ObservedStationIo {
323 fn drop(&mut self) {
324 for (direction, frame) in [
325 (
326 SignalingDirection::StationToServer,
327 self.station_to_server.take_incomplete(),
328 ),
329 (
330 SignalingDirection::ServerToStation,
331 self.server_to_station.take_incomplete(),
332 ),
333 ] {
334 if let Some(frame) = frame {
335 self.sink
336 .observe(ServerObservationKind::Signaling(SignalingObservation {
337 connection_id: self.connection_id,
338 peer: self.peer,
339 local: self.local,
340 transport: self.transport,
341 direction,
342 protocol_header: frame.protocol_header,
343 message_id: frame.message_id,
344 fidelity: frame.fidelity,
345 bytes: frame.bytes,
346 }));
347 }
348 }
349 }
350}
351
352struct ObservedFrame {
353 protocol_header: Option<u32>,
354 message_id: Option<u32>,
355 fidelity: SignalingFidelity,
356 bytes: Vec<u8>,
357}
358
359struct SignalingFrameBuffer {
360 state: SignalingFrameBufferState,
361}
362
363enum SignalingFrameBufferState {
364 Active(Vec<u8>),
365 Disabled,
366}
367
368impl SignalingFrameBuffer {
369 fn new() -> Self {
370 Self {
371 state: SignalingFrameBufferState::Active(Vec::new()),
372 }
373 }
374
375 fn push(&mut self, bytes: &[u8]) -> Vec<ObservedFrame> {
376 let SignalingFrameBufferState::Active(buffer) = &mut self.state else {
377 return Vec::new();
378 };
379 if !append_signaling_bytes(buffer, bytes) {
380 self.state = SignalingFrameBufferState::Disabled;
381 return Vec::new();
382 }
383 let mut frames = Vec::new();
384 let mut disable = false;
385 loop {
386 if buffer.len() < 12 {
387 break;
388 }
389 let wire_length = u32::from_le_bytes([buffer[0], buffer[1], buffer[2], buffer[3]]);
390 let Some(total_bytes) = usize::try_from(wire_length)
391 .ok()
392 .and_then(|length| length.checked_add(8))
393 else {
394 frames.push(sanitize_frame(std::mem::take(buffer), true));
395 disable = true;
396 break;
397 };
398 if wire_length < 4 || total_bytes > crate::message::wire::MAX_FRAME_SIZE {
399 frames.push(sanitize_frame(std::mem::take(buffer), true));
400 disable = true;
401 break;
402 }
403 if buffer.len() < total_bytes {
404 break;
405 }
406 let bytes = buffer[..total_bytes].to_vec();
407 let remaining = buffer[total_bytes..].to_vec();
408 buffer.fill(0);
409 *buffer = remaining;
410 frames.push(sanitize_frame(bytes, false));
411 }
412 if disable {
413 self.state = SignalingFrameBufferState::Disabled;
414 }
415 frames
416 }
417
418 fn take_incomplete(&mut self) -> Option<ObservedFrame> {
419 let state = std::mem::replace(&mut self.state, SignalingFrameBufferState::Disabled);
420 match state {
421 SignalingFrameBufferState::Active(bytes) if !bytes.is_empty() => {
422 Some(sanitize_frame(bytes, true))
423 }
424 SignalingFrameBufferState::Active(_) | SignalingFrameBufferState::Disabled => None,
425 }
426 }
427}
428
429fn append_signaling_bytes(buffer: &mut Vec<u8>, bytes: &[u8]) -> bool {
430 let Some(combined_bytes) = buffer.len().checked_add(bytes.len()) else {
431 buffer.fill(0);
432 return false;
433 };
434 let mut combined = Vec::new();
435 if combined.try_reserve_exact(combined_bytes).is_err() {
436 buffer.fill(0);
437 return false;
438 }
439 combined.extend_from_slice(buffer);
440 combined.extend_from_slice(bytes);
441 buffer.fill(0);
442 *buffer = combined;
443 true
444}
445
446impl Drop for SignalingFrameBuffer {
447 fn drop(&mut self) {
448 if let SignalingFrameBufferState::Active(bytes) = &mut self.state {
449 bytes.fill(0);
450 }
451 }
452}
453
454fn sanitize_frame(mut bytes: Vec<u8>, incomplete: bool) -> ObservedFrame {
455 let protocol_header = read_u32(&bytes, 4);
456 let message_id = read_u32(&bytes, 8);
457 let fidelity = if incomplete {
458 bytes.fill(0);
459 SignalingFidelity::IncompletePayloadSuppressed
460 } else {
461 match observation_sanitization(message_id, protocol_header) {
462 ObservationSanitization::Preserve => SignalingFidelity::Exact,
463 ObservationSanitization::Redact { start, end } => redact_range(&mut bytes, start..end),
464 ObservationSanitization::SuppressPayload => suppress_payload(&mut bytes),
465 }
466 };
467 ObservedFrame {
468 protocol_header,
469 message_id,
470 fidelity,
471 bytes,
472 }
473}
474
475fn suppress_payload(bytes: &mut [u8]) -> SignalingFidelity {
476 if let Some(payload) = bytes.get_mut(12..) {
477 payload.fill(0);
478 }
479 SignalingFidelity::PayloadSuppressed
480}
481
482fn redact_range(bytes: &mut [u8], range: std::ops::Range<usize>) -> SignalingFidelity {
483 if range.start >= bytes.len() {
484 bytes.fill(0);
485 return SignalingFidelity::PayloadSuppressed;
486 }
487 let end = range.end.min(bytes.len());
488 bytes[range.start..end].fill(0);
489 SignalingFidelity::SecretsRedacted
490}
491
492fn read_u32(bytes: &[u8], offset: usize) -> Option<u32> {
493 let value = bytes.get(offset..offset.checked_add(4)?)?;
494 Some(u32::from_le_bytes([value[0], value[1], value[2], value[3]]))
495}
496
497fn unix_time_ms() -> u64 {
498 SystemTime::now()
499 .duration_since(UNIX_EPOCH)
500 .map_or(0, |duration| {
501 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
502 })
503}
504
505pub(super) struct AcceptedStation {
506 pub stream: BoxedStationIo,
507 pub peer: SocketAddr,
508 pub local: SocketAddr,
509 pub transport: StationTransport,
510 pub socket_qos: Option<Box<dyn StationSocketQos>>,
511}
512
513impl fmt::Debug for AcceptedStation {
514 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
515 formatter
516 .debug_struct("AcceptedStation")
517 .field("stream", &"<station I/O>")
518 .field("peer", &self.peer)
519 .field("local", &self.local)
520 .field("transport", &self.transport)
521 .field(
522 "socket_qos",
523 &self.socket_qos.as_ref().map(|_| "<socket QoS control>"),
524 )
525 .finish()
526 }
527}
528
529#[derive(Clone, Debug)]
530pub struct ServerIngress {
538 sender: mpsc::Sender<AcceptedStation>,
539 signaling_qos: SignalingQos,
540}
541
542impl ServerIngress {
543 pub(super) fn channel(
544 capacity: usize,
545 signaling_qos: SignalingQos,
546 ) -> (Self, mpsc::Receiver<AcceptedStation>) {
547 let (sender, receiver) = mpsc::channel(capacity);
548 (
549 Self {
550 sender,
551 signaling_qos,
552 },
553 receiver,
554 )
555 }
556
557 pub async fn accept<S>(
570 &self,
571 stream: S,
572 peer: SocketAddr,
573 local: SocketAddr,
574 transport: StationTransport,
575 ) -> Result<(), ServerError>
576 where
577 S: StationIo + 'static,
578 {
579 self.admit(Box::new(stream), peer, local, transport, None)
580 .await
581 }
582
583 pub async fn accept_with_socket_qos<S, Q>(
589 &self,
590 stream: S,
591 peer: SocketAddr,
592 local: SocketAddr,
593 transport: StationTransport,
594 socket_qos: Q,
595 ) -> Result<(), ServerError>
596 where
597 S: StationIo + 'static,
598 Q: StationSocketQos + 'static,
599 {
600 super::report_socket_qos(None, peer, socket_qos.apply(self.signaling_qos));
601 self.admit(
602 Box::new(stream),
603 peer,
604 local,
605 transport,
606 Some(Box::new(socket_qos)),
607 )
608 .await
609 }
610
611 async fn admit(
612 &self,
613 stream: BoxedStationIo,
614 peer: SocketAddr,
615 local: SocketAddr,
616 transport: StationTransport,
617 socket_qos: Option<Box<dyn StationSocketQos>>,
618 ) -> Result<(), ServerError> {
619 self.sender
620 .send(AcceptedStation {
621 stream,
622 peer,
623 local,
624 transport,
625 socket_qos,
626 })
627 .await
628 .map_err(|_| ServerError::Stopped)
629 }
630}
631
632#[cfg(test)]
633mod observation_tests {
634 use super::*;
635 use crate::message::catalog::MessageId;
636 use crate::message::values::ProtocolVersion;
637
638 fn frame(protocol: u32, message_id: u32, payload_bytes: usize) -> Vec<u8> {
639 let mut bytes = Vec::with_capacity(payload_bytes + 12);
640 bytes.extend_from_slice(&u32::try_from(payload_bytes + 4).unwrap().to_le_bytes());
641 bytes.extend_from_slice(&protocol.to_le_bytes());
642 bytes.extend_from_slice(&message_id.to_le_bytes());
643 bytes.extend(std::iter::repeat_n(0x5a, payload_bytes));
644 bytes
645 }
646
647 #[test]
648 fn frame_buffer_handles_every_transport_boundary() {
649 let first = frame(
650 ProtocolVersion::V22.wire(),
651 MessageId::KeepAlive.wire_value(),
652 0,
653 );
654 let second = frame(
655 ProtocolVersion::V22.wire(),
656 MessageId::Register.wire_value(),
657 16,
658 );
659 let mut stream = first.clone();
660 stream.extend_from_slice(&second);
661 for boundary in 0..=stream.len() {
662 let mut buffer = SignalingFrameBuffer::new();
663 let mut observed = buffer.push(&stream[..boundary]);
664 observed.extend(buffer.push(&stream[boundary..]));
665 assert_eq!(observed.len(), 2, "boundary {boundary}");
666 assert_eq!(observed[0].bytes, first, "boundary {boundary}");
667 assert_eq!(observed[1].bytes, second, "boundary {boundary}");
668 }
669 }
670
671 #[test]
672 fn unknown_frames_remain_exact_and_debug_is_metadata_only() {
673 let bytes = frame(ProtocolVersion::V22.wire(), 0xfeed_beef, 16);
674 let observed = sanitize_frame(bytes.clone(), false);
675 assert_eq!(observed.fidelity, SignalingFidelity::Exact);
676 assert_eq!(observed.bytes, bytes);
677 let observation = SignalingObservation {
678 connection_id: ObservationConnectionId::new(1).unwrap(),
679 peer: "192.0.2.10:2000".parse().unwrap(),
680 local: "192.0.2.1:2000".parse().unwrap(),
681 transport: StationTransport::Secure,
682 direction: SignalingDirection::StationToServer,
683 protocol_header: observed.protocol_header,
684 message_id: observed.message_id,
685 fidelity: observed.fidelity,
686 bytes: observed.bytes,
687 };
688 let debug = format!("{observation:?}");
689 assert!(debug.contains("byte_count: 28"));
690 assert!(!debug.contains("90, 90"));
691 }
692
693 #[test]
694 fn every_media_secret_layout_redacts_only_its_fixed_reservoirs() {
695 for (protocol, message_id, range) in [
696 (
697 ProtocolVersion::V3.wire(),
698 MessageId::OpenReceiveChannel.wire_value(),
699 48..80,
700 ),
701 (
702 ProtocolVersion::V22.wire(),
703 MessageId::OpenReceiveChannel.wire_value(),
704 48..80,
705 ),
706 (
707 ProtocolVersion::V16.wire(),
708 MessageId::StartMediaTransmission.wire_value(),
709 64..96,
710 ),
711 (
712 ProtocolVersion::V17.wire(),
713 MessageId::StartMediaTransmission.wire_value(),
714 80..112,
715 ),
716 (
717 ProtocolVersion::V3.wire(),
718 MessageId::OpenMultimediaChannel.wire_value(),
719 128..160,
720 ),
721 (
722 ProtocolVersion::V22.wire(),
723 MessageId::OpenMultimediaChannel.wire_value(),
724 128..160,
725 ),
726 (
727 ProtocolVersion::V16.wire(),
728 MessageId::StartMultimediaTransmission.wire_value(),
729 132..164,
730 ),
731 (
732 ProtocolVersion::V17.wire(),
733 MessageId::StartMultimediaTransmission.wire_value(),
734 148..180,
735 ),
736 ] {
737 let observed = sanitize_frame(frame(protocol, message_id, 192), false);
738 assert_eq!(observed.fidelity, SignalingFidelity::SecretsRedacted);
739 assert!(observed.bytes[range.clone()].iter().all(|byte| *byte == 0));
740 assert_eq!(observed.bytes[range.start - 1], 0x5a);
741 assert_eq!(observed.bytes[range.end], 0x5a);
742 }
743 }
744
745 #[test]
746 fn fragmented_secret_frame_is_redacted_after_every_transport_boundary() {
747 let bytes = frame(
748 ProtocolVersion::V22.wire(),
749 MessageId::OpenReceiveChannel.wire_value(),
750 192,
751 );
752 for boundary in 0..=bytes.len() {
753 let mut buffer = SignalingFrameBuffer::new();
754 let mut observed = buffer.push(&bytes[..boundary]);
755 observed.extend(buffer.push(&bytes[boundary..]));
756 assert_eq!(observed.len(), 1, "boundary {boundary}");
757 assert_eq!(
758 observed[0].fidelity,
759 SignalingFidelity::SecretsRedacted,
760 "boundary {boundary}"
761 );
762 assert!(
763 observed[0].bytes[48..80].iter().all(|byte| *byte == 0),
764 "boundary {boundary}"
765 );
766 }
767 }
768
769 #[test]
770 fn incomplete_secret_frame_suppresses_all_available_bytes() {
771 let observed = sanitize_frame(
772 frame(
773 ProtocolVersion::V22.wire(),
774 MessageId::OpenReceiveChannel.wire_value(),
775 8,
776 ),
777 true,
778 );
779 assert_eq!(
780 observed.fidelity,
781 SignalingFidelity::IncompletePayloadSuppressed
782 );
783 assert!(observed.bytes.iter().all(|byte| *byte == 0));
784 assert_eq!(observed.protocol_header, Some(22));
785 assert_eq!(
786 observed.message_id,
787 Some(MessageId::OpenReceiveChannel.wire_value())
788 );
789 }
790
791 #[test]
792 fn malformed_framing_suppresses_the_buffer_and_disables_observation() {
793 let mut malformed = vec![0x5a; 32];
794 malformed[..4].copy_from_slice(&u32::MAX.to_le_bytes());
795 let mut buffer = SignalingFrameBuffer::new();
796 let observed = buffer.push(&malformed);
797 assert_eq!(observed.len(), 1);
798 assert_eq!(
799 observed[0].fidelity,
800 SignalingFidelity::IncompletePayloadSuppressed
801 );
802 assert!(observed[0].bytes.iter().all(|byte| *byte == 0));
803
804 let valid = frame(
805 ProtocolVersion::V22.wire(),
806 MessageId::KeepAlive.wire_value(),
807 0,
808 );
809 assert!(buffer.push(&valid).is_empty());
810 }
811
812 #[test]
813 fn station_service_submissions_never_publish_credential_capable_payloads() {
814 for message_id in [
815 MessageId::DeviceToUserData,
816 MessageId::DeviceToUserDataResponse,
817 MessageId::DeviceToUserDataV1,
818 MessageId::DeviceToUserDataResponseV1,
819 ] {
820 let message_id = message_id.wire_value();
821 let observed =
822 sanitize_frame(frame(ProtocolVersion::V22.wire(), message_id, 64), false);
823 assert_eq!(observed.fidelity, SignalingFidelity::PayloadSuppressed);
824 assert!(observed.bytes[12..].iter().all(|byte| *byte == 0));
825 assert_eq!(observed.message_id, Some(message_id));
826 }
827 }
828
829 #[test]
830 fn queue_overflow_is_reported_on_the_next_delivered_observation() {
831 let (sender, mut receiver) = mpsc::channel(1);
832 let sink = ObservationSink::new(sender);
833 let disconnected = || ServerObservationKind::Disconnected {
834 connection_id: ObservationConnectionId::new(1).unwrap(),
835 reason: StationDisconnectReason::PeerClosure,
836 };
837 sink.observe(disconnected());
838 sink.observe(disconnected());
839 assert_eq!(receiver.try_recv().unwrap().dropped_observations, 0);
840 sink.observe(disconnected());
841 assert_eq!(receiver.try_recv().unwrap().dropped_observations, 1);
842 }
843}