kafka_protocol/compression/
lz4.rs

1use crate::protocol::buf::{ByteBuf, ByteBufMut};
2use anyhow::{Context, Result};
3use bytes::{Buf, BufMut, Bytes, BytesMut};
4use lz4::BlockMode;
5use lz4::{Decoder, EncoderBuilder};
6use std::io;
7
8use super::{Compressor, Decompressor};
9
10/// Gzip compression algorithm. See [Kafka's broker configuration](https://kafka.apache.org/documentation/#brokerconfigs_compression.type)
11/// for more information.
12pub struct Lz4;
13
14const COMPRESSION_LEVEL: u32 = 4;
15
16impl<B: ByteBufMut> Compressor<B> for Lz4 {
17    type BufMut = BytesMut;
18    fn compress<R, F>(buf: &mut B, f: F) -> Result<R>
19    where
20        F: FnOnce(&mut Self::BufMut) -> Result<R>,
21    {
22        // Write uncompressed bytes into a temporary buffer
23        let mut tmp = BytesMut::new();
24        let res = f(&mut tmp)?;
25
26        let mut encoder = EncoderBuilder::new()
27            .level(COMPRESSION_LEVEL)
28            .block_mode(BlockMode::Independent)
29            .build(buf.writer())
30            .context("Failed to compress lz4")?;
31
32        io::copy(&mut tmp.reader(), &mut encoder).context("Failed to compress lz4")?;
33        encoder.finish().1.context("Failed to compress lz4")?;
34
35        Ok(res)
36    }
37}
38
39impl<B: ByteBuf> Decompressor<B> for Lz4 {
40    type Buf = Bytes;
41    fn decompress<R, F>(buf: &mut B, f: F) -> Result<R>
42    where
43        F: FnOnce(&mut Self::Buf) -> Result<R>,
44    {
45        let mut tmp = BytesMut::new().writer();
46
47        // Allocate a temporary buffer to hold the uncompressed bytes
48        let buf = buf.copy_to_bytes(buf.remaining());
49
50        let mut decoder = Decoder::new(buf.reader()).context("Failed to decompress lz4")?;
51        io::copy(&mut decoder, &mut tmp).context("Failed to decompress lz4")?;
52
53        f(&mut tmp.into_inner().into())
54    }
55}
56
57#[cfg(test)]
58mod test {
59    use crate::compression::Lz4;
60    use crate::compression::{Compressor, Decompressor};
61    use anyhow::Result;
62    use bytes::BytesMut;
63    use std::fmt::Write;
64    use std::str;
65
66    #[test]
67    fn test_lz4() {
68        let mut compressed = BytesMut::new();
69        Lz4::compress(&mut compressed, |buf| -> Result<()> {
70            buf.write_str("hello lz4").unwrap();
71            Ok(())
72        })
73        .unwrap();
74
75        Lz4::decompress(&mut compressed, |buf| -> Result<()> {
76            let decompressed_str = str::from_utf8(buf.as_ref()).unwrap();
77            assert_eq!(decompressed_str, "hello lz4");
78            Ok(())
79        })
80        .unwrap();
81    }
82}