1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
use std::io::Cursor;
use async_std::prelude::*;
use byteorder::{BigEndian, ReadBytesExt, WriteBytesExt};
use log::debug;
use serde_cbor::de::from_slice;
use serde_cbor::ser::to_vec;
use crate::error::Error;
use crate::network::message::*;
pub use super::platform::socket::Stream;
pub use super::platform::socket::*;
pub async fn send_message(message: Message, stream: &mut GenericStream) -> Result<(), Error> {
debug!("Sending message: {:?}", message);
let payload = to_vec(&message).map_err(|err| Error::MessageDeserialization(err.to_string()))?;
send_bytes(&payload, stream).await
}
pub async fn send_bytes(payload: &[u8], stream: &mut GenericStream) -> Result<(), Error> {
let message_size = payload.len() as u64;
let mut header = vec![];
header.write_u64::<BigEndian>(message_size).unwrap();
stream.write_all(&header).await?;
for chunk in payload.chunks(1400) {
stream.write_all(chunk).await?;
}
Ok(())
}
pub async fn receive_bytes(stream: &mut GenericStream) -> Result<Vec<u8>, Error> {
let mut header = vec![0; 8];
stream.read_exact(&mut header).await?;
let mut header = Cursor::new(header);
let message_size = header.read_u64::<BigEndian>()? as usize;
let mut payload_bytes = Vec::with_capacity(message_size);
let mut chunk_buffer: [u8; 1400] = [0; 1400];
while payload_bytes.len() < message_size {
let received_bytes = stream.read(&mut chunk_buffer).await?;
if received_bytes == 0 {
return Err(Error::Connection(
"Connection went away while receiving payload.".into(),
));
}
payload_bytes.extend_from_slice(&chunk_buffer[0..received_bytes]);
}
Ok(payload_bytes)
}
pub async fn receive_message(stream: &mut GenericStream) -> Result<Message, Error> {
let payload_bytes = receive_bytes(stream).await?;
debug!("Received {} bytes", payload_bytes.len());
if payload_bytes.is_empty() {
return Err(Error::EmptyPayload);
}
let message: Message =
from_slice(&payload_bytes).map_err(|err| Error::MessageDeserialization(err.to_string()))?;
debug!("Received message: {:?}", message);
Ok(message)
}
#[cfg(test)]
mod test {
use super::*;
use async_std::net::{TcpListener, TcpStream};
use async_std::task;
use async_trait::async_trait;
use pretty_assertions::assert_eq;
use crate::network::platform::socket::Stream as PueueStream;
#[async_trait]
impl Listener for TcpListener {
async fn accept<'a>(&'a self) -> Result<GenericStream, Error> {
let (stream, _) = self.accept().await?;
Ok(Box::new(stream))
}
}
impl PueueStream for TcpStream {}
#[async_std::test]
async fn test_single_huge_payload() -> Result<(), Error> {
let listener = TcpListener::bind("127.0.0.1:0").await?;
let addr = listener.local_addr()?;
let payload = "a".repeat(100_000);
let message = create_success_message(payload);
let original_bytes = to_vec(&message).expect("Failed to serialize message.");
let listener: GenericListener = Box::new(listener);
task::spawn(async move {
let mut stream = listener.accept().await.unwrap();
let message_bytes = receive_bytes(&mut stream).await.unwrap();
let message: Message = from_slice(&message_bytes).unwrap();
send_message(message, &mut stream).await.unwrap();
});
let mut client: GenericStream = Box::new(TcpStream::connect(&addr).await?);
send_message(message, &mut client).await?;
let response_bytes = receive_bytes(&mut client).await?;
let _message: Message = from_slice(&response_bytes)
.map_err(|err| Error::MessageDeserialization(err.to_string()))?;
assert_eq!(response_bytes, original_bytes);
Ok(())
}
}