1use 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#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)]
25pub struct MinerTelemetry {
26 pub make: Option<String>,
28 pub model: Option<String>,
30 pub firmware_version: Option<String>,
32 pub reported_hashrate_hs: Option<f64>,
34 pub power_consumption_w: Option<f64>,
36 pub efficiency_j_per_th: Option<f64>,
38 pub average_temperature_c: Option<f64>,
40 pub uptime_secs: Option<u64>,
42 pub is_mining: Option<bool>,
44}
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
48#[serde(rename_all = "snake_case")]
49pub enum MinerTelemetryStatus {
50 Matched,
52 Unmatched,
54 DuplicateWorkerName,
56 FetchFailed,
58}
59
60#[derive(Debug, Clone, Default)]
62pub struct MinerTelemetryDownstreamMatches {
63 pub management_ips_by_downstream_id: HashMap<DownstreamId, IpAddr>,
65 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
106pub struct MinerTelemetryCollector {
108 factory: MinerFactory,
109}
110
111#[derive(Debug, Clone)]
113pub struct DiscoveredMiner {
114 pub ip: IpAddr,
116 pub pools: Vec<DiscoveredMinerPool>,
118}
119
120#[derive(Debug, Clone, PartialEq, Eq)]
122pub struct DiscoveredMinerPool {
123 pub user: String,
125 pub host: String,
127 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}