Skip to main content

logged_stream/
stream.rs

1use crate::ChannelLogger;
2use crate::MemoryStorageLogger;
3use crate::RecordFilter;
4use crate::buffer_formatter::BufferFormatter;
5use crate::logger::Logger;
6use crate::record::Record;
7use crate::record::RecordKind;
8use std::collections;
9use std::fmt;
10use std::io;
11use std::pin::Pin;
12use std::sync::mpsc;
13use std::task::Context;
14use std::task::Poll;
15use tokio::io as tokio_io;
16
17/// Wrapper for an IO object that logs every read, write, error, shutdown and drop that passes
18/// through it.
19///
20/// [`LoggedStream`] wraps an underlying IO object implementing the [`Read`] / [`Write`] traits, or
21/// their asynchronous [`tokio`] analogues [`AsyncRead`] / [`AsyncWrite`], and logs all read and
22/// write operations, errors, shutdowns and drops. It re-implements the same IO trait it wraps, so
23/// it is a drop-in replacement that works transparently in both synchronous and asynchronous code.
24///
25/// # Architecture
26///
27/// [`LoggedStream`] is generic over four independent, pluggable parts. Each logged event flows
28/// through them in order: `event -> Formatter -> Filter -> Logger`.
29///
30/// -   **The inner IO object (`S`).** The stream you are wrapping. Any type implementing [`Read`] /
31///     [`Write`] (or the [`tokio`] async equivalents) works — a socket, a file, an in-memory buffer,
32///     or your own type. [`LoggedStream`] implements the same IO trait `S` does, so it slots in
33///     wherever `S` was used.
34/// -   **Formatter ([`BufferFormatter`]).** Turns the read and written byte buffers into the display
35///     strings you see in the log.
36/// -   **Filter ([`RecordFilter`]).** Decides which records are logged. It runs on every record kind,
37///     including shutdown and drop.
38/// -   **Logger ([`Logger`]).** The sink that consumes accepted records.
39///
40/// All three of [`BufferFormatter`], [`RecordFilter`] and [`Logger`] are public, `Send + 'static`
41/// and object-safe, with blanket implementations for `Box<...>` (and `Arc<T>` where `T: Sync` for
42/// [`BufferFormatter`]). You are free to supply your own implementation of any part.
43///
44/// # Provided implementations
45///
46/// ## Formatters ([`BufferFormatter`])
47///
48/// Control how byte buffers are rendered. Each formatter stores a separator (default `:`) and
49/// exposes parallel constructors: `new`, `new_static`, `new_owned` and `new_default`.
50///
51/// | Formatter | Renders each byte as |
52/// | --- | --- |
53/// | [`LowercaseHexadecimalFormatter`] | lowercase hexadecimal — `0a:ff` |
54/// | [`UppercaseHexadecimalFormatter`] | uppercase hexadecimal — `0A:FF` |
55/// | [`DecimalFormatter`] | decimal — `10:255` |
56/// | [`OctalFormatter`] | octal — `012:377` |
57/// | [`BinaryFormatter`] | binary — `00001010:11111111` |
58///
59/// ## Filters ([`RecordFilter`])
60///
61/// Decide which records reach the logger.
62///
63/// | Filter | Behavior |
64/// | --- | --- |
65/// | [`DefaultFilter`] | Accepts every record. |
66/// | [`RecordKindFilter`] | Accepts only the record kinds in an allow-list given at construction. |
67/// | [`AllFilter`] | AND — a record passes only if every child filter accepts it (an empty list accepts everything). |
68/// | [`AnyFilter`] | OR — a record passes if any child filter accepts it (an empty list rejects everything). |
69///
70/// ## Loggers ([`Logger`])
71///
72/// Consume each accepted record.
73///
74/// | Logger | Destination |
75/// | --- | --- |
76/// | [`ConsoleLogger`] | Emits records through the `log` facade. |
77/// | [`FileLogger`] | Writes records to a file, one line per record. |
78/// | [`MemoryStorageLogger`] | Retains recent records in a bounded in-memory buffer. |
79/// | [`ChannelLogger`] | Sends records over an `mpsc` channel for handling elsewhere. |
80///
81/// [`ConsoleLogger`] and [`FileLogger`] additionally accept an optional prefix (`with_prefix` /
82/// `set_prefix`, none by default), written verbatim before the record kind character of every line.
83/// It tells apart several [`LoggedStream`]s — for example one per connection — that share a single
84/// console or file:
85///
86/// ```text
87/// [2026-07-20T12:34:56Z] [conn 5] > 01:02:03:04
88/// [2026-07-20T12:34:56Z] [conn 7] < 05:06:07:08
89/// ```
90///
91/// To let several [`FileLogger`]s write to one file, construct them with [`FileLogger::open`], which
92/// opens the file in append mode. Each record is rendered up front and written with a single
93/// `write_all` call, so concurrent loggers never interleave parts of a line. Passing independently
94/// opened non-append files (for example from `File::create`) instead gives each logger its own
95/// starting offset, and they will silently overwrite each other's records.
96///
97/// If none of the provided implementations matches your requirements, you can implement
98/// [`BufferFormatter`], [`RecordFilter`] or [`Logger`] yourself and pass your type to
99/// [`LoggedStream`] exactly like a built-in.
100///
101/// [`Read`]: io::Read
102/// [`Write`]: io::Write
103/// [`AsyncRead`]: tokio::io::AsyncRead
104/// [`AsyncWrite`]: tokio::io::AsyncWrite
105/// [`LowercaseHexadecimalFormatter`]: crate::LowercaseHexadecimalFormatter
106/// [`UppercaseHexadecimalFormatter`]: crate::UppercaseHexadecimalFormatter
107/// [`DecimalFormatter`]: crate::DecimalFormatter
108/// [`BinaryFormatter`]: crate::BinaryFormatter
109/// [`OctalFormatter`]: crate::OctalFormatter
110/// [`DefaultFilter`]: crate::DefaultFilter
111/// [`RecordKindFilter`]: crate::RecordKindFilter
112/// [`AllFilter`]: crate::AllFilter
113/// [`AnyFilter`]: crate::AnyFilter
114/// [`ConsoleLogger`]: crate::ConsoleLogger
115/// [`FileLogger`]: crate::FileLogger
116/// [`FileLogger::open`]: crate::FileLogger::open
117pub struct LoggedStream<
118    S: 'static,
119    Formatter: 'static,
120    Filter: RecordFilter + 'static,
121    L: Logger + 'static,
122> {
123    inner_stream: S,
124    formatter: Formatter,
125    filter: Filter,
126    logger: L,
127}
128
129impl<S: 'static, Formatter: 'static, Filter: RecordFilter + 'static, L: Logger + 'static>
130    LoggedStream<S, Formatter, Filter, L>
131{
132    /// Construct a new instance of [`LoggedStream`] using provided arguments.
133    pub fn new(stream: S, formatter: Formatter, filter: Filter, logger: L) -> Self {
134        Self {
135            inner_stream: stream,
136            formatter,
137            filter,
138            logger,
139        }
140    }
141
142    /// Emit a custom [`RecordKind::Open`] record carrying `message`.
143    ///
144    /// [`RecordKind::Open`] is never produced automatically by the read, write, shutdown and drop
145    /// machinery — this method is the way to emit one. Use it to annotate the start of a stream,
146    /// for example to record the peer of a freshly established connection
147    /// (`"Established connection with 127.0.0.1:8080"`) or other per-stream metadata.
148    ///
149    /// Like every other record, the `Open` record is passed through the filter before it reaches
150    /// the logger, so a `RecordKindFilter` that does not allow `Open` will suppress it. The message
151    /// is logged verbatim; it is not run through the formatter, which only applies to byte buffers.
152    ///
153    /// For asynchronous streams, call this before splitting the wrapper with `tokio::io::split`,
154    /// since the resulting halves do not expose it.
155    pub fn log_open(&mut self, message: impl Into<String>) {
156        self.emit(Record::new(RecordKind::Open, message.into()));
157    }
158
159    /// Route a record through the filter, logging it only if the filter accepts it.
160    ///
161    /// Every record — reads, writes, errors, shutdowns, drops and manual `Open` markers — is
162    /// emitted through this single method, so the filter is applied consistently to all of them.
163    #[inline]
164    fn emit(&mut self, record: Record) {
165        if self.filter.check(&record) {
166            self.logger.log(record);
167        }
168    }
169}
170
171impl<S: 'static, Formatter: 'static, Filter: RecordFilter + 'static>
172    LoggedStream<S, Formatter, Filter, MemoryStorageLogger>
173{
174    #[inline]
175    pub fn get_log_records(&self) -> collections::VecDeque<Record> {
176        self.logger.get_log_records()
177    }
178
179    #[inline]
180    pub fn clear_log_records(&mut self) {
181        self.logger.clear_log_records()
182    }
183}
184
185impl<S: 'static, Formatter: 'static, Filter: RecordFilter + 'static>
186    LoggedStream<S, Formatter, Filter, ChannelLogger>
187{
188    #[inline]
189    pub fn take_receiver(&mut self) -> Option<mpsc::Receiver<Record>> {
190        self.logger.take_receiver()
191    }
192
193    #[inline]
194    pub fn take_receiver_unchecked(&mut self) -> mpsc::Receiver<Record> {
195        self.logger.take_receiver_unchecked()
196    }
197}
198
199impl<
200    S: fmt::Debug + 'static,
201    Formatter: fmt::Debug + 'static,
202    Filter: RecordFilter + fmt::Debug + 'static,
203    L: Logger + fmt::Debug + 'static,
204> fmt::Debug for LoggedStream<S, Formatter, Filter, L>
205{
206    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
207        f.debug_struct("LoggedStream")
208            .field("inner_stream", &self.inner_stream)
209            .field("formatter", &self.formatter)
210            .field("filter", &self.filter)
211            .field("logger", &self.logger)
212            .finish()
213    }
214}
215
216impl<
217    S: io::Read + 'static,
218    Formatter: BufferFormatter + 'static,
219    Filter: RecordFilter + 'static,
220    L: Logger + 'static,
221> io::Read for LoggedStream<S, Formatter, Filter, L>
222{
223    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
224        let result = self.inner_stream.read(buf);
225
226        match &result {
227            Ok(length) => {
228                let record = Record::new(
229                    RecordKind::Read,
230                    self.formatter.format_buffer(&buf[0..*length]),
231                );
232                self.emit(record);
233            }
234            Err(e) if matches!(e.kind(), io::ErrorKind::WouldBlock) => {}
235            Err(e) => {
236                self.emit(Record::new(
237                    RecordKind::Error,
238                    format!("Error during read: {e}"),
239                ));
240            }
241        };
242
243        result
244    }
245}
246
247impl<
248    S: tokio_io::AsyncRead + Unpin + 'static,
249    Formatter: BufferFormatter + Unpin + 'static,
250    Filter: RecordFilter + Unpin + 'static,
251    L: Logger + Unpin + 'static,
252> tokio_io::AsyncRead for LoggedStream<S, Formatter, Filter, L>
253{
254    fn poll_read(
255        self: Pin<&mut Self>,
256        cx: &mut Context<'_>,
257        buf: &mut tokio_io::ReadBuf<'_>,
258    ) -> Poll<io::Result<()>> {
259        let mut_self = self.get_mut();
260        let length_before_read = buf.filled().len();
261        let result = Pin::new(&mut mut_self.inner_stream).poll_read(cx, buf);
262        let length_after_read = buf.filled().len();
263        let diff = length_after_read - length_before_read;
264
265        match &result {
266            Poll::Ready(Ok(())) if diff == 0 => {}
267            Poll::Ready(Ok(())) => {
268                let record = Record::new(
269                    RecordKind::Read,
270                    mut_self
271                        .formatter
272                        .format_buffer(&(buf.filled())[length_before_read..length_after_read]),
273                );
274                mut_self.emit(record);
275            }
276            Poll::Ready(Err(e)) => {
277                mut_self.emit(Record::new(
278                    RecordKind::Error,
279                    format!("Error during async read: {e}"),
280                ));
281            }
282            Poll::Pending => {}
283        }
284
285        result
286    }
287}
288
289impl<
290    S: io::Write + 'static,
291    Formatter: BufferFormatter + 'static,
292    Filter: RecordFilter + 'static,
293    L: Logger + 'static,
294> io::Write for LoggedStream<S, Formatter, Filter, L>
295{
296    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
297        let result = self.inner_stream.write(buf);
298
299        match &result {
300            Ok(length) => {
301                let record = Record::new(
302                    RecordKind::Write,
303                    self.formatter.format_buffer(&buf[0..*length]),
304                );
305                self.emit(record);
306            }
307            Err(e)
308                if matches!(
309                    e.kind(),
310                    io::ErrorKind::WriteZero | io::ErrorKind::WouldBlock
311                ) => {}
312            Err(e) => {
313                self.emit(Record::new(
314                    RecordKind::Error,
315                    format!("Error during write: {e}"),
316                ));
317            }
318        };
319
320        result
321    }
322
323    fn flush(&mut self) -> io::Result<()> {
324        self.inner_stream.flush()
325    }
326}
327
328impl<
329    S: tokio_io::AsyncWrite + Unpin + 'static,
330    Formatter: BufferFormatter + Unpin + 'static,
331    Filter: RecordFilter + Unpin + 'static,
332    L: Logger + Unpin + 'static,
333> tokio_io::AsyncWrite for LoggedStream<S, Formatter, Filter, L>
334{
335    fn poll_write(
336        self: Pin<&mut Self>,
337        cx: &mut Context<'_>,
338        buf: &[u8],
339    ) -> Poll<Result<usize, io::Error>> {
340        let mut_self = self.get_mut();
341        let result = Pin::new(&mut mut_self.inner_stream).poll_write(cx, buf);
342
343        match &result {
344            Poll::Ready(Ok(length)) => {
345                let record = Record::new(
346                    RecordKind::Write,
347                    mut_self.formatter.format_buffer(&buf[0..*length]),
348                );
349                mut_self.emit(record);
350            }
351            Poll::Ready(Err(e)) => {
352                mut_self.emit(Record::new(
353                    RecordKind::Error,
354                    format!("Error during async write: {e}"),
355                ));
356            }
357            Poll::Pending => {}
358        }
359
360        result
361    }
362
363    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> {
364        Pin::new(&mut self.get_mut().inner_stream).poll_flush(cx)
365    }
366
367    fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> {
368        let mut_self = self.get_mut();
369        let result = Pin::new(&mut mut_self.inner_stream).poll_shutdown(cx);
370
371        mut_self.emit(Record::new(
372            RecordKind::Shutdown,
373            String::from("Writer shutdown request."),
374        ));
375
376        result
377    }
378}
379
380impl<S: 'static, Formatter: 'static, Filter: RecordFilter + 'static, L: Logger + 'static> Drop
381    for LoggedStream<S, Formatter, Filter, L>
382{
383    fn drop(&mut self) {
384        self.emit(Record::new(RecordKind::Drop, String::from("Deallocated.")));
385    }
386}
387
388#[cfg(test)]
389mod tests {
390    use crate::ChannelLogger;
391    use crate::DecimalFormatter;
392    use crate::DefaultFilter;
393    use crate::LoggedStream;
394    use crate::LowercaseHexadecimalFormatter;
395    use crate::MemoryStorageLogger;
396    use crate::RecordKind;
397    use crate::RecordKindFilter;
398    use crate::UppercaseHexadecimalFormatter;
399    use std::cell::RefCell;
400    use std::io::Cursor;
401    use std::io::ErrorKind;
402    use std::io::Read;
403    use std::io::Write;
404    use std::pin::Pin;
405    use std::rc::Rc;
406    use std::task::Context;
407    use std::task::Poll;
408
409    //////////////////////////////////////////////////////////////////////////////////////////////////////////
410    // Test doubles
411    //////////////////////////////////////////////////////////////////////////////////////////////////////////
412
413    /// A synchronous reader whose every `read` fails with the given [`ErrorKind`].
414    struct ErrReader(ErrorKind);
415
416    impl Read for ErrReader {
417        fn read(&mut self, _buf: &mut [u8]) -> std::io::Result<usize> {
418            Err(std::io::Error::new(self.0, "boom"))
419        }
420    }
421
422    /// A synchronous writer whose every `write` fails with the given [`ErrorKind`].
423    struct ErrWriter(ErrorKind);
424
425    impl Write for ErrWriter {
426        fn write(&mut self, _buf: &[u8]) -> std::io::Result<usize> {
427            Err(std::io::Error::new(self.0, "boom"))
428        }
429
430        fn flush(&mut self) -> std::io::Result<()> {
431            Ok(())
432        }
433    }
434
435    /// An asynchronous reader whose every `poll_read` fails with the given [`ErrorKind`].
436    struct ErrAsyncReader(ErrorKind);
437
438    impl tokio::io::AsyncRead for ErrAsyncReader {
439        fn poll_read(
440            self: Pin<&mut Self>,
441            _cx: &mut Context<'_>,
442            _buf: &mut tokio::io::ReadBuf<'_>,
443        ) -> Poll<std::io::Result<()>> {
444            Poll::Ready(Err(std::io::Error::new(self.0, "boom")))
445        }
446    }
447
448    /// An asynchronous writer whose every `poll_write` fails with the given [`ErrorKind`].
449    struct ErrAsyncWriter(ErrorKind);
450
451    impl tokio::io::AsyncWrite for ErrAsyncWriter {
452        fn poll_write(
453            self: Pin<&mut Self>,
454            _cx: &mut Context<'_>,
455            _buf: &[u8],
456        ) -> Poll<std::io::Result<usize>> {
457            Poll::Ready(Err(std::io::Error::new(self.0, "boom")))
458        }
459
460        fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
461            Poll::Ready(Ok(()))
462        }
463
464        fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<std::io::Result<()>> {
465            Poll::Ready(Ok(()))
466        }
467    }
468
469    /// A synchronous writer that appends everything written to it into a shared buffer, so a test
470    /// can assert that bytes actually reach the wrapped stream.
471    struct RecordingWriter(Rc<RefCell<Vec<u8>>>);
472
473    impl Write for RecordingWriter {
474        fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
475            self.0.borrow_mut().extend_from_slice(buf);
476            Ok(buf.len())
477        }
478
479        fn flush(&mut self) -> std::io::Result<()> {
480            Ok(())
481        }
482    }
483
484    //////////////////////////////////////////////////////////////////////////////////////////////////////////
485    // Construction & transparency
486    //////////////////////////////////////////////////////////////////////////////////////////////////////////
487
488    #[test]
489    fn test_read_returns_inner_bytes_unchanged() {
490        let data = vec![0x01, 0x02, 0x03, 0x04];
491        let mut stream = LoggedStream::new(
492            Cursor::new(data.clone()),
493            LowercaseHexadecimalFormatter::new_default(),
494            DefaultFilter,
495            MemoryStorageLogger::new(16),
496        );
497
498        let mut buf = [0u8; 4];
499        stream.read_exact(&mut buf).unwrap();
500
501        assert_eq!(buf.to_vec(), data);
502    }
503
504    #[test]
505    fn test_write_reaches_inner_stream() {
506        let sink = Rc::new(RefCell::new(Vec::new()));
507        let mut stream = LoggedStream::new(
508            RecordingWriter(Rc::clone(&sink)),
509            LowercaseHexadecimalFormatter::new_default(),
510            DefaultFilter,
511            MemoryStorageLogger::new(16),
512        );
513
514        stream.write_all(&[0xde, 0xad, 0xbe, 0xef]).unwrap();
515
516        assert_eq!(*sink.borrow(), vec![0xde, 0xad, 0xbe, 0xef]);
517    }
518
519    //////////////////////////////////////////////////////////////////////////////////////////////////////////
520    // Read logging
521    //////////////////////////////////////////////////////////////////////////////////////////////////////////
522
523    #[test]
524    fn test_read_logs_read_record_with_formatted_content() {
525        let mut stream = LoggedStream::new(
526            Cursor::new(vec![0x0a, 0xff]),
527            LowercaseHexadecimalFormatter::new_default(),
528            DefaultFilter,
529            MemoryStorageLogger::new(16),
530        );
531
532        let mut buf = [0u8; 2];
533        stream.read_exact(&mut buf).unwrap();
534
535        let records = stream.get_log_records();
536        assert_eq!(records.len(), 1);
537        assert_eq!(records[0].kind, RecordKind::Read);
538        assert_eq!(records[0].message, "0a:ff");
539    }
540
541    #[test]
542    fn test_read_uses_configured_formatter() {
543        let mut stream = LoggedStream::new(
544            Cursor::new(vec![10, 255]),
545            DecimalFormatter::new_default(),
546            DefaultFilter,
547            MemoryStorageLogger::new(16),
548        );
549
550        let mut buf = [0u8; 2];
551        stream.read_exact(&mut buf).unwrap();
552
553        assert_eq!(stream.get_log_records()[0].message, "10:255");
554    }
555
556    #[test]
557    fn test_multiple_reads_log_in_order() {
558        let mut stream = LoggedStream::new(
559            Cursor::new(vec![0x01, 0x02, 0x03, 0x04]),
560            LowercaseHexadecimalFormatter::new_default(),
561            DefaultFilter,
562            MemoryStorageLogger::new(16),
563        );
564
565        let mut buf = [0u8; 2];
566        stream.read_exact(&mut buf).unwrap();
567        stream.read_exact(&mut buf).unwrap();
568
569        let records = stream.get_log_records();
570        assert_eq!(records.len(), 2);
571        assert_eq!(records[0].message, "01:02");
572        assert_eq!(records[1].message, "03:04");
573    }
574
575    //////////////////////////////////////////////////////////////////////////////////////////////////////////
576    // Write logging
577    //////////////////////////////////////////////////////////////////////////////////////////////////////////
578
579    #[test]
580    fn test_write_logs_write_record_with_formatted_content() {
581        let mut stream = LoggedStream::new(
582            Cursor::new(Vec::<u8>::new()),
583            UppercaseHexadecimalFormatter::new_default(),
584            DefaultFilter,
585            MemoryStorageLogger::new(16),
586        );
587
588        stream.write_all(&[0x0a, 0xff]).unwrap();
589
590        let records = stream.get_log_records();
591        assert_eq!(records.len(), 1);
592        assert_eq!(records[0].kind, RecordKind::Write);
593        assert_eq!(records[0].message, "0A:FF");
594    }
595
596    #[test]
597    fn test_flush_does_not_log() {
598        let mut stream = LoggedStream::new(
599            Cursor::new(Vec::<u8>::new()),
600            LowercaseHexadecimalFormatter::new_default(),
601            DefaultFilter,
602            MemoryStorageLogger::new(16),
603        );
604
605        stream.flush().unwrap();
606
607        assert!(stream.get_log_records().is_empty());
608    }
609
610    //////////////////////////////////////////////////////////////////////////////////////////////////////////
611    // Record filtering
612    //////////////////////////////////////////////////////////////////////////////////////////////////////////
613
614    #[test]
615    fn test_read_record_suppressed_by_filter() {
616        // The filter allows only Write, so the Read record is dropped.
617        let mut stream = LoggedStream::new(
618            Cursor::new(vec![0x01, 0x02]),
619            LowercaseHexadecimalFormatter::new_default(),
620            RecordKindFilter::new(&[RecordKind::Write]),
621            MemoryStorageLogger::new(16),
622        );
623
624        let mut buf = [0u8; 2];
625        stream.read_exact(&mut buf).unwrap();
626
627        assert!(stream.get_log_records().is_empty());
628    }
629
630    #[test]
631    fn test_write_record_suppressed_by_filter() {
632        // The filter allows only Read, so the Write record is dropped.
633        let mut stream = LoggedStream::new(
634            Cursor::new(Vec::<u8>::new()),
635            LowercaseHexadecimalFormatter::new_default(),
636            RecordKindFilter::new(&[RecordKind::Read]),
637            MemoryStorageLogger::new(16),
638        );
639
640        stream.write_all(&[0x01, 0x02]).unwrap();
641
642        assert!(stream.get_log_records().is_empty());
643    }
644
645    //////////////////////////////////////////////////////////////////////////////////////////////////////////
646    // Error handling (message content, filtering and swallowed error kinds)
647    //////////////////////////////////////////////////////////////////////////////////////////////////////////
648
649    #[test]
650    fn test_read_error_logs_error_record() {
651        let mut stream = LoggedStream::new(
652            ErrReader(ErrorKind::Other),
653            LowercaseHexadecimalFormatter::new_default(),
654            DefaultFilter,
655            MemoryStorageLogger::new(16),
656        );
657
658        let mut buf = [0u8; 4];
659        let _ = stream.read(&mut buf);
660
661        let records = stream.get_log_records();
662        assert_eq!(records.len(), 1);
663        assert_eq!(records[0].kind, RecordKind::Error);
664        assert!(records[0].message.starts_with("Error during read:"));
665    }
666
667    #[test]
668    fn test_read_error_suppressed_by_filter_without_error() {
669        let mut stream = LoggedStream::new(
670            ErrReader(ErrorKind::Other),
671            LowercaseHexadecimalFormatter::new_default(),
672            RecordKindFilter::new(&[RecordKind::Read, RecordKind::Write]),
673            MemoryStorageLogger::new(16),
674        );
675
676        let mut buf = [0u8; 4];
677        let _ = stream.read(&mut buf);
678
679        assert!(stream.get_log_records().is_empty());
680    }
681
682    #[test]
683    fn test_read_would_block_is_not_logged() {
684        // WouldBlock is a transient non-event and must not produce a record.
685        let mut stream = LoggedStream::new(
686            ErrReader(ErrorKind::WouldBlock),
687            LowercaseHexadecimalFormatter::new_default(),
688            DefaultFilter,
689            MemoryStorageLogger::new(16),
690        );
691
692        let mut buf = [0u8; 4];
693        let _ = stream.read(&mut buf);
694
695        assert!(stream.get_log_records().is_empty());
696    }
697
698    #[test]
699    fn test_write_error_logs_error_record() {
700        let mut stream = LoggedStream::new(
701            ErrWriter(ErrorKind::Other),
702            LowercaseHexadecimalFormatter::new_default(),
703            DefaultFilter,
704            MemoryStorageLogger::new(16),
705        );
706
707        let _ = stream.write(&[0x01, 0x02]);
708
709        let records = stream.get_log_records();
710        assert_eq!(records.len(), 1);
711        assert_eq!(records[0].kind, RecordKind::Error);
712        assert!(records[0].message.starts_with("Error during write:"));
713    }
714
715    #[test]
716    fn test_write_error_suppressed_by_filter_without_error() {
717        let mut stream = LoggedStream::new(
718            ErrWriter(ErrorKind::Other),
719            LowercaseHexadecimalFormatter::new_default(),
720            RecordKindFilter::new(&[RecordKind::Read, RecordKind::Write]),
721            MemoryStorageLogger::new(16),
722        );
723
724        let _ = stream.write(&[0x01, 0x02]);
725
726        assert!(stream.get_log_records().is_empty());
727    }
728
729    #[test]
730    fn test_write_would_block_is_not_logged() {
731        let mut stream = LoggedStream::new(
732            ErrWriter(ErrorKind::WouldBlock),
733            LowercaseHexadecimalFormatter::new_default(),
734            DefaultFilter,
735            MemoryStorageLogger::new(16),
736        );
737
738        let _ = stream.write(&[0x01, 0x02]);
739
740        assert!(stream.get_log_records().is_empty());
741    }
742
743    #[test]
744    fn test_write_write_zero_is_not_logged() {
745        let mut stream = LoggedStream::new(
746            ErrWriter(ErrorKind::WriteZero),
747            LowercaseHexadecimalFormatter::new_default(),
748            DefaultFilter,
749            MemoryStorageLogger::new(16),
750        );
751
752        let _ = stream.write(&[0x01, 0x02]);
753
754        assert!(stream.get_log_records().is_empty());
755    }
756
757    //////////////////////////////////////////////////////////////////////////////////////////////////////////
758    // Manual Open marker (log_open)
759    //////////////////////////////////////////////////////////////////////////////////////////////////////////
760
761    #[test]
762    fn test_log_open_emits_open_record() {
763        let mut stream = LoggedStream::new(
764            Cursor::new(Vec::<u8>::new()),
765            LowercaseHexadecimalFormatter::new_default(),
766            DefaultFilter,
767            MemoryStorageLogger::new(16),
768        );
769        stream.log_open("Established connection with 127.0.0.1:8080");
770
771        let records = stream.get_log_records();
772        assert_eq!(records.len(), 1);
773        assert_eq!(records[0].kind, RecordKind::Open);
774        assert_eq!(
775            records[0].message,
776            "Established connection with 127.0.0.1:8080"
777        );
778    }
779
780    #[test]
781    fn test_log_open_passes_filter_allowing_open() {
782        let mut stream = LoggedStream::new(
783            Cursor::new(Vec::<u8>::new()),
784            LowercaseHexadecimalFormatter::new_default(),
785            RecordKindFilter::new(&[RecordKind::Open]),
786            MemoryStorageLogger::new(16),
787        );
788        stream.log_open("kept");
789
790        let records = stream.get_log_records();
791        assert_eq!(records.len(), 1);
792        assert_eq!(records[0].kind, RecordKind::Open);
793        assert_eq!(records[0].message, "kept");
794    }
795
796    #[test]
797    fn test_log_open_suppressed_by_filter_without_open() {
798        let mut stream = LoggedStream::new(
799            Cursor::new(Vec::<u8>::new()),
800            LowercaseHexadecimalFormatter::new_default(),
801            RecordKindFilter::new(&[RecordKind::Read, RecordKind::Write]),
802            MemoryStorageLogger::new(16),
803        );
804        stream.log_open("should be filtered out");
805        assert!(stream.get_log_records().is_empty());
806    }
807
808    //////////////////////////////////////////////////////////////////////////////////////////////////////////
809    // Lifecycle: drop & shutdown
810    //////////////////////////////////////////////////////////////////////////////////////////////////////////
811
812    #[test]
813    fn test_drop_logs_drop_record() {
814        // The ChannelLogger's receiver outlives the stream, so we can observe the Drop record.
815        let mut stream = LoggedStream::new(
816            Cursor::new(Vec::<u8>::new()),
817            LowercaseHexadecimalFormatter::new_default(),
818            DefaultFilter,
819            ChannelLogger::new(),
820        );
821        let receiver = stream.take_receiver_unchecked();
822
823        drop(stream);
824
825        let record = receiver
826            .recv()
827            .expect("dropping the stream should emit a record");
828        assert_eq!(record.kind, RecordKind::Drop);
829        assert_eq!(record.message, "Deallocated.");
830    }
831
832    #[tokio::test]
833    async fn test_async_shutdown_logs_shutdown_record() {
834        use tokio::io::AsyncWriteExt;
835
836        let (client, _server) = tokio::io::duplex(64);
837        let mut stream = LoggedStream::new(
838            client,
839            LowercaseHexadecimalFormatter::new_default(),
840            DefaultFilter,
841            MemoryStorageLogger::new(16),
842        );
843
844        stream.shutdown().await.unwrap();
845
846        let records = stream.get_log_records();
847        assert_eq!(records.len(), 1);
848        assert_eq!(records[0].kind, RecordKind::Shutdown);
849        assert_eq!(records[0].message, "Writer shutdown request.");
850    }
851
852    //////////////////////////////////////////////////////////////////////////////////////////////////////////
853    // Asynchronous IO
854    //////////////////////////////////////////////////////////////////////////////////////////////////////////
855
856    #[tokio::test]
857    async fn test_async_read_logs_read_record() {
858        use tokio::io::AsyncReadExt;
859        use tokio::io::AsyncWriteExt;
860
861        let (client, mut server) = tokio::io::duplex(64);
862        server.write_all(&[0xaa, 0xbb]).await.unwrap();
863
864        let mut stream = LoggedStream::new(
865            client,
866            LowercaseHexadecimalFormatter::new_default(),
867            DefaultFilter,
868            MemoryStorageLogger::new(16),
869        );
870
871        let mut buf = [0u8; 2];
872        stream.read_exact(&mut buf).await.unwrap();
873
874        assert_eq!(buf, [0xaa, 0xbb]);
875        let records = stream.get_log_records();
876        assert_eq!(records.len(), 1);
877        assert_eq!(records[0].kind, RecordKind::Read);
878        assert_eq!(records[0].message, "aa:bb");
879    }
880
881    #[tokio::test]
882    async fn test_async_write_logs_write_record() {
883        use tokio::io::AsyncWriteExt;
884
885        let (client, _server) = tokio::io::duplex(64);
886        let mut stream = LoggedStream::new(
887            client,
888            LowercaseHexadecimalFormatter::new_default(),
889            DefaultFilter,
890            MemoryStorageLogger::new(16),
891        );
892
893        stream.write_all(&[0x01, 0x02, 0x03, 0x04]).await.unwrap();
894
895        let records = stream.get_log_records();
896        assert_eq!(records.len(), 1);
897        assert_eq!(records[0].kind, RecordKind::Write);
898        assert_eq!(records[0].message, "01:02:03:04");
899    }
900
901    #[tokio::test]
902    async fn test_async_read_error_suppressed_by_filter_without_error() {
903        use tokio::io::AsyncReadExt;
904
905        let mut stream = LoggedStream::new(
906            ErrAsyncReader(ErrorKind::Other),
907            LowercaseHexadecimalFormatter::new_default(),
908            RecordKindFilter::new(&[RecordKind::Read, RecordKind::Write]),
909            MemoryStorageLogger::new(16),
910        );
911
912        let mut buf = [0u8; 4];
913        let _ = stream.read(&mut buf).await;
914
915        assert!(stream.get_log_records().is_empty());
916    }
917
918    #[tokio::test]
919    async fn test_async_write_error_suppressed_by_filter_without_error() {
920        use tokio::io::AsyncWriteExt;
921
922        let mut stream = LoggedStream::new(
923            ErrAsyncWriter(ErrorKind::Other),
924            LowercaseHexadecimalFormatter::new_default(),
925            RecordKindFilter::new(&[RecordKind::Read, RecordKind::Write]),
926            MemoryStorageLogger::new(16),
927        );
928
929        let _ = stream.write(&[0x01, 0x02]).await;
930
931        assert!(stream.get_log_records().is_empty());
932    }
933
934    #[tokio::test]
935    async fn test_async_read_error_logs_error_record() {
936        use tokio::io::AsyncReadExt;
937
938        let mut stream = LoggedStream::new(
939            ErrAsyncReader(ErrorKind::Other),
940            LowercaseHexadecimalFormatter::new_default(),
941            DefaultFilter,
942            MemoryStorageLogger::new(16),
943        );
944
945        let mut buf = [0u8; 4];
946        let _ = stream.read(&mut buf).await;
947
948        let records = stream.get_log_records();
949        assert_eq!(records.len(), 1);
950        assert_eq!(records[0].kind, RecordKind::Error);
951        assert!(records[0].message.starts_with("Error during async read:"));
952    }
953
954    #[tokio::test]
955    async fn test_async_write_error_logs_error_record() {
956        use tokio::io::AsyncWriteExt;
957
958        let mut stream = LoggedStream::new(
959            ErrAsyncWriter(ErrorKind::Other),
960            LowercaseHexadecimalFormatter::new_default(),
961            DefaultFilter,
962            MemoryStorageLogger::new(16),
963        );
964
965        let _ = stream.write(&[0x01, 0x02]).await;
966
967        let records = stream.get_log_records();
968        assert_eq!(records.len(), 1);
969        assert_eq!(records[0].kind, RecordKind::Error);
970        assert!(records[0].message.starts_with("Error during async write:"));
971    }
972
973    //////////////////////////////////////////////////////////////////////////////////////////////////////////
974    // Logger accessors
975    //////////////////////////////////////////////////////////////////////////////////////////////////////////
976
977    #[test]
978    fn test_memory_storage_get_and_clear() {
979        let mut stream = LoggedStream::new(
980            Cursor::new(vec![0x01, 0x02]),
981            LowercaseHexadecimalFormatter::new_default(),
982            DefaultFilter,
983            MemoryStorageLogger::new(16),
984        );
985
986        let mut buf = [0u8; 2];
987        stream.read_exact(&mut buf).unwrap();
988        assert_eq!(stream.get_log_records().len(), 1);
989
990        stream.clear_log_records();
991        assert!(stream.get_log_records().is_empty());
992    }
993
994    #[test]
995    fn test_channel_logger_delivers_records() {
996        let mut stream = LoggedStream::new(
997            Cursor::new(vec![0x01, 0x02]),
998            LowercaseHexadecimalFormatter::new_default(),
999            DefaultFilter,
1000            ChannelLogger::new(),
1001        );
1002        let receiver = stream.take_receiver_unchecked();
1003
1004        let mut buf = [0u8; 2];
1005        stream.read_exact(&mut buf).unwrap();
1006
1007        let record = receiver.recv().unwrap();
1008        assert_eq!(record.kind, RecordKind::Read);
1009        assert_eq!(record.message, "01:02");
1010    }
1011
1012    //////////////////////////////////////////////////////////////////////////////////////////////////////////
1013    // Trait assertions
1014    //////////////////////////////////////////////////////////////////////////////////////////////////////////
1015
1016    fn assert_send<T: Send>() {}
1017
1018    fn assert_unpin<T: Unpin>() {}
1019
1020    #[test]
1021    fn test_send() {
1022        assert_send::<
1023            LoggedStream<
1024                Cursor<Vec<u8>>,
1025                LowercaseHexadecimalFormatter,
1026                DefaultFilter,
1027                MemoryStorageLogger,
1028            >,
1029        >();
1030    }
1031
1032    #[test]
1033    fn test_unpin() {
1034        assert_unpin::<
1035            LoggedStream<
1036                Cursor<Vec<u8>>,
1037                LowercaseHexadecimalFormatter,
1038                DefaultFilter,
1039                MemoryStorageLogger,
1040            >,
1041        >();
1042    }
1043
1044    #[test]
1045    fn test_debug() {
1046        let stream = LoggedStream::new(
1047            Cursor::new(Vec::<u8>::new()),
1048            LowercaseHexadecimalFormatter::new_default(),
1049            DefaultFilter,
1050            MemoryStorageLogger::new(4),
1051        );
1052
1053        let debug = format!("{stream:?}");
1054        assert!(debug.contains("LoggedStream"));
1055        assert!(debug.contains("formatter"));
1056    }
1057}