Skip to content

Commit 89e8a94

Browse files
committed
Address codec review: fix gaps and add missing functionality
- LZW: propagate errors from Deflater::encode() in compress() instead of silently swallowing them and returning truncated output - checksum: add Hasher trait implementation for XXH64 (CRC32, Adler32, and XXH32 already had it), with checksum() returning lower 32 bits - flate: add compress_with_dict() convenience function to match the existing decompress_with_dict() - gzip: always write FHCRC (header CRC-16) in encoder output; add compress_members() for multi-member archive encoding to match decompress_members() - snappy: add streaming FramedInflater for framed format decoding, symmetric with the existing FramedDeflater - async: add async compress/decompress wrappers for zstd, brotli, lz4, and snappy (flate, gzip, zlib, bzip2, lzw already had them)
1 parent 0491ac6 commit 89e8a94

19 files changed

Lines changed: 501 additions & 36 deletions

‎brotli/async/compress.mbt‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
///|
2+
/// Reads raw data from `reader`, compresses it using Brotli, and writes
3+
/// the compressed stream to `writer`.
4+
pub async fn compress(
5+
reader : &@io.Reader,
6+
writer : &@io.Writer,
7+
level? : @brotli.CompressionLevel = Default,
8+
) -> Unit {
9+
let def = @brotli.Deflater::new(level~)
10+
let buf = FixedArray::make(32768, b'\x00')
11+
while true {
12+
let n = reader.read(buf)
13+
if n == 0 {
14+
break
15+
}
16+
let chunk = Bytes::makei(n, fn(i) { buf[i] })
17+
match def.encode(Some(chunk[:])) {
18+
Data(out) => writer.write(out)
19+
Ok => ()
20+
End => return
21+
Error(e) => raise e
22+
}
23+
}
24+
match def.encode(None) {
25+
Data(out) => writer.write(out)
26+
End | Ok => ()
27+
Error(e) => raise e
28+
}
29+
}

‎brotli/async/decompress.mbt‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
///|
2+
/// Reads Brotli-compressed data from `reader`, decompresses it, and writes
3+
/// the decompressed stream to `writer`.
4+
pub async fn decompress(reader : &@io.Reader, writer : &@io.Writer) -> Unit {
5+
let inf = @brotli.Inflater::new()
6+
let buf = FixedArray::make(32768, b'\x00')
7+
loop inf.decode() {
8+
Await => {
9+
let n = reader.read(buf)
10+
if n == 0 {
11+
raise @brotli.BrotliError::UnexpectedEOF
12+
}
13+
let chunk = Bytes::makei(n, fn(i) { buf[i] })
14+
inf.src(chunk[:])
15+
continue inf.decode()
16+
}
17+
Data(out) => {
18+
writer.write(out)
19+
continue inf.decode()
20+
}
21+
End => ()
22+
Error(e) => raise e
23+
}
24+
}

‎brotli/async/moon.pkg‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
import {
2+
"bikallem/compress/brotli",
3+
"moonbitlang/async/io",
4+
}

‎checksum/xxh64.mbt‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -193,3 +193,23 @@ pub fn XXH64::checksum64(self : XXH64) -> UInt64 {
193193
}
194194
xxh64_avalanche(h)
195195
}
196+
197+
///|
198+
pub impl Hasher for XXH64 with size(self) -> Int {
199+
self.size()
200+
}
201+
202+
///|
203+
pub impl Hasher for XXH64 with reset(self) -> Unit {
204+
self.reset()
205+
}
206+
207+
///|
208+
pub impl Hasher for XXH64 with update(self, data) -> Unit {
209+
self.update(data)
210+
}
211+
212+
///|
213+
pub impl Hasher for XXH64 with checksum(self) -> UInt {
214+
self.checksum64().to_uint()
215+
}

‎flate/convenience.mbt‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,20 @@ pub fn compress(
4444
buf.to_bytes()
4545
}
4646

47+
///|
48+
/// Compress data using DEFLATE with a preset dictionary.
49+
pub fn compress_with_dict(
50+
data : Bytes,
51+
dict : Bytes,
52+
level? : CompressionLevel = DefaultCompression,
53+
) -> Bytes raise CompressError {
54+
let e = Deflater::new(level~, dict~)
55+
let buf = @blit.Buffer::new()
56+
drain(e, buf, Some(data[:]))
57+
drain(e, buf, None)
58+
buf.to_bytes()
59+
}
60+
4761
///|
4862
fn drain(
4963
e : Deflater,

‎gzip/compress.mbt‎

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ fn build_gzip_header(header : Header, level : @flate.CompressionLevel) -> Bytes
3030
// Compression method: deflate
3131
buf.write_byte(GZIP_DEFLATE)
3232
// Build FLG byte
33-
let mut flg = 0
33+
let mut flg = GZIP_FHCRC
3434
if header.name.length() > 0 {
3535
flg = flg | GZIP_FNAME
3636
}
@@ -73,6 +73,26 @@ fn build_gzip_header(header : Header, level : @flate.CompressionLevel) -> Bytes
7373
}
7474
buf.write_byte(b'\x00')
7575
}
76+
// FHCRC: CRC-32 of header so far, truncated to 16 bits (LE)
77+
let header_bytes = buf.to_bytes()
78+
let hcrc = @checksum.crc32(header_bytes[:])
79+
buf.write_byte((hcrc & 0xFFU).to_byte())
80+
buf.write_byte(((hcrc >> 8) & 0xFFU).to_byte())
81+
buf.to_bytes()
82+
}
83+
84+
///|
85+
/// Compresses multiple data segments as a multi-member gzip archive (RFC 1952).
86+
/// Each entry produces a separate gzip member; concatenated members form a valid archive.
87+
pub fn compress_members(
88+
members : Array[(Bytes, Header)],
89+
level? : @flate.CompressionLevel = DefaultCompression,
90+
) -> Bytes raise GzipError {
91+
let buf = @blit.Buffer::new()
92+
for entry in members {
93+
let (data, header) = entry
94+
buf.write_bytes(compress(data, level~, header~))
95+
}
7696
buf.to_bytes()
7797
}
7898

‎gzip/deflater_inflater_test.mbt‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -271,7 +271,8 @@ test "streaming gzip inflater preserves repeat state across byte chunks" {
271271
///|
272272
test "streaming gzip inflater validates FHCRC across header chunks" {
273273
let original = b"header crc fixture"
274-
let compressed = gzip_with_fhcrc(@gzip.compress(original))
274+
// Encoder now always writes FHCRC — use directly.
275+
let compressed = @gzip.compress(original)
275276
let d = @gzip.Inflater::new()
276277
let buf = @blit.Buffer::new()
277278
let mut pos = 0

‎gzip/gzip_test.mbt‎

Lines changed: 13 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -126,35 +126,6 @@ fn gzip_named_header(name : String) -> @gzip.Header {
126126
{ name, comment: "", extra: Bytes::default(), mod_time: 0U, os: b'\xFF' }
127127
}
128128

129-
///|
130-
fn gzip_with_fhcrc(data : Bytes) -> Bytes {
131-
let header = Bytes::makei(10, fn(i) {
132-
if i == 3 {
133-
(data[i].to_int() | 0x02).to_byte()
134-
} else {
135-
data[i]
136-
}
137-
})
138-
let hcrc = @checksum.crc32(header[:]) & 0xFFFFU
139-
let out = @blit.Buffer::new(size_hint=data.length() + 2)
140-
out.write_bytes(header)
141-
out.write_byte((hcrc & 0xFFU).to_byte())
142-
out.write_byte(((hcrc >> 8) & 0xFFU).to_byte())
143-
out.write_bytesview(data[10:])
144-
out.to_bytes()
145-
}
146-
147-
///|
148-
fn gzip_with_bad_fhcrc(data : Bytes) -> Bytes {
149-
let with_crc = gzip_with_fhcrc(data)
150-
Bytes::makei(with_crc.length(), fn(i) {
151-
if i == 10 {
152-
(with_crc[i].to_int() ^ 0xFF).to_byte()
153-
} else {
154-
with_crc[i]
155-
}
156-
})
157-
}
158129

159130
///|
160131
fn gzip_hex_digit(c : UInt16) -> Int {
@@ -217,16 +188,26 @@ test "gzip checksum mismatch" {
217188
///|
218189
test "gzip accepts valid FHCRC" {
219190
let data = b"header crc fixture"
220-
let compressed = gzip_with_fhcrc(@gzip.compress(data))
191+
// Encoder now always writes FHCRC — just round-trip directly.
192+
let compressed = @gzip.compress(data)
221193
let (decompressed, _header) = @gzip.decompress(compressed)
222194
assert_eq(decompressed, data)
223195
}
224196

225197
///|
226198
test "gzip rejects invalid FHCRC" {
227199
let data = b"header crc fixture"
228-
let compressed = gzip_with_bad_fhcrc(@gzip.compress(data))
229-
let result = try? @gzip.decompress(compressed)
200+
// Encoder writes FHCRC as last 2 bytes of header (bytes 10-11 for a default header).
201+
// Corrupt byte 10 to invalidate the header CRC.
202+
let compressed = @gzip.compress(data)
203+
let corrupted = Bytes::makei(compressed.length(), fn(i) {
204+
if i == 10 {
205+
(compressed[i].to_int() ^ 0xFF).to_byte()
206+
} else {
207+
compressed[i]
208+
}
209+
})
210+
let result = try? @gzip.decompress(corrupted)
230211
guard result is Err(_) else { fail("expected error for bad header checksum") }
231212
}
232213

‎lz4/async/compress.mbt‎

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
///|
2+
/// Reads raw data from `reader`, compresses it using LZ4 frame format, and writes
3+
/// the compressed stream to `writer`.
4+
pub async fn compress(
5+
reader : &@io.Reader,
6+
writer : &@io.Writer,
7+
options? : @lz4.FrameOptions = { ..@lz4.FrameOptions::default() },
8+
) -> Unit {
9+
let def = @lz4.Deflater::new(options~)
10+
let buf = FixedArray::make(32768, b'\x00')
11+
while true {
12+
let n = reader.read(buf)
13+
if n == 0 {
14+
break
15+
}
16+
let chunk = Bytes::makei(n, fn(i) { buf[i] })
17+
match def.encode(Some(chunk[:])) {
18+
Data(out) => writer.write(out)
19+
Ok => ()
20+
End => return
21+
Error(e) => raise e
22+
}
23+
}
24+
match def.encode(None) {
25+
Data(out) => writer.write(out)
26+
End | Ok => ()
27+
Error(e) => raise e
28+
}
29+
}

‎lz4/async/decompress.mbt‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
///|
2+
/// Reads LZ4-compressed data from `reader`, decompresses it, and writes
3+
/// the decompressed stream to `writer`.
4+
pub async fn decompress(reader : &@io.Reader, writer : &@io.Writer) -> Unit {
5+
let inf = @lz4.Inflater::new()
6+
let buf = FixedArray::make(32768, b'\x00')
7+
loop inf.decode() {
8+
Await => {
9+
let n = reader.read(buf)
10+
if n == 0 {
11+
inf.finish()
12+
continue inf.decode()
13+
}
14+
let chunk = Bytes::makei(n, fn(i) { buf[i] })
15+
inf.src(chunk[:])
16+
continue inf.decode()
17+
}
18+
Data(out) => {
19+
writer.write(out)
20+
continue inf.decode()
21+
}
22+
End => ()
23+
Error(e) => raise e
24+
}
25+
}

0 commit comments

Comments
 (0)