Skip to main content

stratum_apps/monitoring/
client.rs

1//! Sv2 client monitoring types
2//!
3//! These types are for monitoring **Sv2 clients** (downstream connections).
4//! Each client can have multiple channels opened with the app.
5
6use serde::{Deserialize, Serialize};
7use std::collections::HashMap;
8#[cfg(feature = "asic-rs-telemetry")]
9use std::net::IpAddr;
10use utoipa::ToSchema;
11
12#[cfg(feature = "asic-rs-telemetry")]
13use super::miner_telemetry::{MinerTelemetry, MinerTelemetryStatus};
14
15/// Kind of SV2 downstream client connected to this node.
16#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
17#[serde(rename_all = "snake_case")]
18pub enum Sv2ClientKind {
19    /// A mining device connected directly as an SV2 downstream client.
20    Miner,
21    /// A Translator Proxy connected as the SV2 client for one or more SV1 miners.
22    TranslatorProxy,
23    /// The downstream client type could not be inferred from SetupConnection metadata.
24    Unknown,
25}
26
27impl Default for Sv2ClientKind {
28    fn default() -> Self {
29        Self::Unknown
30    }
31}
32
33impl Sv2ClientKind {
34    /// Infer the downstream client kind from SetupConnection vendor and hardware version fields.
35    pub fn from_setup_connection(vendor: &str, hardware_version: &str) -> Self {
36        if vendor == "SRI" && hardware_version == "Translator Proxy" {
37            Self::TranslatorProxy
38        } else if vendor.is_empty() && hardware_version.is_empty() {
39            Self::Unknown
40        } else {
41            Self::Miner
42        }
43    }
44}
45
46/// Information about an extended channel
47#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
48pub struct ExtendedChannelInfo {
49    pub channel_id: u32,
50    pub user_identity: String,
51    pub nominal_hashrate: f32,
52    pub stable_hashrate: bool,
53    pub target_hex: String,
54    pub requested_max_target_hex: String,
55    pub extranonce_prefix_hex: String,
56    pub full_extranonce_size: usize,
57    pub rollable_extranonce_size: u16,
58    pub expected_shares_per_minute: f32,
59    pub shares_accepted: u32,
60    pub shares_rejected: u32,
61    pub shares_rejected_by_reason: HashMap<String, u32>,
62    pub share_work_sum: f64,
63    pub last_share_sequence_number: u32,
64    pub best_diff: f64,
65    pub last_batch_accepted: u32,
66    pub last_batch_work_sum: u64,
67    pub share_batch_size: usize,
68    pub blocks_found: u32,
69}
70
71/// Information about a standard channel
72#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
73pub struct StandardChannelInfo {
74    pub channel_id: u32,
75    pub user_identity: String,
76    pub nominal_hashrate: f32,
77    pub stable_hashrate: bool,
78    pub target_hex: String,
79    pub requested_max_target_hex: String,
80    pub extranonce_prefix_hex: String,
81    pub expected_shares_per_minute: f32,
82    pub shares_accepted: u32,
83    pub shares_rejected: u32,
84    pub shares_rejected_by_reason: HashMap<String, u32>,
85    pub share_work_sum: f64,
86    pub last_share_sequence_number: u32,
87    pub best_diff: f64,
88    pub last_batch_accepted: u32,
89    pub last_batch_work_sum: u64,
90    pub share_batch_size: usize,
91    pub blocks_found: u32,
92}
93
94/// Full information about a single Sv2 client including all channels
95#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
96pub struct Sv2ClientInfo {
97    pub client_id: usize,
98    /// Classification inferred from the client's SV2 SetupConnection metadata.
99    pub client_kind: Sv2ClientKind,
100    pub extended_channels: Vec<ExtendedChannelInfo>,
101    pub standard_channels: Vec<StandardChannelInfo>,
102    #[cfg(feature = "asic-rs-telemetry")]
103    /// Miner management IP used for matched telemetry, if discovery found one.
104    #[schema(value_type = Option<String>)]
105    pub management_ip: Option<IpAddr>,
106    #[cfg(feature = "asic-rs-telemetry")]
107    /// Latest telemetry fetched from the matched miner management interface.
108    pub miner_telemetry: Option<MinerTelemetry>,
109    #[cfg(feature = "asic-rs-telemetry")]
110    /// Current discovery and fetch status for miner telemetry matching.
111    pub miner_telemetry_status: Option<MinerTelemetryStatus>,
112}
113
114impl Sv2ClientInfo {
115    pub fn new(
116        client_id: usize,
117        extended_channels: Vec<ExtendedChannelInfo>,
118        standard_channels: Vec<StandardChannelInfo>,
119    ) -> Self {
120        Self {
121            client_id,
122            client_kind: Sv2ClientKind::Unknown,
123            extended_channels,
124            standard_channels,
125            #[cfg(feature = "asic-rs-telemetry")]
126            management_ip: None,
127            #[cfg(feature = "asic-rs-telemetry")]
128            miner_telemetry: None,
129            #[cfg(feature = "asic-rs-telemetry")]
130            miner_telemetry_status: None,
131        }
132    }
133
134    /// Get total number of channels for this client
135    pub fn total_channels(&self) -> usize {
136        self.extended_channels.len() + self.standard_channels.len()
137    }
138
139    /// Get total hashrate for this client
140    pub fn total_hashrate(&self) -> f32 {
141        self.extended_channels
142            .iter()
143            .map(|c| c.nominal_hashrate)
144            .sum::<f32>()
145            + self
146                .standard_channels
147                .iter()
148                .map(|c| c.nominal_hashrate)
149                .sum::<f32>()
150    }
151
152    /// Convert to metadata (without channel arrays)
153    pub fn to_metadata(&self) -> Sv2ClientMetadata {
154        Sv2ClientMetadata {
155            client_id: self.client_id,
156            client_kind: self.client_kind,
157            extended_channels_count: self.extended_channels.len(),
158            standard_channels_count: self.standard_channels.len(),
159            total_hashrate: self.total_hashrate(),
160            #[cfg(feature = "asic-rs-telemetry")]
161            management_ip: self.management_ip,
162            #[cfg(feature = "asic-rs-telemetry")]
163            miner_telemetry: self.miner_telemetry.clone(),
164            #[cfg(feature = "asic-rs-telemetry")]
165            miner_telemetry_status: self.miner_telemetry_status,
166        }
167    }
168}
169
170/// Sv2 client metadata without channel arrays (for listings)
171#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
172pub struct Sv2ClientMetadata {
173    pub client_id: usize,
174    /// Classification inferred from the client's SV2 SetupConnection metadata.
175    pub client_kind: Sv2ClientKind,
176    pub extended_channels_count: usize,
177    pub standard_channels_count: usize,
178    pub total_hashrate: f32,
179    #[cfg(feature = "asic-rs-telemetry")]
180    /// Miner management IP used for matched telemetry, if discovery found one.
181    #[schema(value_type = Option<String>)]
182    pub management_ip: Option<IpAddr>,
183    #[cfg(feature = "asic-rs-telemetry")]
184    /// Latest telemetry fetched from the matched miner management interface.
185    pub miner_telemetry: Option<MinerTelemetry>,
186    #[cfg(feature = "asic-rs-telemetry")]
187    /// Current discovery and fetch status for miner telemetry matching.
188    pub miner_telemetry_status: Option<MinerTelemetryStatus>,
189}
190
191/// Aggregate information about all Sv2 clients
192#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
193pub struct Sv2ClientsSummary {
194    pub total_clients: usize,
195    pub total_channels: usize,
196    pub extended_channels: usize,
197    pub standard_channels: usize,
198    pub total_hashrate: f32,
199}
200
201/// Trait for monitoring Sv2 clients (downstream connections)
202pub trait Sv2ClientsMonitoring: Send + Sync {
203    /// Get all Sv2 clients with their channels
204    fn get_sv2_clients(&self) -> Vec<Sv2ClientInfo>;
205
206    /// Get a single Sv2 client by client_id
207    ///
208    /// Default implementation does O(n) scan. Override for O(1) lookup
209    /// if your implementation uses a HashMap internally.
210    fn get_sv2_client_by_id(&self, client_id: usize) -> Option<Sv2ClientInfo> {
211        self.get_sv2_clients()
212            .into_iter()
213            .find(|c| c.client_id == client_id)
214    }
215
216    /// Get summary of all Sv2 clients
217    fn get_sv2_clients_summary(&self) -> Sv2ClientsSummary {
218        let clients = self.get_sv2_clients();
219        let extended: usize = clients.iter().map(|c| c.extended_channels.len()).sum();
220        let standard: usize = clients.iter().map(|c| c.standard_channels.len()).sum();
221
222        Sv2ClientsSummary {
223            total_clients: clients.len(),
224            total_channels: extended + standard,
225            extended_channels: extended,
226            standard_channels: standard,
227            total_hashrate: clients.iter().map(|c| c.total_hashrate()).sum(),
228        }
229    }
230}
231
232#[cfg(test)]
233mod tests {
234    use super::*;
235    use stratum_core::mining_sv2::ERROR_CODE_SUBMIT_SHARES_DUPLICATE_SHARE;
236
237    // ── helpers ──────────────────────────────────────────────────────
238
239    fn create_extended_channel_info(channel_id: u32, hashrate: f32) -> ExtendedChannelInfo {
240        ExtendedChannelInfo {
241            channel_id,
242            user_identity: format!("user-ext-{}", channel_id),
243            nominal_hashrate: hashrate,
244            stable_hashrate: false,
245            target_hex: "00ff".into(),
246            requested_max_target_hex: "00ff".into(),
247            extranonce_prefix_hex: "aa".into(),
248            full_extranonce_size: 16,
249            rollable_extranonce_size: 4,
250            expected_shares_per_minute: 1.0,
251            shares_accepted: 10,
252            shares_rejected: 0,
253            shares_rejected_by_reason: HashMap::new(),
254            share_work_sum: 100.0,
255            last_share_sequence_number: 5,
256            best_diff: 50.0,
257            last_batch_accepted: 3,
258            last_batch_work_sum: 30,
259            share_batch_size: 10,
260            blocks_found: 0,
261        }
262    }
263
264    fn create_standard_channel_info(channel_id: u32, hashrate: f32) -> StandardChannelInfo {
265        StandardChannelInfo {
266            channel_id,
267            user_identity: format!("user-std-{}", channel_id),
268            nominal_hashrate: hashrate,
269            stable_hashrate: false,
270            target_hex: "00ff".into(),
271            requested_max_target_hex: "00ff".into(),
272            extranonce_prefix_hex: "bb".into(),
273            expected_shares_per_minute: 2.0,
274            shares_accepted: 20,
275            shares_rejected: 1,
276            shares_rejected_by_reason: HashMap::from([(
277                ERROR_CODE_SUBMIT_SHARES_DUPLICATE_SHARE.to_string(),
278                1,
279            )]),
280            share_work_sum: 200.0,
281            last_share_sequence_number: 8,
282            best_diff: 80.0,
283            last_batch_accepted: 5,
284            last_batch_work_sum: 50,
285            share_batch_size: 20,
286            blocks_found: 0,
287        }
288    }
289
290    fn create_sv2_client_info(
291        id: usize,
292        ext: Vec<ExtendedChannelInfo>,
293        std: Vec<StandardChannelInfo>,
294    ) -> Sv2ClientInfo {
295        Sv2ClientInfo {
296            client_id: id,
297            client_kind: Sv2ClientKind::Miner,
298            extended_channels: ext,
299            standard_channels: std,
300            #[cfg(feature = "asic-rs-telemetry")]
301            management_ip: None,
302            #[cfg(feature = "asic-rs-telemetry")]
303            miner_telemetry: None,
304            #[cfg(feature = "asic-rs-telemetry")]
305            miner_telemetry_status: None,
306        }
307    }
308
309    // ── ClientInfo unit tests ───────────────────────────────────────
310
311    #[test]
312    fn client_info_empty_channels() {
313        let client = create_sv2_client_info(1, vec![], vec![]);
314        assert_eq!(client.total_channels(), 0);
315        assert_eq!(client.total_hashrate(), 0.0);
316    }
317
318    #[test]
319    fn client_info_aggregates_both_channel_types() {
320        let client = create_sv2_client_info(
321            1,
322            vec![
323                create_extended_channel_info(1, 100.0),
324                create_extended_channel_info(2, 200.0),
325            ],
326            vec![create_standard_channel_info(3, 50.0)],
327        );
328        assert_eq!(client.total_channels(), 3);
329        assert_eq!(client.total_hashrate(), 350.0);
330    }
331
332    #[test]
333    fn client_info_to_metadata() {
334        let client = create_sv2_client_info(
335            42,
336            vec![create_extended_channel_info(1, 100.0)],
337            vec![
338                create_standard_channel_info(2, 50.0),
339                create_standard_channel_info(3, 75.0),
340            ],
341        );
342        let meta = client.to_metadata();
343
344        assert_eq!(meta.client_id, 42);
345        assert_eq!(meta.extended_channels_count, 1);
346        assert_eq!(meta.standard_channels_count, 2);
347        assert_eq!(meta.total_hashrate, 225.0);
348        assert_eq!(meta.client_kind, Sv2ClientKind::Miner);
349        #[cfg(feature = "asic-rs-telemetry")]
350        assert!(meta.miner_telemetry.is_none());
351    }
352
353    #[test]
354    fn client_kind_from_setup_connection_classifies_known_clients() {
355        assert_eq!(
356            Sv2ClientKind::from_setup_connection("SRI", "Translator Proxy"),
357            Sv2ClientKind::TranslatorProxy
358        );
359        assert_eq!(
360            Sv2ClientKind::from_setup_connection("Bitaxe", "Gamma"),
361            Sv2ClientKind::Miner
362        );
363        assert_eq!(
364            Sv2ClientKind::from_setup_connection("", ""),
365            Sv2ClientKind::Unknown
366        );
367    }
368
369    // ── ClientsMonitoring trait default implementations ─────────────
370
371    struct MockClients(Vec<Sv2ClientInfo>);
372    impl Sv2ClientsMonitoring for MockClients {
373        fn get_sv2_clients(&self) -> Vec<Sv2ClientInfo> {
374            self.0.clone()
375        }
376    }
377
378    #[test]
379    fn clients_monitoring_get_client_by_id_found() {
380        let monitor = MockClients(vec![
381            create_sv2_client_info(1, vec![create_extended_channel_info(1, 10.0)], vec![]),
382            create_sv2_client_info(2, vec![], vec![create_standard_channel_info(1, 20.0)]),
383        ]);
384        let found = monitor.get_sv2_client_by_id(2);
385        assert!(found.is_some());
386        assert_eq!(found.unwrap().client_id, 2);
387    }
388
389    #[test]
390    fn clients_monitoring_get_client_by_id_not_found() {
391        let monitor = MockClients(vec![create_sv2_client_info(1, vec![], vec![])]);
392        assert!(monitor.get_sv2_client_by_id(999).is_none());
393    }
394
395    #[test]
396    fn clients_monitoring_summary_empty() {
397        let monitor = MockClients(vec![]);
398        let summary = monitor.get_sv2_clients_summary();
399
400        assert_eq!(summary.total_clients, 0);
401        assert_eq!(summary.total_channels, 0);
402        assert_eq!(summary.extended_channels, 0);
403        assert_eq!(summary.standard_channels, 0);
404        assert_eq!(summary.total_hashrate, 0.0);
405    }
406
407    #[test]
408    fn clients_monitoring_summary_aggregates_correctly() {
409        let monitor = MockClients(vec![
410            create_sv2_client_info(
411                1,
412                vec![create_extended_channel_info(1, 100.0)],
413                vec![create_standard_channel_info(2, 50.0)],
414            ),
415            create_sv2_client_info(
416                2,
417                vec![
418                    create_extended_channel_info(3, 200.0),
419                    create_extended_channel_info(4, 300.0),
420                ],
421                vec![],
422            ),
423        ]);
424        let summary = monitor.get_sv2_clients_summary();
425
426        assert_eq!(summary.total_clients, 2);
427        assert_eq!(summary.extended_channels, 3);
428        assert_eq!(summary.standard_channels, 1);
429        assert_eq!(summary.total_channels, 4);
430        assert_eq!(summary.total_hashrate, 650.0);
431    }
432}