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/// Why an admitted station connection stopped.
119///
120/// This classification deliberately stays independent of error text so
121/// telemetry consumers can aggregate connection outcomes without retaining
122/// transport- or protocol-specific details.
123#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
124#[non_exhaustive]
125pub enum StationDisconnectReason {
126    /// The station closed its side of the transport cleanly.
127    PeerClosure,
128    /// Reading from or writing to the station transport failed.
129    IoFailure,
130    /// The station sent no valid traffic before its keepalive deadline.
131    KeepaliveExpiry,
132    /// The server deliberately retired the session.
133    ServerRetirement,
134    /// The station explicitly requested that its session end.
135    StationRequest,
136    /// The server rejected the station during registration.
137    RegistrationRejected,
138    /// Malformed framing or another protocol error ended the session.
139    ProtocolFailure,
140    /// A non-I/O server failure ended the session.
141    ServerFailure,
142}
143
144/// Connection lifecycle and signaling records emitted by the server.
145#[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)]
530/// Cloneable admission endpoint returned by [`super::Server::with_ingress`].
531///
532/// A listener owner performs transport-specific setup first, then submits the
533/// ready byte stream with its actual peer address, accepted local address, and
534/// transport classification. Clones share one bounded queue, so awaiting
535/// [`Self::accept`] propagates server backpressure instead of creating
536/// unbounded session work.
537pub 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    /// Transfer ownership of an accepted stream to the server run loop.
558    ///
559    /// `peer` identifies the remote station for events and address policy;
560    /// `local` is the concrete local endpoint used for server-list responses.
561    /// `transport` must describe the already-established stream because it is
562    /// checked against the device definition during registration. For secure
563    /// admission, complete the handshake and any certificate policy before
564    /// calling this method.
565    ///
566    /// The method waits for capacity in the ingress queue. It returns
567    /// [`ServerError::Stopped`] without starting a session if the run loop has
568    /// ended.
569    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    /// Admit a stream while retaining control of its underlying TCP markings.
584    ///
585    /// The server reapplies the selected station's signaling policy after the
586    /// registration message identifies it. Marking failures are logged while
587    /// registration and subsequent protocol traffic continue normally.
588    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}