Browse Source

Implement `Stream` trait for `Ros2Subscription`

tags/v0.2.5-alpha.2
Philipp Oppermann 2 years ago
parent
commit
42dcd67b24
Failed to extract signature
2 changed files with 26 additions and 1 deletions
  1. +1
    -0
      python/Cargo.toml
  2. +25
    -1
      python/src/lib.rs

+ 1
- 0
python/Cargo.toml View File

@@ -14,6 +14,7 @@ serde_yaml = "0.8.23"
flume = "0.10.14"
arrow = { version = "45.0.0", features = ["pyarrow"] }
pythonize = "0.19.0"
futures = "0.3.28"

[lib]
name = "dora_ros2_bridge"


+ 25
- 1
python/src/lib.rs View File

@@ -8,10 +8,11 @@ use std::{
use ::dora_ros2_bridge::{ros2_client, rustdds};
use dora_ros2_bridge_msg_gen::types::Message;
use eyre::{eyre, Context, ContextCompat};
use futures::{Stream, StreamExt};
use pyo3::{
prelude::{pyclass, pymethods, pymodule},
types::PyModule,
PyAny, PyObject, PyResult, Python,
wrap_pyfunction, PyAny, PyErr, PyObject, PyResult, Python, ToPyObject,
};
use typed::{
deserialize::{Ros2Value, TypedDeserializer},
@@ -226,6 +227,29 @@ impl Ros2Subscription {
}
}

impl Ros2Subscription {
fn as_stream(
&self,
) -> impl Stream<Item = Result<(Ros2Value, ros2_client::MessageInfo), rustdds::dds::ReadError>> + '_
{
self.subscription
.async_stream_seed(self.deserializer.clone())
}
}

impl Stream for Ros2Subscription {
type Item = Result<(Ros2Value, ros2_client::MessageInfo), rustdds::dds::ReadError>;

fn poll_next(
self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
let s = self.as_stream();
futures::pin_mut!(s);
s.poll_next_unpin(cx)
}
}

#[pymodule]
fn dora_ros2_bridge(_py: Python, m: &PyModule) -> PyResult<()> {
m.add_class::<Ros2Context>()?;


Loading…
Cancel
Save