media_pp\elements\sink/
packet_counter.rs1use std::sync::{
2 Arc,
3 atomic::{AtomicUsize, Ordering},
4};
5
6use crate::pp_log::{PpLog, pp_info};
7
8use crate::{
9 buffer::MediaBuffer,
10 control::ControlMsg,
11 element::{Element, ElementType, Sink, element_pp_log},
12 error::Result,
13};
14
15pub struct PacketCounter {
19 pp_log: PpLog,
20 name: Arc<str>,
21 count: Arc<AtomicUsize>,
22}
23
24impl PacketCounter {
25 pub fn new(name: impl Into<String>) -> (Self, Arc<AtomicUsize>) {
26 let count = Arc::new(AtomicUsize::new(0));
27 let name: Arc<str> = name.into().into();
28 let pp_log = element_pp_log(ElementType::PacketCounter, &name, None);
29 pp_info!(pp_log: &pp_log, "created");
30 (
31 Self {
32 name,
33 pp_log,
34 count: count.clone(),
35 },
36 count,
37 )
38 }
39}
40
41impl Element for PacketCounter {
42 fn name(&self) -> Arc<str> {
43 self.name.clone()
44 }
45
46 fn element_type(&self) -> ElementType {
47 ElementType::PacketCounter
48 }
49
50 fn pp_log(&self) -> &PpLog {
51 &self.pp_log
52 }
53
54 fn pp_log_mut(&mut self) -> &mut PpLog {
55 &mut self.pp_log
56 }
57}
58
59impl Sink for PacketCounter {
60 fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
61 if let MediaBuffer::Packet(_) = buf {
62 self.count.fetch_add(1, Ordering::Relaxed);
63 }
64 Ok(())
65 }
66
67 fn control(&mut self, _msg: ControlMsg) -> Result<()> {
68 Ok(())
70 }
71}