Skip to content
Merged
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
22 changes: 9 additions & 13 deletions src/io/formats/bag/parallel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,30 +65,26 @@ impl BagFormat {

/// Open a BAG reader from a transport source.
#[cfg(feature = "remote")]
pub fn open_from_transport(
mut transport: Box<dyn crate::io::transport::Transport>,
pub async fn open_from_transport(
transport: Box<dyn crate::io::transport::Transport>,
path: String,
) -> Result<ParallelBagReader> {
use std::pin::Pin;
use std::task::{Context, Poll, Waker};
use std::future::poll_fn;

let mut data = Vec::new();
let mut buffer = vec![0u8; 64 * 1024];
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
let mut pinned_transport = unsafe { Pin::new_unchecked(transport.as_mut()) };
let mut pinned_transport = Box::into_pin(transport);

loop {
match pinned_transport.as_mut().poll_read(&mut cx, &mut buffer) {
Poll::Ready(Ok(0)) => break,
Poll::Ready(Ok(n)) => data.extend_from_slice(&buffer[..n]),
Poll::Ready(Err(e)) => {
match poll_fn(|cx| pinned_transport.as_mut().poll_read(cx, &mut buffer)).await {
Ok(0) => break,
Ok(n) => data.extend_from_slice(&buffer[..n]),
Err(e) => {
return Err(CodecError::encode(
"Transport",
format!("Failed to read from {path}: {e}"),
));
}
Poll::Pending => std::thread::yield_now(),
}
}

Expand Down Expand Up @@ -346,7 +342,7 @@ impl ParallelBagReader {

impl FormatReader for ParallelBagReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
2 changes: 1 addition & 1 deletion src/io/formats/bag/sequential.rs
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,7 @@ impl SequentialBagReader {

impl FormatReader for SequentialBagReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
2 changes: 1 addition & 1 deletion src/io/formats/mcap/parallel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -685,7 +685,7 @@ impl ParallelMcapReader {

impl FormatReader for ParallelMcapReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
22 changes: 9 additions & 13 deletions src/io/formats/mcap/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,30 +59,26 @@ impl McapFormat {

/// Open an MCAP reader from a transport source.
#[cfg(feature = "remote")]
pub fn open_from_transport(
mut transport: Box<dyn crate::io::transport::Transport>,
pub async fn open_from_transport(
transport: Box<dyn crate::io::transport::Transport>,
path: String,
) -> Result<McapReader> {
use std::pin::Pin;
use std::task::{Context, Poll, Waker};
use std::future::poll_fn;

let mut data = Vec::new();
let mut buffer = vec![0u8; 64 * 1024];
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
let mut pinned_transport = unsafe { Pin::new_unchecked(transport.as_mut()) };
let mut pinned_transport = Box::into_pin(transport);

loop {
match pinned_transport.as_mut().poll_read(&mut cx, &mut buffer) {
Poll::Ready(Ok(0)) => break,
Poll::Ready(Ok(n)) => data.extend_from_slice(&buffer[..n]),
Poll::Ready(Err(e)) => {
match poll_fn(|cx| pinned_transport.as_mut().poll_read(cx, &mut buffer)).await {
Ok(0) => break,
Ok(n) => data.extend_from_slice(&buffer[..n]),
Err(e) => {
return Err(CodecError::encode(
"Transport",
format!("Failed to read from {path}: {e}"),
));
}
Poll::Pending => std::thread::yield_now(),
}
}

Expand Down Expand Up @@ -271,7 +267,7 @@ impl McapReader {

impl FormatReader for McapReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
2 changes: 1 addition & 1 deletion src/io/formats/mcap/sequential.rs
Original file line number Diff line number Diff line change
Expand Up @@ -233,7 +233,7 @@ impl SequentialMcapReader {

impl FormatReader for SequentialMcapReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
2 changes: 1 addition & 1 deletion src/io/formats/mcap/two_pass.rs
Original file line number Diff line number Diff line change
Expand Up @@ -583,7 +583,7 @@ impl TwoPassMcapReader {

impl FormatReader for TwoPassMcapReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
2 changes: 1 addition & 1 deletion src/io/formats/rrd/parallel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -470,7 +470,7 @@ impl Iterator for RrdDecodedMessageWithTimestampStream<'_> {

impl FormatReader for ParallelRrdReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
22 changes: 9 additions & 13 deletions src/io/formats/rrd/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,30 +56,26 @@ impl RrdFormat {

/// Open an RRD reader from a transport source.
#[cfg(feature = "remote")]
pub fn open_from_transport(
mut transport: Box<dyn crate::io::transport::Transport>,
pub async fn open_from_transport(
transport: Box<dyn crate::io::transport::Transport>,
path: String,
) -> Result<ParallelRrdReader> {
use std::pin::Pin;
use std::task::{Context, Poll, Waker};
use std::future::poll_fn;

let mut data = Vec::new();
let mut buffer = vec![0u8; 64 * 1024];
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
let mut pinned_transport = unsafe { Pin::new_unchecked(transport.as_mut()) };
let mut pinned_transport = Box::into_pin(transport);

loop {
match pinned_transport.as_mut().poll_read(&mut cx, &mut buffer) {
Poll::Ready(Ok(0)) => break,
Poll::Ready(Ok(n)) => data.extend_from_slice(&buffer[..n]),
Poll::Ready(Err(e)) => {
match poll_fn(|cx| pinned_transport.as_mut().poll_read(cx, &mut buffer)).await {
Ok(0) => break,
Ok(n) => data.extend_from_slice(&buffer[..n]),
Err(e) => {
return Err(CodecError::encode(
"Transport",
format!("Failed to read from {path}: {e}"),
));
}
Poll::Pending => std::thread::yield_now(),
}
}

Expand Down Expand Up @@ -394,7 +390,7 @@ impl RrdReader {

impl FormatReader for RrdReader {
#[cfg(feature = "remote")]
fn open_from_transport(
async fn open_from_transport(
_transport: Box<dyn crate::io::transport::Transport>,
_path: String,
) -> Result<Self>
Expand Down
Loading