use crate::models::LogLine;
use flate2::read::{DeflateDecoder, GzDecoder};
use std::io::Read;
use std::str::FromStr;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ContentEncoding {
Zstd,
Gzip,
Deflate,
Identity,
}
impl FromStr for ContentEncoding {
type Err = ();
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s.trim().to_lowercase().as_str() {
"zstd" => Ok(ContentEncoding::Zstd),
"gzip" => Ok(ContentEncoding::Gzip),
"deflate" => Ok(ContentEncoding::Deflate),
"identity" | "" => Ok(ContentEncoding::Identity),
_ => Err(()),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ContentType {
MsgPack,
Json,
}
impl FromStr for ContentType {
type Err = ();
fn from_str(s: &str) -> Result<Self, Self::Err> {
let s_lower = s.to_lowercase();
if s_lower.contains("msgpack") || s_lower.contains("octet-stream") {
Ok(ContentType::MsgPack)
} else if s_lower.contains("json") {
Ok(ContentType::Json)
} else {
Err(())
}
}
}
pub fn decompress_zstd(bytes: &[u8]) -> Result<Vec<u8>, std::io::Error> {
let mut decoder = zstd::stream::read::Decoder::new(bytes)?;
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed)?;
Ok(decompressed)
}
pub fn decompress_gzip(bytes: &[u8]) -> Result<Vec<u8>, std::io::Error> {
let mut decoder = GzDecoder::new(bytes);
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed)?;
Ok(decompressed)
}
pub fn decompress_deflate(bytes: &[u8]) -> Result<Vec<u8>, std::io::Error> {
let mut decoder = DeflateDecoder::new(bytes);
let mut decompressed = Vec::new();
decoder.read_to_end(&mut decompressed)?;
Ok(decompressed)
}
pub fn parse_payload(
bytes: &[u8],
encoding: ContentEncoding,
content_type: ContentType,
) -> Result<Vec<LogLine>, Box<dyn std::error::Error>> {
// Decompress with fallback to identity if decompression fails (e.g. if already decompressed)
let decompressed_bytes = match encoding {
ContentEncoding::Zstd => match decompress_zstd(bytes) {
Ok(decompressed) => decompressed,
Err(_) => bytes.to_vec(),
},
ContentEncoding::Gzip => match decompress_gzip(bytes) {
Ok(decompressed) => decompressed,
Err(_) => bytes.to_vec(),
},
ContentEncoding::Deflate => match decompress_deflate(bytes) {
Ok(decompressed) => decompressed,
Err(_) => bytes.to_vec(),
},
ContentEncoding::Identity => bytes.to_vec(),
};
// Deserialize
let logs = match content_type {
ContentType::MsgPack => match rmp_serde::from_slice::<Vec<LogLine>>(&decompressed_bytes) {
Ok(list) => list,
Err(_) => {
let single = rmp_serde::from_slice::<LogLine>(&decompressed_bytes)?;
vec![single]
}
},
ContentType::Json => match serde_json::from_slice::<Vec<LogLine>>(&decompressed_bytes) {
Ok(list) => list,
Err(_) => {
let single = serde_json::from_slice::<LogLine>(&decompressed_bytes)?;
vec![single]
}
},
};
Ok(logs)
}