Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions python/pyasic_rs/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
TuningConfig,
)
from .factory import MinerFactory
from .listener import MinerListener
from .miner import Miner
from .data import TuningTarget

Expand All @@ -23,6 +24,7 @@
"FanMode",
"Miner",
"MinerFactory",
"MinerListener",
"Pool",
"PoolGroup",
"ScalingConfig",
Expand Down
6 changes: 6 additions & 0 deletions python/pyasic_rs/asic_rs.pyi

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

9 changes: 9 additions & 0 deletions python/pyasic_rs/listener.py
Original file line number Diff line number Diff line change
@@ -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"]
12 changes: 6 additions & 6 deletions src/listener.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ impl MinerListener {
/// ```
pub async fn listen(
&self,
) -> Pin<Box<dyn Stream<Item = anyhow::Result<Option<Box<dyn Miner>>>> + '_>> {
) -> Pin<Box<dyn Stream<Item = anyhow::Result<Option<Box<dyn Miner>>>> + Send + '_>> {
let am_stream = self.antminer_listener.listen().await;
let wm_stream = self.whatsminer_listener.listen().await;

Expand All @@ -60,7 +60,7 @@ impl MinerListener {
}
pub async fn listen_ip_only(
&self,
) -> Pin<Box<dyn Stream<Item = anyhow::Result<Option<IpAddr>>> + '_>> {
) -> Pin<Box<dyn Stream<Item = anyhow::Result<Option<IpAddr>>> + Send + '_>> {
let am_stream = self.antminer_listener.listen_ip_only().await;
let wm_stream = self.whatsminer_listener.listen_ip_only().await;

Expand All @@ -79,7 +79,7 @@ impl AntMinerListener {

pub(crate) async fn listen(
&self,
) -> impl Stream<Item = anyhow::Result<Option<Box<dyn Miner>>>> {
) -> impl Stream<Item = anyhow::Result<Option<Box<dyn Miner>>>> + Send {
stream! {
let factory = MinerFactory::new();
let sock = match UdpSocket::bind("0.0.0.0:14235").await {
Expand All @@ -105,7 +105,7 @@ impl AntMinerListener {
}
pub(crate) async fn listen_ip_only(
&self,
) -> impl Stream<Item = anyhow::Result<Option<IpAddr>>> {
) -> impl Stream<Item = anyhow::Result<Option<IpAddr>>> + Send {
stream! {
let sock = match UdpSocket::bind("0.0.0.0:14235").await {
Ok(s) => s,
Expand Down Expand Up @@ -139,7 +139,7 @@ impl WhatsMinerListener {

pub(crate) async fn listen(
&self,
) -> impl Stream<Item = anyhow::Result<Option<Box<dyn Miner>>>> {
) -> impl Stream<Item = anyhow::Result<Option<Box<dyn Miner>>>> + Send {
stream! {
let factory = MinerFactory::new();
let sock = match UdpSocket::bind("0.0.0.0:8888").await {
Expand All @@ -165,7 +165,7 @@ impl WhatsMinerListener {
}
pub(crate) async fn listen_ip_only(
&self,
) -> impl Stream<Item = anyhow::Result<Option<IpAddr>>> {
) -> impl Stream<Item = anyhow::Result<Option<IpAddr>>> + Send {
stream! {
let sock = match UdpSocket::bind("0.0.0.0:8888").await {
Ok(s) => s,
Expand Down
266 changes: 266 additions & 0 deletions src/python/listener.rs
Original file line number Diff line number Diff line change
@@ -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<Box<dyn Stream<Item = anyhow::Result<Option<Box<dyn MinerTrait>>>> + Send>>;
type ListenerIpStream = Pin<Box<dyn Stream<Item = anyhow::Result<Option<IpAddr>>> + Send>>;

enum StreamState<S> {
Ready(S),
InUse,
Closed,
}

struct StreamLease<S> {
state: Arc<Mutex<StreamState<S>>>,
stream: Option<S>,
}

impl<S> StreamLease<S> {
fn take(state: Arc<Mutex<StreamState<S>>>) -> PyResult<Self> {
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<S>) -> 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<S> Drop for StreamLease<S> {
fn drop(&mut self) {
if self.stream.is_some()
&& let Ok(mut guard) = self.state.lock()
{
*guard = StreamState::Closed;
}
}
}

fn close_stream_state<S>(state: &Arc<Mutex<StreamState<S>>>) {
let _previous = {
let Ok(mut guard) = state.lock() else {
return;
};
std::mem::replace(&mut *guard, StreamState::Closed)
};
}

fn close_stream_on_cancel<S: Send + 'static>(state: Arc<Mutex<StreamState<S>>>) -> CancelAction {
Box::new(move || close_stream_state(&state))
}

fn miner_stream(listener: Arc<MinerListenerBase>) -> ListenerMinerStream {
Box::pin(stream! {
let mut inner = listener.listen().await;
while let Some(item) = inner.next().await {
yield item;
}
})
}

fn ip_stream(listener: Arc<MinerListenerBase>) -> 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<Mutex<StreamState<ListenerMinerStream>>>,
}

impl PyListenerMinerStream {
fn new(inner: ListenerMinerStream) -> Self {
Self {
inner: Arc::new(Mutex::new(StreamState::Ready(inner))),
}
}
}

#[pymethods]
impl PyListenerMinerStream {
pub fn __aiter__(slf: PyRef<Self>) -> PyRef<Self> {
slf
}

pub fn __anext__<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
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<Mutex<StreamState<ListenerIpStream>>>,
}

impl PyListenerIpStream {
fn new(inner: ListenerIpStream) -> Self {
Self {
inner: Arc::new(Mutex::new(StreamState::Ready(inner))),
}
}
}

#[pymethods]
impl PyListenerIpStream {
pub fn __aiter__(slf: PyRef<Self>) -> PyRef<Self> {
slf
}

pub fn __anext__<'py>(&self, py: Python<'py>) -> PyResult<Bound<'py, PyAny>> {
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<MinerListenerBase>,
}

#[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<PyAsyncIterator<Miner>> {
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<PyAsyncIterator<IpAddr>> {
Bound::new(py, PyListenerIpStream::new(ip_stream(self.inner.clone())))
.map(Bound::into_any)
.map(PyAsyncIterator::new)
}
}
3 changes: 3 additions & 0 deletions src/python/mod.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
use pyo3::prelude::*;

mod factory;
mod listener;
mod miner;
mod typing;

Expand All @@ -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::{
Expand Down