Skip to main content

stratum_apps/monitoring/
miner_telemetry.rs

1//! Miner telemetry monitoring types and helpers.
2
3use crate::utils::types::DownstreamId;
4use asic_rs::{
5    core::data::{collector::DataField, hashrate::HashRateUnit, miner::MinerData, pool::PoolData},
6    MinerFactory,
7};
8use futures::{stream::FuturesUnordered, StreamExt};
9use serde::{Deserialize, Serialize};
10use std::{
11    collections::HashMap,
12    net::{IpAddr, Ipv4Addr},
13    time::Duration,
14};
15use tokio::time::timeout;
16use tracing::{debug, warn};
17use utoipa::ToSchema;
18
19const MINER_DISCOVERY_PROBE_TIMEOUT: Duration = Duration::from_secs(2);
20const MINER_DISCOVERY_MAX_CONCURRENCY: usize = 64;
21const MINER_DISCOVERY_MIN_IPV4_PREFIX: u8 = 24;
22
23/// Telemetry reported by the miner's management interface.
24#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
25pub struct MinerTelemetry {
26    /// Miner manufacturer or brand reported by the management interface.
27    pub make: Option<String>,
28    /// Miner model reported by the management interface.
29    pub model: Option<String>,
30    /// Firmware version reported by the miner, when exposed.
31    pub firmware_version: Option<String>,
32    /// Miner-reported hashrate in hashes per second.
33    pub reported_hashrate_hs: Option<f64>,
34    /// Current miner power consumption in watts.
35    pub power_consumption_w: Option<f64>,
36    /// Miner efficiency in joules per terahash.
37    pub efficiency_j_per_th: Option<f64>,
38    /// Average miner temperature in degrees Celsius.
39    pub average_temperature_c: Option<f64>,
40    /// Total miner system uptime in seconds.
41    pub uptime_secs: Option<u64>,
42    /// Whether the miner reports that hashing is currently running.
43    pub is_mining: Option<bool>,
44}
45
46/// Status of matching a connected miner to discovered management telemetry.
47#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
48#[serde(rename_all = "snake_case")]
49pub enum MinerTelemetryStatus {
50    /// Telemetry was matched to this connected miner.
51    Matched,
52    /// No discovered miner management interface matched this connection.
53    Unmatched,
54    /// More than one connected miner used the same worker name, so telemetry was not assigned.
55    DuplicateWorkerName,
56    /// A matching miner was found, but fetching telemetry from its management interface failed.
57    FetchFailed,
58}
59
60/// Matches between active downstream connections and discovered miner management interfaces.
61#[derive(Debug, Clone, Default)]
62pub struct MinerTelemetryDownstreamMatches {
63    /// Matched miner management IP for each downstream connection id.
64    pub management_ips_by_downstream_id: HashMap<DownstreamId, IpAddr>,
65    /// Telemetry matching status for each active downstream connection id.
66    pub statuses_by_downstream_id: HashMap<DownstreamId, MinerTelemetryStatus>,
67}
68
69impl From<MinerData> for MinerTelemetry {
70    fn from(data: MinerData) -> Self {
71        Self {
72            make: Some(data.device_info.make),
73            model: Some(data.device_info.model),
74            firmware_version: data.firmware_version,
75            reported_hashrate_hs: data
76                .hashrate
77                .map(|hashrate| hashrate.as_unit(HashRateUnit::Hash).value),
78            power_consumption_w: data.wattage.map(|power| power.as_watts()),
79            efficiency_j_per_th: data.efficiency,
80            average_temperature_c: data
81                .average_temperature
82                .map(|temperature| temperature.as_celsius()),
83            uptime_secs: data.uptime.map(|uptime| uptime.as_secs()),
84            is_mining: Some(data.is_mining),
85        }
86    }
87}
88
89const MINER_TELEMETRY_EXCLUDED_FIELDS: &[DataField] = &[
90    DataField::Mac,
91    DataField::SerialNumber,
92    DataField::Hostname,
93    DataField::ApiVersion,
94    DataField::ControlBoardVersion,
95    DataField::Chips,
96    DataField::ExpectedHashrate,
97    DataField::Fans,
98    DataField::PsuFans,
99    DataField::FluidTemperature,
100    DataField::TuningTarget,
101    DataField::LightFlashing,
102    DataField::Messages,
103    DataField::Pools,
104];
105
106/// Collects miner telemetry and discovers miner management interfaces on configured LAN ranges.
107pub struct MinerTelemetryCollector {
108    factory: MinerFactory,
109}
110
111/// Miner management interface discovered on the LAN.
112#[derive(Debug, Clone)]
113pub struct DiscoveredMiner {
114    /// Management IP address of the discovered miner.
115    pub ip: IpAddr,
116    /// Pools configured on the discovered miner.
117    pub pools: Vec<DiscoveredMinerPool>,
118}
119
120/// Pool configuration reported by a discovered miner.
121#[derive(Debug, Clone, PartialEq, Eq)]
122pub struct DiscoveredMinerPool {
123    /// Worker name or user configured for the pool.
124    pub user: String,
125    /// Pool host configured on the miner.
126    pub host: String,
127    /// Pool port configured on the miner.
128    pub port: u16,
129}
130
131#[derive(Debug, Clone, Copy)]
132struct Ipv4Cidr {
133    network: Ipv4Addr,
134    prefix: u8,
135}
136
137impl MinerTelemetryCollector {
138    pub fn new() -> Self {
139        Self {
140            factory: MinerFactory::new(),
141        }
142    }
143
144    pub async fn fetch(&self, ip: IpAddr) -> Option<MinerTelemetry> {
145        let miner = match self.factory.get_miner(ip).await {
146            Ok(Some(miner)) => miner,
147            Ok(None) => {
148                debug!("No miner management interface found at {ip}");
149                return None;
150            }
151            Err(error) => {
152                debug!("Failed to get miner management interface at {ip}: {error}");
153                return None;
154            }
155        };
156
157        Some(MinerTelemetry::from(
158            miner
159                .get_data_filtered(MINER_TELEMETRY_EXCLUDED_FIELDS.to_vec())
160                .await,
161        ))
162    }
163
164    pub async fn discover(&self, cidrs: &[String]) -> Vec<DiscoveredMiner> {
165        let ips = cidrs
166            .iter()
167            .filter_map(|cidr| parse_private_ipv4_cidr(cidr))
168            .flat_map(|cidr| cidr.host_ips())
169            .collect::<Vec<_>>();
170
171        debug!(
172            cidrs = ?cidrs,
173            hosts = ips.len(),
174            "Starting miner telemetry discovery scan"
175        );
176
177        let mut discovered = Vec::new();
178
179        for chunk in ips.chunks(MINER_DISCOVERY_MAX_CONCURRENCY) {
180            let mut probes = FuturesUnordered::new();
181            for ip in chunk {
182                probes.push(self.discover_ip(*ip));
183            }
184
185            while let Some(result) = probes.next().await {
186                if let Some(miner) = result {
187                    discovered.push(miner);
188                }
189            }
190        }
191
192        discovered
193    }
194
195    async fn discover_ip(&self, ip: IpAddr) -> Option<DiscoveredMiner> {
196        let miner = match timeout(MINER_DISCOVERY_PROBE_TIMEOUT, self.factory.get_miner(ip)).await {
197            Ok(Ok(Some(miner))) => miner,
198            Ok(Ok(None)) => return None,
199            Ok(Err(error)) => {
200                debug!("Failed to discover miner management interface at {ip}: {error}");
201                return None;
202            }
203            Err(_) => return None,
204        };
205
206        let pools = match timeout(MINER_DISCOVERY_PROBE_TIMEOUT, miner.get_pools()).await {
207            Ok(pools) => pools,
208            Err(_) => {
209                debug!("Timed out reading pool users from miner management interface at {ip}");
210                Vec::new()
211            }
212        };
213
214        let mut miner_pools = pools
215            .into_iter()
216            .flat_map(|group| group.pools)
217            .filter_map(discovered_miner_pool)
218            .collect::<Vec<_>>();
219        miner_pools.sort_by(|a, b| {
220            a.user
221                .cmp(&b.user)
222                .then_with(|| a.host.cmp(&b.host))
223                .then_with(|| a.port.cmp(&b.port))
224        });
225        miner_pools.dedup();
226
227        debug!(
228            pools = ?miner_pools,
229            "Discovered miner management interface at {ip}"
230        );
231
232        Some(DiscoveredMiner {
233            ip,
234            pools: miner_pools,
235        })
236    }
237}
238
239fn discovered_miner_pool(pool: PoolData) -> Option<DiscoveredMinerPool> {
240    if pool.active == Some(false) {
241        return None;
242    }
243
244    let user = pool.user?.trim().to_owned();
245    if user.is_empty() {
246        return None;
247    }
248
249    let url = pool.url?;
250    Some(DiscoveredMinerPool {
251        user,
252        host: url.host,
253        port: url.port,
254    })
255}
256
257pub fn match_discovered_miners_to_downstreams_by_worker_and_port(
258    downstream_workers: &[(DownstreamId, String)],
259    discovered_miners: &[DiscoveredMiner],
260    expected_pool_port: u16,
261) -> MinerTelemetryDownstreamMatches {
262    let mut result = MinerTelemetryDownstreamMatches::default();
263    let mut downstream_worker_counts = HashMap::new();
264
265    for (_, worker_name) in downstream_workers {
266        if !worker_name.is_empty() {
267            *downstream_worker_counts
268                .entry(worker_name.as_str())
269                .or_insert(0usize) += 1;
270        }
271    }
272
273    for (downstream_id, worker_name) in downstream_workers {
274        if worker_name.is_empty() {
275            result
276                .statuses_by_downstream_id
277                .insert(*downstream_id, MinerTelemetryStatus::Unmatched);
278            continue;
279        }
280
281        let matches = discovered_miners
282            .iter()
283            .filter(|miner| {
284                miner.pools.iter().any(|pool| {
285                    pool.user == worker_name.as_str() && pool.port == expected_pool_port
286                })
287            })
288            .collect::<Vec<_>>();
289
290        if downstream_worker_counts
291            .get(worker_name.as_str())
292            .copied()
293            .unwrap_or_default()
294            > 1
295        {
296            result
297                .statuses_by_downstream_id
298                .insert(*downstream_id, MinerTelemetryStatus::DuplicateWorkerName);
299            debug!(
300                "Miner telemetry discovery found multiple active downstreams for worker {worker_name}; leaving downstream {downstream_id} unmatched"
301            );
302            continue;
303        }
304
305        if matches.len() == 1 {
306            result
307                .management_ips_by_downstream_id
308                .insert(*downstream_id, matches[0].ip);
309            result
310                .statuses_by_downstream_id
311                .insert(*downstream_id, MinerTelemetryStatus::Matched);
312        } else if matches.len() > 1 {
313            result
314                .statuses_by_downstream_id
315                .insert(*downstream_id, MinerTelemetryStatus::DuplicateWorkerName);
316            debug!(
317                "Miner telemetry discovery found multiple miners for worker {worker_name}; leaving downstream {downstream_id} unmatched"
318            );
319        } else {
320            result
321                .statuses_by_downstream_id
322                .insert(*downstream_id, MinerTelemetryStatus::Unmatched);
323        }
324    }
325
326    result
327}
328
329fn parse_private_ipv4_cidr(value: &str) -> Option<Ipv4Cidr> {
330    let (addr, prefix) = match value.split_once('/') {
331        Some(parts) => parts,
332        None => {
333            warn!("Ignoring miner telemetry CIDR {value}: expected IPv4 CIDR notation");
334            return None;
335        }
336    };
337
338    let addr = match addr.parse::<Ipv4Addr>() {
339        Ok(addr) => addr,
340        Err(error) => {
341            warn!("Ignoring miner telemetry CIDR {value}: invalid IPv4 address: {error}");
342            return None;
343        }
344    };
345
346    let prefix = match prefix.parse::<u8>() {
347        Ok(prefix) => prefix,
348        Err(error) => {
349            warn!("Ignoring miner telemetry CIDR {value}: invalid prefix: {error}");
350            return None;
351        }
352    };
353
354    if !(MINER_DISCOVERY_MIN_IPV4_PREFIX..=32).contains(&prefix) {
355        warn!(
356            "Ignoring miner telemetry CIDR {value}: prefix must be /{MINER_DISCOVERY_MIN_IPV4_PREFIX} or narrower"
357        );
358        return None;
359    }
360
361    if !addr.is_private() {
362        warn!("Ignoring miner telemetry CIDR {value}: only private IPv4 ranges are supported");
363        return None;
364    }
365
366    let network = Ipv4Addr::from(u32::from(addr) & ipv4_mask(prefix));
367    Some(Ipv4Cidr { network, prefix })
368}
369
370fn ipv4_mask(prefix: u8) -> u32 {
371    if prefix == 0 {
372        0
373    } else {
374        u32::MAX << (32 - prefix)
375    }
376}
377
378impl Ipv4Cidr {
379    fn host_ips(self) -> Vec<IpAddr> {
380        let network = u32::from(self.network);
381        let host_count = 1u64 << (32 - self.prefix);
382        let broadcast = network + host_count as u32 - 1;
383
384        let (first, last) = if self.prefix <= 30 {
385            (network + 1, broadcast - 1)
386        } else {
387            (network, broadcast)
388        };
389
390        (first..=last)
391            .map(|ip| IpAddr::V4(Ipv4Addr::from(ip)))
392            .collect()
393    }
394}
395
396impl Default for MinerTelemetryCollector {
397    fn default() -> Self {
398        Self::new()
399    }
400}
401
402#[cfg(test)]
403mod tests {
404    use super::*;
405    use asic_rs::core::data::pool::{PoolScheme, PoolURL};
406
407    fn discovered_miner(ip: [u8; 4], user: &str, port: u16) -> DiscoveredMiner {
408        DiscoveredMiner {
409            ip: IpAddr::V4(Ipv4Addr::from(ip)),
410            pools: vec![DiscoveredMinerPool {
411                user: user.to_string(),
412                host: "192.168.1.10".to_string(),
413                port,
414            }],
415        }
416    }
417
418    fn pool_data(user: &str, port: u16, active: Option<bool>) -> PoolData {
419        PoolData {
420            position: Some(0),
421            url: Some(PoolURL {
422                scheme: PoolScheme::StratumV1,
423                host: "192.168.1.10".to_string(),
424                port,
425                pubkey: None,
426            }),
427            accepted_shares: None,
428            rejected_shares: None,
429            active,
430            alive: None,
431            user: Some(user.to_string()),
432        }
433    }
434
435    #[test]
436    fn parses_private_ipv4_cidr_hosts() {
437        let cidr = parse_private_ipv4_cidr("192.168.1.0/30").unwrap();
438        let hosts = cidr.host_ips();
439
440        assert_eq!(
441            hosts,
442            vec![
443                IpAddr::V4(Ipv4Addr::new(192, 168, 1, 1)),
444                IpAddr::V4(Ipv4Addr::new(192, 168, 1, 2))
445            ]
446        );
447    }
448
449    #[test]
450    fn rejects_non_private_or_broad_cidrs() {
451        assert!(parse_private_ipv4_cidr("8.8.8.0/24").is_none());
452        assert!(parse_private_ipv4_cidr("192.168.0.0/16").is_none());
453        assert!(parse_private_ipv4_cidr("192.168.1.0").is_none());
454    }
455
456    #[test]
457    fn ignores_inactive_discovered_pool_entries() {
458        assert_eq!(
459            discovered_miner_pool(pool_data("worker-a", 34255, Some(false))),
460            None
461        );
462        assert_eq!(
463            discovered_miner_pool(pool_data("worker-a", 34255, Some(true))),
464            Some(DiscoveredMinerPool {
465                user: "worker-a".to_string(),
466                host: "192.168.1.10".to_string(),
467                port: 34255
468            })
469        );
470    }
471
472    #[test]
473    fn matches_unique_worker_names() {
474        let downstream_workers = vec![(10, "worker-a".to_string()), (11, "worker-b".to_string())];
475        let discovered_miners = vec![
476            discovered_miner([192, 168, 1, 20], "worker-a", 34255),
477            discovered_miner([192, 168, 1, 21], "worker-b", 34255),
478        ];
479
480        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
481            &downstream_workers,
482            &discovered_miners,
483            34255,
484        );
485
486        assert_eq!(
487            result.management_ips_by_downstream_id.get(&10),
488            Some(&IpAddr::V4(Ipv4Addr::new(192, 168, 1, 20)))
489        );
490        assert_eq!(
491            result.management_ips_by_downstream_id.get(&11),
492            Some(&IpAddr::V4(Ipv4Addr::new(192, 168, 1, 21)))
493        );
494    }
495
496    #[test]
497    fn reports_duplicate_worker_status_for_discovered_duplicates() {
498        let downstream_workers = vec![(10, "worker-a".to_string())];
499        let discovered_miners = vec![
500            discovered_miner([192, 168, 1, 20], "worker-a", 34255),
501            discovered_miner([192, 168, 1, 21], "worker-a", 34255),
502        ];
503
504        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
505            &downstream_workers,
506            &discovered_miners,
507            34255,
508        );
509
510        assert!(result.management_ips_by_downstream_id.is_empty());
511        assert_eq!(
512            result.statuses_by_downstream_id.get(&10),
513            Some(&MinerTelemetryStatus::DuplicateWorkerName)
514        );
515    }
516
517    #[test]
518    fn leaves_duplicate_downstream_worker_unmatched() {
519        let downstream_workers = vec![(10, "worker-a".to_string()), (11, "worker-a".to_string())];
520        let discovered_miners = vec![discovered_miner([192, 168, 1, 20], "worker-a", 34255)];
521
522        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
523            &downstream_workers,
524            &discovered_miners,
525            34255,
526        );
527
528        assert!(result.management_ips_by_downstream_id.is_empty());
529    }
530
531    #[test]
532    fn reports_unmatched_worker_status() {
533        let downstream_workers = vec![(10, "worker-a".to_string()), (11, String::new())];
534        let discovered_miners = vec![discovered_miner([192, 168, 1, 20], "worker-b", 34255)];
535
536        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
537            &downstream_workers,
538            &discovered_miners,
539            34255,
540        );
541
542        assert!(result.management_ips_by_downstream_id.is_empty());
543        assert_eq!(
544            result.statuses_by_downstream_id.get(&10),
545            Some(&MinerTelemetryStatus::Unmatched)
546        );
547        assert_eq!(
548            result.statuses_by_downstream_id.get(&11),
549            Some(&MinerTelemetryStatus::Unmatched)
550        );
551    }
552
553    #[test]
554    fn reports_duplicate_downstream_worker_status() {
555        let downstream_workers = vec![(10, "worker-a".to_string()), (11, "worker-a".to_string())];
556        let discovered_miners = vec![discovered_miner([192, 168, 1, 63], "worker-a", 34255)];
557
558        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
559            &downstream_workers,
560            &discovered_miners,
561            34255,
562        );
563
564        assert!(result.management_ips_by_downstream_id.is_empty());
565        assert_eq!(
566            result.statuses_by_downstream_id.get(&10),
567            Some(&MinerTelemetryStatus::DuplicateWorkerName)
568        );
569        assert_eq!(
570            result.statuses_by_downstream_id.get(&11),
571            Some(&MinerTelemetryStatus::DuplicateWorkerName)
572        );
573    }
574
575    #[test]
576    fn ignores_miner_with_matching_worker_on_different_pool_port() {
577        let downstream_workers = vec![(10, "worker-a".to_string())];
578        let discovered_miners = vec![discovered_miner([192, 168, 1, 63], "worker-a", 34265)];
579
580        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
581            &downstream_workers,
582            &discovered_miners,
583            34255,
584        );
585
586        assert!(result.management_ips_by_downstream_id.is_empty());
587        assert_eq!(
588            result.statuses_by_downstream_id.get(&10),
589            Some(&MinerTelemetryStatus::Unmatched)
590        );
591    }
592
593    #[test]
594    fn matches_worker_on_expected_pool_port() {
595        let downstream_workers = vec![(10, "worker-a".to_string())];
596        let discovered_miners = vec![discovered_miner([192, 168, 1, 63], "worker-a", 34255)];
597
598        let result = match_discovered_miners_to_downstreams_by_worker_and_port(
599            &downstream_workers,
600            &discovered_miners,
601            34255,
602        );
603
604        assert_eq!(
605            result.management_ips_by_downstream_id.get(&10),
606            Some(&IpAddr::V4(Ipv4Addr::new(192, 168, 1, 63)))
607        );
608        assert_eq!(
609            result.statuses_by_downstream_id.get(&10),
610            Some(&MinerTelemetryStatus::Matched)
611        );
612    }
613}