Skip to main content

deaddrop_core/event/
mod.rs

1use crate::{ObjectId, PeerId};
2use serde::{Deserialize, Serialize};
3use std::sync::atomic::{AtomicU64, Ordering};
4use std::sync::{Arc, Mutex};
5
6#[derive(Debug, Clone, Serialize, Deserialize)]
7#[serde(tag = "event", rename_all = "snake_case")]
8pub enum Event {
9    PeerDiscovered {
10        peer: String,
11    },
12    PeerAuthenticated {
13        peer: String,
14    },
15    PeerDisconnected {
16        peer: String,
17    },
18    DropCreated {
19        object: String,
20    },
21    DropAccepted {
22        object: String,
23    },
24    DropForwarded {
25        object: String,
26        peer: String,
27    },
28    DropDelivered {
29        object: String,
30    },
31    DropReceived {
32        object: String,
33    },
34    DropExpired {
35        object: String,
36    },
37    TransferStarted {
38        object: String,
39    },
40    TransferPaused {
41        object: String,
42    },
43    TransferCompleted {
44        object: String,
45    },
46    RouteChanged {
47        object: String,
48        peer: String,
49    },
50    SpaceUpdated {
51        name: String,
52    },
53    ChannelUpdated {
54        name: String,
55    },
56    ChunkReceived {
57        object: String,
58        index: u32,
59    },
60    ChunkVerified {
61        object: String,
62        index: u32,
63    },
64    ChunkRejected {
65        object: String,
66        index: u32,
67    },
68    RouteSelected {
69        object: String,
70        peer: String,
71        score: f64,
72    },
73    RouteRejected {
74        object: String,
75        peer: String,
76    },
77    StoragePressure {
78        used: u64,
79        limit: u64,
80    },
81}
82
83impl Event {
84    pub fn peer_discovered(p: PeerId) -> Self {
85        Self::PeerDiscovered {
86            peer: p.to_string(),
87        }
88    }
89    pub fn drop_created(o: ObjectId) -> Self {
90        Self::DropCreated {
91            object: o.to_string(),
92        }
93    }
94    pub fn peer_connected(p: PeerId) -> Self {
95        Self::PeerAuthenticated {
96            peer: p.to_string(),
97        }
98    }
99}
100
101#[derive(Debug, Default)]
102pub struct Metrics {
103    pub peers_discovered: AtomicU64,
104    pub active_sessions: AtomicU64,
105    pub drops_created: AtomicU64,
106    pub drops_delivered: AtomicU64,
107    pub drops_relayed: AtomicU64,
108    pub chunks_transferred: AtomicU64,
109    pub bytes_transferred: AtomicU64,
110    pub route_decisions: AtomicU64,
111    pub failed_auth: AtomicU64,
112    pub expired: AtomicU64,
113}
114
115impl Metrics {
116    pub fn snapshot(&self) -> MetricsSnap {
117        MetricsSnap {
118            peers_discovered: self.peers_discovered.load(Ordering::Relaxed),
119            active_sessions: self.active_sessions.load(Ordering::Relaxed),
120            drops_created: self.drops_created.load(Ordering::Relaxed),
121            drops_delivered: self.drops_delivered.load(Ordering::Relaxed),
122            drops_relayed: self.drops_relayed.load(Ordering::Relaxed),
123            chunks_transferred: self.chunks_transferred.load(Ordering::Relaxed),
124            bytes_transferred: self.bytes_transferred.load(Ordering::Relaxed),
125            route_decisions: self.route_decisions.load(Ordering::Relaxed),
126            failed_auth: self.failed_auth.load(Ordering::Relaxed),
127            expired: self.expired.load(Ordering::Relaxed),
128        }
129    }
130}
131
132#[derive(Debug, Clone, Serialize)]
133pub struct MetricsSnap {
134    pub peers_discovered: u64,
135    pub active_sessions: u64,
136    pub drops_created: u64,
137    pub drops_delivered: u64,
138    pub drops_relayed: u64,
139    pub chunks_transferred: u64,
140    pub bytes_transferred: u64,
141    pub route_decisions: u64,
142    pub failed_auth: u64,
143    pub expired: u64,
144}
145
146#[derive(Clone, Default)]
147pub struct Bus {
148    events: Arc<Mutex<Vec<Event>>>,
149}
150
151impl Bus {
152    pub fn emit(&self, e: Event) {
153        tracing::info!(target: "ddp", event = ?e, "event");
154        self.events.lock().expect("bus").push(e);
155    }
156    pub fn take(&self) -> Vec<Event> {
157        std::mem::take(&mut *self.events.lock().expect("bus"))
158    }
159}