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