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}