Start work on compressing RPC messages
Co-Authored-By: Antonio Scandurra <me@as-cii.com>
This commit is contained in:
committed by
Antonio Scandurra
co-authored by
Antonio Scandurra
parent
286846cafd
commit
8bfee93be4
@@ -20,6 +20,7 @@ prost = "0.7"
|
||||
rand = "0.8"
|
||||
rsa = "0.4"
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
zstd = "0.9"
|
||||
|
||||
[build-dependencies]
|
||||
prost-build = { git = "https://github.com/tokio-rs/prost", rev = "6cf97ea422b09d98de34643c4dda2d4f8b7e23e6" }
|
||||
|
||||
+28
-4
@@ -4,6 +4,7 @@ use async_tungstenite::tungstenite::{Error as WebSocketError, Message as WebSock
|
||||
use futures::{SinkExt as _, StreamExt as _};
|
||||
use prost::Message;
|
||||
use std::any::{Any, TypeId};
|
||||
use std::time::Instant;
|
||||
use std::{
|
||||
io,
|
||||
time::{Duration, SystemTime, UNIX_EPOCH},
|
||||
@@ -192,11 +193,15 @@ entity_messages!(channel_id, ChannelMessageSent);
|
||||
/// A stream of protobuf messages.
|
||||
pub struct MessageStream<S> {
|
||||
stream: S,
|
||||
encoding_buffer: Vec<u8>,
|
||||
}
|
||||
|
||||
impl<S> MessageStream<S> {
|
||||
pub fn new(stream: S) -> Self {
|
||||
Self { stream }
|
||||
Self {
|
||||
stream,
|
||||
encoding_buffer: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn inner_mut(&mut self) -> &mut S {
|
||||
@@ -210,10 +215,19 @@ where
|
||||
{
|
||||
/// Write a given protobuf message to the stream.
|
||||
pub async fn write_message(&mut self, message: &Envelope) -> Result<(), WebSocketError> {
|
||||
let mut buffer = Vec::with_capacity(message.encoded_len());
|
||||
self.encoding_buffer.resize(message.encoded_len(), 0);
|
||||
self.encoding_buffer.clear();
|
||||
message
|
||||
.encode(&mut buffer)
|
||||
.encode(&mut self.encoding_buffer)
|
||||
.map_err(|err| io::Error::from(err))?;
|
||||
let t0 = Instant::now();
|
||||
let buffer = zstd::stream::encode_all(self.encoding_buffer.as_slice(), 4).unwrap();
|
||||
eprintln!(
|
||||
"write_message. len: {}, compressed_len: {}, compression time: {:?}",
|
||||
message.encoded_len(),
|
||||
buffer.len(),
|
||||
t0.elapsed(),
|
||||
);
|
||||
self.stream.send(WebSocketMessage::Binary(buffer)).await?;
|
||||
Ok(())
|
||||
}
|
||||
@@ -228,7 +242,17 @@ where
|
||||
while let Some(bytes) = self.stream.next().await {
|
||||
match bytes? {
|
||||
WebSocketMessage::Binary(bytes) => {
|
||||
let envelope = Envelope::decode(bytes.as_slice()).map_err(io::Error::from)?;
|
||||
let t0 = Instant::now();
|
||||
self.encoding_buffer.clear();
|
||||
zstd::stream::copy_decode(bytes.as_slice(), &mut self.encoding_buffer).unwrap();
|
||||
eprintln!(
|
||||
"read_message. len: {}, compressed_len: {}, decompression time: {:?}",
|
||||
self.encoding_buffer.len(),
|
||||
bytes.len(),
|
||||
t0.elapsed()
|
||||
);
|
||||
let envelope = Envelope::decode(self.encoding_buffer.as_slice())
|
||||
.map_err(io::Error::from)?;
|
||||
return Ok(envelope);
|
||||
}
|
||||
WebSocketMessage::Close(_) => break,
|
||||
|
||||
Reference in New Issue
Block a user