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
7 changes: 6 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,8 @@ categories = ["asynchronous", "network-programming"]
[features]
json = [ "serde", "serde_json" ]
cbor = [ "serde", "serde_cbor" ]
default = [ "json", "cbor" ]
serde_bincode = [ "serde", "bincode" ]
default = [ "json", "cbor", "serde_bincode" ]

[dependencies]

Expand All @@ -36,3 +37,7 @@ optional = true
[dependencies.serde_cbor]
version = '0.10.2'
optional = true

[dependencies.bincode]
version = "1.2.1"
optional = true
27 changes: 26 additions & 1 deletion src/codec/length.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,11 +67,36 @@ impl Decoder for LengthCodec {
return Ok(None);
}

let len = src.get_u64() as usize;
let mut len_bytes = [0u8; U64_LENGTH];
len_bytes.copy_from_slice(&src[..U64_LENGTH]);
let len = u64::from_be_bytes(len_bytes) as usize;

if src.len() - U64_LENGTH >= len {
// Skip the length header we already read.
src.advance(U64_LENGTH);
Ok(Some(src.split_to(len).freeze()))
} else {
Ok(None)
}
}
}

#[cfg(test)]
mod tests {
use super::*;

mod decode {
use super::*;

#[test]
fn it_returns_bytes_withouth_length_header() {
let mut codec = LengthCodec{ };

let mut src = BytesMut::with_capacity(5);
src.put(&[0, 0, 0, 0, 0, 0, 0, 3u8, 1, 2, 3, 4][..]);
let item = codec.decode(&mut src).unwrap();

assert!(item == Some(Bytes::from(&[1u8, 2, 3][..])));
}
}
}
3 changes: 3 additions & 0 deletions src/codec/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,3 +12,6 @@ pub use self::lines::LinesCodec;

#[cfg(feature = "cbor")] mod cbor;
#[cfg(feature = "cbor")] pub use self::cbor::{CborCodec, CborCodecError};

#[cfg(feature = "bincode")] mod serde;
#[cfg(feature = "bincode")] pub use self::serde::SerdeCodec;
60 changes: 60 additions & 0 deletions src/codec/serde.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
use std::io;
use std::marker::PhantomData;
use bytes::{BytesMut, Bytes};
use serde::Serialize;
use serde::de::DeserializeOwned;
use bincode;

use crate::encoder::Encoder;
use crate::decoder::Decoder;
use crate::codec::LengthCodec;

/// Encodes/decodes types implementing Serde Serialize/Deserialize traits.
/// It is built on top of `LengthCodec`.
pub struct SerdeCodec<T: Serialize + DeserializeOwned> {
inner: LengthCodec,
phantom: PhantomData<T>,
}

impl<T: Serialize + DeserializeOwned> Default for SerdeCodec<T> {
fn default() -> Self {
Self {
inner: LengthCodec {},
phantom: PhantomData,
}
}
}

impl<T: Serialize + DeserializeOwned> Encoder for SerdeCodec<T> {
type Item = T;
type Error = io::Error;

fn encode(&mut self, src: Self::Item, dst: &mut BytesMut) -> Result<(), Self::Error> {
let data = bincode::serialize(&src).map_err(to_io_err)?;
let bytes = Bytes::from(data);
self.inner.encode(bytes, dst)
}
}

impl<T: Serialize + DeserializeOwned> Decoder for SerdeCodec<T> {
type Item = T;
type Error = io::Error;

fn decode(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
match self.inner.decode(src)? {
Some(bytes) => {
bincode::deserialize(&bytes)
.map_err(to_io_err)
.map(|item| Some(item))
}
None => Ok(None),
}
}
}

fn to_io_err(err: Box<bincode::ErrorKind>) -> io::Error {
match *err {
bincode::ErrorKind::Io(e) => e,
other => io::Error::new(io::ErrorKind::Other, other),
}
}
2 changes: 1 addition & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
//! ```

mod codec;
pub use codec::{BytesCodec, LengthCodec, LinesCodec};
pub use codec::{BytesCodec, LengthCodec, LinesCodec, SerdeCodec};

#[cfg(feature = "json")] pub use codec::{JsonCodec, JsonCodecError};
#[cfg(feature = "cbor")] pub use codec::{CborCodec, CborCodecError};
Expand Down
29 changes: 29 additions & 0 deletions tests/length_delimited.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
use bytes::Bytes;
use futures::io::Cursor;
use futures::{executor, SinkExt, StreamExt};
use futures_codec::{Framed, LengthCodec};

#[test]
fn same_msgs_are_received_as_were_sent() {
let cur = Cursor::new(vec![0; 256]);
let mut framed = Framed::new(cur, LengthCodec {});

let send_msgs = async {
framed.send(Bytes::from("msg1")).await.unwrap();
framed.send(Bytes::from("msg2")).await.unwrap();
framed.send(Bytes::from("msg3")).await.unwrap();
};
executor::block_on(send_msgs);

let (mut cur, _) = framed.release();
cur.set_position(0);
let framed = Framed::new(cur, LengthCodec {});

let recv_msgs = framed.take(3)
.map(|res| res.unwrap())
.map(|buf| String::from_utf8(buf.to_vec()).unwrap())
.collect::<Vec<_>>();
let msgs: Vec<String> = executor::block_on(recv_msgs);

assert!(msgs == vec!["msg1", "msg2", "msg3"]);
}
47 changes: 47 additions & 0 deletions tests/serde.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
use futures::io::Cursor;
use futures::{executor, SinkExt, StreamExt};
use futures_codec::{Framed, SerdeCodec};
use serde::{Serialize, Deserialize};

#[derive(Serialize, Deserialize, Debug, PartialEq)]
struct Person {
name: String,
age: u8,
}

impl Person {
fn new(name: &str, age: u8) -> Self {
Self {
name: name.into(),
age,
}
}
}

#[test]
fn serializes_serde_enabled_structures() {
let cur = Cursor::new(vec![0; 4096]);
let mut framed = Framed::new(cur, SerdeCodec::default());

let send_msgs = async {
framed.send(Person::new("John", 11)).await.unwrap();
framed.send(Person::new("Paul", 12)).await.unwrap();
framed.send(Person::new("Mike", 13)).await.unwrap();
};
executor::block_on(send_msgs);

let (mut cur, _) = framed.release();
cur.set_position(0);
let framed = Framed::new(cur, SerdeCodec::default());

let recv_msgs = framed.take(3)
.map(|res| res.unwrap())
.collect::<Vec<_>>();
let items: Vec<Person> = executor::block_on(recv_msgs);

assert!(items == vec![
Person::new("John", 11),
Person::new("Paul", 12),
Person::new("Mike", 13),
])
}