Skip to main content

sccp_protocol/server/
transport.rs

1//! Transport-neutral station admission.
2//!
3//! Listener owners establish the underlying connection and its security
4//! policy, then hand the ready byte stream to the protocol server. Session
5//! framing and lifecycle remain independent of the transport implementation.
6
7use 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
25/// Bidirectional asynchronous byte stream accepted by a station session.
26///
27/// The protocol server owns the stream after admission and applies identical
28/// framing, backpressure, registration, and shutdown behavior regardless of
29/// the underlying transport. A transport adapter may implement this trait with
30/// a plain socket, a decrypted secure stream, or an in-memory test stream.
31pub 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/// Identifies one admitted connection for the lifetime of a [`super::Server`].
38#[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/// Direction of one complete decrypted signaling frame.
55#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
56pub enum SignalingDirection {
57    StationToServer,
58    ServerToStation,
59}
60
61/// Describes whether the observed bytes were preserved or sanitized.
62#[derive(Clone, Copy, Debug, Eq, PartialEq)]
63pub enum SignalingFidelity {
64    Exact,
65    SecretsRedacted,
66    PayloadSuppressed,
67    IncompletePayloadSuppressed,
68}
69
70/// One bounded frame observed at the decrypted station transport boundary.
71///
72/// Complete unknown frames are retained exactly. Known credential-capable
73/// payloads, media key reservoirs, and every incomplete frame are sanitized
74/// before entering the observation queue. `Debug` never renders wire bytes.
75#[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/// One item from the server's bounded, nonblocking observation stream.
106///
107/// `observation_id` is unique and monotonic but does not define ordering across
108/// concurrent connections. `dropped_observations` is a batched loss counter
109/// carried by the next item admitted after queue saturation.
110#[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/// Connection lifecycle and signaling records emitted by the server.
119#[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)]
503/// Cloneable admission endpoint returned by [`super::Server::with_ingress`].
504///
505/// A listener owner performs transport-specific setup first, then submits the
506/// ready byte stream with its actual peer address, accepted local address, and
507/// transport classification. Clones share one bounded queue, so awaiting
508/// [`Self::accept`] propagates server backpressure instead of creating
509/// unbounded session work.
510pub 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    /// Transfer ownership of an accepted stream to the server run loop.
531    ///
532    /// `peer` identifies the remote station for events and address policy;
533    /// `local` is the concrete local endpoint used for server-list responses.
534    /// `transport` must describe the already-established stream because it is
535    /// checked against the device definition during registration. For secure
536    /// admission, complete the handshake and any certificate policy before
537    /// calling this method.
538    ///
539    /// The method waits for capacity in the ingress queue. It returns
540    /// [`ServerError::Stopped`] without starting a session if the run loop has
541    /// ended.
542    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    /// Admit a stream while retaining control of its underlying TCP markings.
557    ///
558    /// The server reapplies the selected station's signaling policy after the
559    /// registration message identifies it. Marking failures are logged while
560    /// registration and subsequent protocol traffic continue normally.
561    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}