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
17pub 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 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 pub fn log_open(&mut self, message: impl Into<String>) {
156 self.emit(Record::new(RecordKind::Open, message.into()));
157 }
158
159 #[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 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 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 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 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 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 #[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 #[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 #[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 #[test]
615 fn test_read_record_suppressed_by_filter() {
616 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 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 #[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 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 #[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 #[test]
813 fn test_drop_logs_drop_record() {
814 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 #[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 #[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 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}