-
Notifications
You must be signed in to change notification settings - Fork 74
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Add transport compression for UploadPath on server
- Loading branch information
Showing
5 changed files
with
100 additions
and
1 deletion.
There are no files selected for viewing
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,85 @@ | ||
//! Module for implementing streaming decompression across multiple | ||
//! algorithms | ||
|
||
use std::{ | ||
io, | ||
pin::Pin, | ||
task::{Context, Poll}, | ||
}; | ||
|
||
use anyhow::anyhow; | ||
use async_compression::tokio::bufread::{ | ||
BrotliDecoder, DeflateDecoder, GzipDecoder, XzDecoder, ZstdDecoder, | ||
}; | ||
use pin_project::pin_project; | ||
use tokio::io::{AsyncBufRead, AsyncRead, BufReader, ReadBuf}; | ||
|
||
use crate::error::{ErrorKind, ServerResult}; | ||
|
||
/// A streaming multi-codec decompressor | ||
#[pin_project(project = SDProj)] | ||
pub enum StreamingDecompressor<S: AsyncBufRead> { | ||
/// None decompression | ||
None(#[pin] S), | ||
/// Brotli decompression | ||
Brotli(#[pin] BrotliDecoder<S>), | ||
/// Deflate decompression | ||
Deflate(#[pin] DeflateDecoder<S>), | ||
/// Gzip decompression | ||
Gzip(#[pin] GzipDecoder<S>), | ||
/// XZ decompression | ||
Xz(#[pin] XzDecoder<S>), | ||
/// Zstd decompression | ||
Zstd(#[pin] ZstdDecoder<S>), | ||
} | ||
|
||
impl<S: AsyncBufRead> StreamingDecompressor<S> { | ||
/// Creates a new streaming decompressor from a buffered stream and compression type. | ||
/// | ||
/// An empty string or "identity" corresponds to no decompression. | ||
/// | ||
/// # Errors | ||
/// This function will return an error if the compression type is invalid | ||
pub fn new(inner: S, kind: &str) -> ServerResult<Self> { | ||
match kind { | ||
"" | "identity" => Ok(Self::None(inner)), | ||
"br" => Ok(Self::Brotli(BrotliDecoder::new(inner))), | ||
"deflate" => Ok(Self::Deflate(DeflateDecoder::new(inner))), | ||
"gzip" => Ok(Self::Gzip(GzipDecoder::new(inner))), | ||
"xz" => Ok(Self::Xz(XzDecoder::new(inner))), | ||
"zstd" => Ok(Self::Zstd(ZstdDecoder::new(inner))), | ||
_ => Err(ErrorKind::RequestError(anyhow!( | ||
"{} is unsupported transport compression", | ||
kind | ||
)) | ||
.into()), | ||
} | ||
} | ||
} | ||
|
||
impl<U: AsyncRead> StreamingDecompressor<BufReader<U>> { | ||
/// Creates a new streaming decompressor from an unbuffered stream and compression type. | ||
/// | ||
/// # Errors | ||
/// This function will return an error if the compression type is invalid | ||
pub fn new_unbuffered(inner: U, kind: &str) -> ServerResult<Self> { | ||
Self::new(BufReader::new(inner), kind) | ||
} | ||
} | ||
|
||
impl<S: AsyncBufRead> AsyncRead for StreamingDecompressor<S> { | ||
fn poll_read( | ||
self: Pin<&mut Self>, | ||
cx: &mut Context<'_>, | ||
buf: &mut ReadBuf<'_>, | ||
) -> Poll<io::Result<()>> { | ||
match self.project() { | ||
SDProj::None(i) => i.poll_read(cx, buf), | ||
SDProj::Brotli(i) => i.poll_read(cx, buf), | ||
SDProj::Deflate(i) => i.poll_read(cx, buf), | ||
SDProj::Gzip(i) => i.poll_read(cx, buf), | ||
SDProj::Xz(i) => i.poll_read(cx, buf), | ||
SDProj::Zstd(i) => i.poll_read(cx, buf), | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters