Skip to content
Merged
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
95 changes: 37 additions & 58 deletions src/substream/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -369,7 +369,34 @@ impl Substream {

io.write_all(&payload)
.await
.map_err(|_| Error::SubstreamError(SubstreamError::ConnectionClosed))
.map_err(|_| Error::SubstreamError(SubstreamError::ConnectionClosed))?;

// Flush the stream.
io.flush().await.map_err(From::from)
}

/// Send unsigned varint payload to remote peer.
async fn send_unsigned_varint_payload<T: AsyncWrite + Unpin>(
io: &mut T,
bytes: Bytes,
max_size: Option<usize>,
) -> crate::Result<()> {
if let Some(max_size) = max_size {
if bytes.len() > max_size {
return Err(Error::IoError(ErrorKind::PermissionDenied));
}
}

// Write the length of the frame.
let mut buffer = unsigned_varint::encode::usize_buffer();
let encoded_len = unsigned_varint::encode::usize(bytes.len(), &mut buffer).len();
io.write_all(&buffer[..encoded_len]).await?;

// Write the frame.
io.write_all(bytes.as_ref()).await?;

// Flush the stream.
io.flush().await.map_err(From::from)
}

/// Send framed data to remote peer.
Expand All @@ -386,7 +413,7 @@ impl Substream {
/// # Panics
///
/// Panics if no codec is provided.
pub async fn send_framed(&mut self, mut bytes: Bytes) -> crate::Result<()> {
pub async fn send_framed(&mut self, bytes: Bytes) -> crate::Result<()> {
tracing::trace!(
target: LOG_TARGET,
peer = ?self.peer,
Expand All @@ -403,48 +430,16 @@ impl Substream {
ProtocolCodec::Unspecified => panic!("codec is unspecified"),
ProtocolCodec::Identity(payload_size) =>
Self::send_identity_payload(substream, payload_size, bytes).await,
ProtocolCodec::UnsignedVarint(max_size) => {
check_size!(max_size, bytes.len());

let mut buffer = [0u8; 10];
let len = unsigned_varint::encode::usize(bytes.len(), &mut buffer);
let mut offset = 0;

while offset < len.len() {
offset += substream.write(&len[offset..]).await?;
}

while bytes.has_remaining() {
let nwritten = substream.write(&bytes).await?;
bytes.advance(nwritten);
}

substream.flush().await.map_err(From::from)
}
ProtocolCodec::UnsignedVarint(max_size) =>
Self::send_unsigned_varint_payload(substream, bytes, max_size).await,
},
#[cfg(feature = "websocket")]
SubstreamType::WebSocket(ref mut substream) => match self.codec {
ProtocolCodec::Unspecified => panic!("codec is unspecified"),
ProtocolCodec::Identity(payload_size) =>
Self::send_identity_payload(substream, payload_size, bytes).await,
ProtocolCodec::UnsignedVarint(max_size) => {
check_size!(max_size, bytes.len());

let mut buffer = [0u8; 10];
let len = unsigned_varint::encode::usize(bytes.len(), &mut buffer);
let mut offset = 0;

while offset < len.len() {
offset += substream.write(&len[offset..]).await?;
}

while bytes.has_remaining() {
let nwritten = substream.write(&bytes).await?;
bytes.advance(nwritten);
}

substream.flush().await.map_err(From::from)
}
ProtocolCodec::UnsignedVarint(max_size) =>
Self::send_unsigned_varint_payload(substream, bytes, max_size).await,
},
#[cfg(feature = "quic")]
SubstreamType::Quic(ref mut substream) => match self.codec {
Expand All @@ -454,7 +449,7 @@ impl Substream {
ProtocolCodec::UnsignedVarint(max_size) => {
check_size!(max_size, bytes.len());

let mut buffer = [0u8; 10];
let mut buffer = unsigned_varint::encode::usize_buffer();
let len = unsigned_varint::encode::usize(bytes.len(), &mut buffer);
let len = BytesMut::from(len);

Expand All @@ -466,24 +461,8 @@ impl Substream {
ProtocolCodec::Unspecified => panic!("codec is unspecified"),
ProtocolCodec::Identity(payload_size) =>
Self::send_identity_payload(substream, payload_size, bytes).await,
ProtocolCodec::UnsignedVarint(max_size) => {
check_size!(max_size, bytes.len());

let mut buffer = [0u8; 10];
let len = unsigned_varint::encode::usize(bytes.len(), &mut buffer);
let mut offset = 0;

while offset < len.len() {
offset += substream.write(&len[offset..]).await?;
}

while bytes.has_remaining() {
let nwritten = substream.write(&bytes).await?;
bytes.advance(nwritten);
}

substream.flush().await.map_err(From::from)
}
ProtocolCodec::UnsignedVarint(max_size) =>
Self::send_unsigned_varint_payload(substream, bytes, max_size).await,
},
}
}
Expand Down Expand Up @@ -722,7 +701,7 @@ impl Sink<Bytes> for Substream {
check_size!(max_size, item.len());

let len = {
let mut buffer = [0u8; 10];
let mut buffer = unsigned_varint::encode::usize_buffer();
let len = unsigned_varint::encode::usize(item.len(), &mut buffer);
BytesMut::from(len)
};
Expand Down