From 147fece0e272ae6fe4f925ad99489deac9f6ddf4 Mon Sep 17 00:00:00 2001 From: Brett Rowan <121075405+b-rowan@users.noreply.github.com> Date: Mon, 20 Jul 2026 14:06:35 -0600 Subject: [PATCH] feat: add miner listener to python bindings --- python/pyasic_rs/__init__.py | 2 + python/pyasic_rs/asic_rs.pyi | 6 + python/pyasic_rs/listener.py | 9 ++ src/listener.rs | 12 +- src/python/listener.rs | 266 +++++++++++++++++++++++++++++++++++ src/python/mod.rs | 3 + 6 files changed, 292 insertions(+), 6 deletions(-) create mode 100644 python/pyasic_rs/listener.py create mode 100644 src/python/listener.rs diff --git a/python/pyasic_rs/__init__.py b/python/pyasic_rs/__init__.py index b6a886ba..c2ecbdaf 100644 --- a/python/pyasic_rs/__init__.py +++ b/python/pyasic_rs/__init__.py @@ -15,6 +15,7 @@ TuningConfig, ) from .factory import MinerFactory +from .listener import MinerListener from .miner import Miner from .data import TuningTarget @@ -23,6 +24,7 @@ "FanMode", "Miner", "MinerFactory", + "MinerListener", "Pool", "PoolGroup", "ScalingConfig", diff --git a/python/pyasic_rs/asic_rs.pyi b/python/pyasic_rs/asic_rs.pyi index 27ab09e7..c2184974 100644 --- a/python/pyasic_rs/asic_rs.pyi +++ b/python/pyasic_rs/asic_rs.pyi @@ -636,6 +636,12 @@ class MinerHardware: @classmethod def model_validate(cls, /, obj: "object", **_kwargs: "object") -> "MinerHardware": ... +@final +class MinerListener: + def __new__(cls, /) -> MinerListener: ... + def listen(self, /) -> AsyncIterator[Miner]: ... + def listen_ip_only(self, /) -> AsyncIterator[IPv4Address |IPv6Address]: ... + @final class MinerMessage: @classmethod diff --git a/python/pyasic_rs/listener.py b/python/pyasic_rs/listener.py new file mode 100644 index 00000000..fa301d8f --- /dev/null +++ b/python/pyasic_rs/listener.py @@ -0,0 +1,9 @@ +"""Miner broadcast listener API. + +Import `MinerListener` from this module to listen for miner broadcast packets +and receive miners or IP addresses as async iterators. +""" + +from pyasic_rs.asic_rs import MinerListener + +__all__ = ["MinerListener"] diff --git a/src/listener.rs b/src/listener.rs index c13a3a7d..9f64c4a6 100644 --- a/src/listener.rs +++ b/src/listener.rs @@ -50,7 +50,7 @@ impl MinerListener { /// ``` pub async fn listen( &self, - ) -> Pin>>> + '_>> { + ) -> Pin>>> + Send + '_>> { let am_stream = self.antminer_listener.listen().await; let wm_stream = self.whatsminer_listener.listen().await; @@ -60,7 +60,7 @@ impl MinerListener { } pub async fn listen_ip_only( &self, - ) -> Pin>> + '_>> { + ) -> Pin>> + Send + '_>> { let am_stream = self.antminer_listener.listen_ip_only().await; let wm_stream = self.whatsminer_listener.listen_ip_only().await; @@ -79,7 +79,7 @@ impl AntMinerListener { pub(crate) async fn listen( &self, - ) -> impl Stream>>> { + ) -> impl Stream>>> + Send { stream! { let factory = MinerFactory::new(); let sock = match UdpSocket::bind("0.0.0.0:14235").await { @@ -105,7 +105,7 @@ impl AntMinerListener { } pub(crate) async fn listen_ip_only( &self, - ) -> impl Stream>> { + ) -> impl Stream>> + Send { stream! { let sock = match UdpSocket::bind("0.0.0.0:14235").await { Ok(s) => s, @@ -139,7 +139,7 @@ impl WhatsMinerListener { pub(crate) async fn listen( &self, - ) -> impl Stream>>> { + ) -> impl Stream>>> + Send { stream! { let factory = MinerFactory::new(); let sock = match UdpSocket::bind("0.0.0.0:8888").await { @@ -165,7 +165,7 @@ impl WhatsMinerListener { } pub(crate) async fn listen_ip_only( &self, - ) -> impl Stream>> { + ) -> impl Stream>> + Send { stream! { let sock = match UdpSocket::bind("0.0.0.0:8888").await { Ok(s) => s, diff --git a/src/python/listener.rs b/src/python/listener.rs new file mode 100644 index 00000000..55fbd012 --- /dev/null +++ b/src/python/listener.rs @@ -0,0 +1,266 @@ +use std::{ + net::IpAddr, + pin::Pin, + sync::{Arc, Mutex}, +}; + +use crate::{ + listener::MinerListener as MinerListenerBase, + python::{ + miner::Miner, + typing::{CancelAction, PyAsyncIterator, abortable_future_into_py_with_cancel}, + }, +}; +use asic_rs_core::traits::miner::Miner as MinerTrait; +use async_stream::stream; +use futures::Stream; +use pyo3::{ + exceptions::{PyRuntimeError, PyStopAsyncIteration}, + prelude::*, +}; +use tokio_stream::StreamExt; + +type ListenerMinerStream = + Pin>>> + Send>>; +type ListenerIpStream = Pin>> + Send>>; + +enum StreamState { + Ready(S), + InUse, + Closed, +} + +struct StreamLease { + state: Arc>>, + stream: Option, +} + +impl StreamLease { + fn take(state: Arc>>) -> PyResult { + let stream = { + let mut guard = state + .lock() + .map_err(|_| PyRuntimeError::new_err("stream state lock poisoned"))?; + match std::mem::replace(&mut *guard, StreamState::InUse) { + StreamState::Ready(stream) => stream, + StreamState::InUse => { + return Err(PyRuntimeError::new_err("stream is already being polled")); + } + StreamState::Closed => { + *guard = StreamState::Closed; + return Err(PyStopAsyncIteration::new_err("stream complete")); + } + } + }; + + Ok(Self { + state, + stream: Some(stream), + }) + } + + fn stream_mut(&mut self) -> PyResult<&mut S> { + self.stream + .as_mut() + .ok_or_else(|| PyRuntimeError::new_err("stream lease missing stream")) + } + + fn store(mut self, state: StreamState) -> PyResult<()> { + self.stream = None; + let mut guard = self + .state + .lock() + .map_err(|_| PyRuntimeError::new_err("stream state lock poisoned"))?; + if matches!(*guard, StreamState::Closed) && matches!(state, StreamState::Ready(_)) { + return Ok(()); + } + *guard = state; + Ok(()) + } + + fn store_ready(mut self) -> PyResult<()> { + let stream = self + .stream + .take() + .ok_or_else(|| PyRuntimeError::new_err("stream lease missing stream"))?; + self.store(StreamState::Ready(stream)) + } + + fn close(self) -> PyResult<()> { + self.store(StreamState::Closed) + } +} + +impl Drop for StreamLease { + fn drop(&mut self) { + if self.stream.is_some() + && let Ok(mut guard) = self.state.lock() + { + *guard = StreamState::Closed; + } + } +} + +fn close_stream_state(state: &Arc>>) { + let _previous = { + let Ok(mut guard) = state.lock() else { + return; + }; + std::mem::replace(&mut *guard, StreamState::Closed) + }; +} + +fn close_stream_on_cancel(state: Arc>>) -> CancelAction { + Box::new(move || close_stream_state(&state)) +} + +fn miner_stream(listener: Arc) -> ListenerMinerStream { + Box::pin(stream! { + let mut inner = listener.listen().await; + while let Some(item) = inner.next().await { + yield item; + } + }) +} + +fn ip_stream(listener: Arc) -> ListenerIpStream { + Box::pin(stream! { + let mut inner = listener.listen_ip_only().await; + while let Some(item) = inner.next().await { + yield item; + } + }) +} + +#[pyclass] +struct PyListenerMinerStream { + inner: Arc>>, +} + +impl PyListenerMinerStream { + fn new(inner: ListenerMinerStream) -> Self { + Self { + inner: Arc::new(Mutex::new(StreamState::Ready(inner))), + } + } +} + +#[pymethods] +impl PyListenerMinerStream { + pub fn __aiter__(slf: PyRef) -> PyRef { + slf + } + + pub fn __anext__<'py>(&self, py: Python<'py>) -> PyResult> { + let inner = self.inner.clone(); + let cancel_inner = inner.clone(); + abortable_future_into_py_with_cancel( + py, + async move { + let mut lease = StreamLease::take(inner)?; + loop { + match lease.stream_mut()?.next().await { + Some(Ok(Some(miner))) => { + let miner = Miner::from(miner); + lease.store_ready()?; + return Ok(miner); + } + Some(Ok(None)) => continue, + Some(Err(error)) => { + lease.close()?; + return Err(PyRuntimeError::new_err(error.to_string())); + } + None => { + lease.close()?; + return Err(PyStopAsyncIteration::new_err("stream complete")); + } + } + } + }, + Some(close_stream_on_cancel(cancel_inner)), + ) + } +} + +#[pyclass] +struct PyListenerIpStream { + inner: Arc>>, +} + +impl PyListenerIpStream { + fn new(inner: ListenerIpStream) -> Self { + Self { + inner: Arc::new(Mutex::new(StreamState::Ready(inner))), + } + } +} + +#[pymethods] +impl PyListenerIpStream { + pub fn __aiter__(slf: PyRef) -> PyRef { + slf + } + + pub fn __anext__<'py>(&self, py: Python<'py>) -> PyResult> { + let inner = self.inner.clone(); + let cancel_inner = inner.clone(); + abortable_future_into_py_with_cancel( + py, + async move { + let mut lease = StreamLease::take(inner)?; + loop { + match lease.stream_mut()?.next().await { + Some(Ok(Some(ip))) => { + lease.store_ready()?; + return Ok(ip); + } + Some(Ok(None)) => continue, + Some(Err(error)) => { + lease.close()?; + return Err(PyRuntimeError::new_err(error.to_string())); + } + None => { + lease.close()?; + return Err(PyStopAsyncIteration::new_err("stream complete")); + } + } + } + }, + Some(close_stream_on_cancel(cancel_inner)), + ) + } +} + +/// Python listener for miner broadcast packets. +#[pyclass(module = "asic_rs")] +pub(crate) struct MinerListener { + inner: Arc, +} + +#[pymethods] +impl MinerListener { + /// Create a listener for supported miner broadcast packets. + #[new] + pub fn new() -> Self { + Self { + inner: Arc::new(MinerListenerBase::new()), + } + } + + /// Return an async iterator over miners as broadcast packets arrive. + pub fn listen<'py>(&self, py: Python<'py>) -> PyResult> { + Bound::new( + py, + PyListenerMinerStream::new(miner_stream(self.inner.clone())), + ) + .map(Bound::into_any) + .map(PyAsyncIterator::new) + } + + /// Return an async iterator over miner IP addresses as broadcast packets arrive. + pub fn listen_ip_only<'py>(&self, py: Python<'py>) -> PyResult> { + Bound::new(py, PyListenerIpStream::new(ip_stream(self.inner.clone()))) + .map(Bound::into_any) + .map(PyAsyncIterator::new) + } +} diff --git a/src/python/mod.rs b/src/python/mod.rs index 6fc876c5..048117d3 100644 --- a/src/python/mod.rs +++ b/src/python/mod.rs @@ -1,6 +1,7 @@ use pyo3::prelude::*; mod factory; +mod listener; mod miner; mod typing; @@ -26,6 +27,8 @@ mod asic_rs { #[pymodule_export] use super::factory::MinerFactory; #[pymodule_export] + use super::listener::MinerListener; + #[pymodule_export] use super::miner::Miner; #[pymodule_export] use asic_rs_core::config::{