Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

fix(transport): discard async data if receiver is closed #419

Merged
merged 2 commits into from
Oct 23, 2024
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
67 changes: 27 additions & 40 deletions crates/transport/src/value.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,13 @@ use bytes::{Buf as _, BufMut as _, Bytes, BytesMut};
use futures::stream::{self, FuturesUnordered};
use futures::{Stream, StreamExt as _, TryStreamExt as _};
use tokio::io::{AsyncRead, AsyncReadExt as _, AsyncWrite, AsyncWriteExt as _};
use tokio::select;
use tokio::sync::{mpsc, oneshot};
use tokio::task::JoinSet;
use tokio::{select, try_join};
use tokio_stream::wrappers::ReceiverStream;
use tokio_util::codec::{Encoder as _, FramedRead};
use tokio_util::io::StreamReader;
use tracing::{error, instrument, trace, Instrument as _, Span};
use tracing::{debug, error, instrument, trace, Instrument as _, Span};
use wasm_tokio::cm::{
BoolCodec, F32Codec, F64Codec, OptionDecoder, OptionEncoder, PrimValEncoder, ResultDecoder,
ResultEncoder, S16Codec, S32Codec, S64Codec, S8Codec, TupleDecoder, TupleEncoder, U16Codec,
Expand Down Expand Up @@ -1617,30 +1617,20 @@ where
return Err(std::io::ErrorKind::UnexpectedEof.into());
};
let item = item?;
try_join!(
async {
tx.send(item).map_err(|_| {
std::io::Error::new(
std::io::ErrorKind::BrokenPipe,
"future receiver closed",
)
})
},
async {
if let Some(rx) = dec.decoder_mut().take_deferred() {
let buf = mem::take(dec.read_buffer_mut());
let mut r = dec.into_inner();
if !r.buffer.is_empty() {
r.buffer.unsplit(buf);
} else {
r.buffer = buf;
}
rx(r, Vec::default()).await
} else {
Ok(())
}
if tx.send(item).is_err() {
debug!("future receiver closed, discard data");
return Ok(());
}
if let Some(rx) = dec.decoder_mut().take_deferred() {
let buf = mem::take(dec.read_buffer_mut());
let mut r = dec.into_inner();
if !r.buffer.is_empty() {
r.buffer.unsplit(buf);
} else {
r.buffer = buf;
}
)?;
rx(r, Vec::default()).await?;
}
Ok(())
}
.instrument(span),
Expand Down Expand Up @@ -2042,9 +2032,10 @@ where
)
})?;
trace!(i, end, "received stream chunk");
tx.send(chunk).await.map_err(|_| {
std::io::Error::new(std::io::ErrorKind::BrokenPipe, "stream receiver closed")
})?;
if tx.send(chunk).await.is_err() {
debug!("stream receiver closed, discard data");
return Ok(())
}
for (i, deferred) in zip(i.., mem::take(&mut framed.decoder_mut().deferred)) {
if let Some(deferred) = deferred {
let r = framed.get_ref().index(&[i]).map_err(std::io::Error::other)?;
Expand Down Expand Up @@ -2162,12 +2153,10 @@ where
return Ok(());
}
trace!(?chunk, "received pending byte stream chunk");
tx.send(chunk).await.map_err(|_| {
std::io::Error::new(
std::io::ErrorKind::BrokenPipe,
"stream receiver closed",
)
})?;
if tx.send(chunk).await.is_err() {
debug!("stream receiver closed, discard data");
return Ok(());
}
}
Ok(())
}
Expand Down Expand Up @@ -2240,12 +2229,10 @@ where
return Ok(());
}
trace!(?chunk, "received byte stream chunk");
tx.send(std::io::Result::Ok(chunk)).await.map_err(|_| {
std::io::Error::new(
std::io::ErrorKind::BrokenPipe,
"stream receiver closed",
)
})?;
if tx.send(std::io::Result::Ok(chunk)).await.is_err() {
debug!("stream receiver closed, discard data");
return Ok(());
}
}
Ok(())
})
Expand Down
2 changes: 0 additions & 2 deletions flake.nix
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,13 @@
"https://nixify.cachix.org"
"https://crane.cachix.org"
"https://nix-community.cachix.org"
"https://cache.garnix.io"
];
nixConfig.extra-trusted-public-keys = [
"bytecodealliance.cachix.org-1:0SBgh//n2n0heh0sDFhTm+ZKBRy2sInakzFGfzN531Y="
"wasmcloud.cachix.org-1:9gRBzsKh+x2HbVVspreFg/6iFRiD4aOcUQfXVDl3hiM="
"nixify.cachix.org-1:95SiUQuf8Ij0hwDweALJsLtnMyv/otZamWNRp1Q1pXw="
"crane.cachix.org-1:8Scfpmn9w+hGdXH/Q9tTLiYAE/2dnJYRJP7kl80GuRk="
"nix-community.cachix.org-1:mB9FSh9qf2dCimDSUo8Zy7bkq5CX+/rkCWyvRCYg3Fs="
"cache.garnix.io:CTFPyKSLcx5RMJKfLo5EEPUObbA78b0YQ2DTCJXqr9g="
];

inputs.nixify.inputs.nixlib.follows = "nixlib";
Expand Down
Loading