Proper decompression of data object payloads. (#21578)
* Remove unused functions. * Rename HashableObject::get_payload to raw_payload. * Handle compressed payloads in data objects. * Support all compression schemes of journal files. * Bump patch version. * Remove stale comment.
vkalintiris committed
Jan 16, 2026 at 11:19 UTC
7f4ce41c142b38ccbe2a6bc0268858e38b541736
11 files changed
+212
-144
src/crates/Cargo.lock
+108
-21
@@ -334,6 +334,15 @@ version = "2.10.0"
334
source = "registry+https://github.com/rust-lang/crates.io-index"
335
checksum = "812e12b5285cc515a9c72a5c1d3b6d46a19dac5acfef5265968c166106e31dd3"
336
337
+[[package]]
338
+name = "block-buffer"
339
+version = "0.10.4"
340
+source = "registry+https://github.com/rust-lang/crates.io-index"
341
+checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71"
342
+dependencies = [
343
+ "generic-array",
344
+]
345
+
346
[[package]]
347
name = "block2"
348
version = "0.6.2"
@@ -345,7 +354,7 @@ dependencies = [
354
355
[[package]]
356
name = "bridge"
348
-version = "0.1.0"
357
+version = "0.1.1"
358
dependencies = [
359
"async-trait",
360
"console-subscriber",
@@ -578,6 +587,21 @@ dependencies = [
587
"libc",
588
]
589
590
+[[package]]
591
+name = "crc"
592
+version = "3.3.0"
593
+source = "registry+https://github.com/rust-lang/crates.io-index"
594
+checksum = "9710d3b3739c2e349eb44fe848ad0b7c8cb1e42bd87ee49371df2f7acaf3e675"
595
+dependencies = [
596
+ "crc-catalog",
597
+]
598
+
599
+[[package]]
600
+name = "crc-catalog"
601
+version = "2.4.0"
602
+source = "registry+https://github.com/rust-lang/crates.io-index"
603
+checksum = "19d374276b40fb8bbdee95aef7c7fa6b5316ec764510eb64b8dd0e2ed0d7e7f5"
604
+
605
[[package]]
606
name = "crc32fast"
607
version = "1.5.0"
@@ -621,6 +645,16 @@ version = "0.8.21"
645
source = "registry+https://github.com/rust-lang/crates.io-index"
646
checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28"
647
648
+[[package]]
649
+name = "crypto-common"
650
+version = "0.1.7"
651
+source = "registry+https://github.com/rust-lang/crates.io-index"
652
+checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a"
653
+dependencies = [
654
+ "generic-array",
655
+ "typenum",
656
+]
657
+
658
[[package]]
659
name = "ctor"
660
version = "0.1.26"
@@ -685,6 +719,16 @@ dependencies = [
719
"powerfmt",
720
]
721
722
+[[package]]
723
+name = "digest"
724
+version = "0.10.7"
725
+source = "registry+https://github.com/rust-lang/crates.io-index"
726
+checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292"
727
+dependencies = [
728
+ "block-buffer",
729
+ "crypto-common",
730
+]
731
+
732
[[package]]
733
name = "dispatch2"
734
version = "0.3.0"
@@ -803,7 +847,7 @@ dependencies = [
847
848
[[package]]
849
name = "flatten_otel"
806
-version = "0.1.0"
850
+version = "0.1.1"
851
dependencies = [
852
"flatten-serde-json",
853
"opentelemetry-proto",
@@ -852,7 +896,7 @@ dependencies = [
896
897
[[package]]
898
name = "foundation"
855
-version = "0.1.0"
899
+version = "0.1.1"
900
dependencies = [
901
"tokio",
902
]
@@ -1081,6 +1125,16 @@ dependencies = [
1125
"byteorder",
1126
]
1127
1128
+[[package]]
1129
+name = "generic-array"
1130
+version = "0.14.7"
1131
+source = "registry+https://github.com/rust-lang/crates.io-index"
1132
+checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a"
1133
+dependencies = [
1134
+ "typenum",
1135
+ "version_check",
1136
+]
1137
+
1138
[[package]]
1139
name = "getrandom"
1140
version = "0.2.16"
@@ -1586,7 +1640,7 @@ dependencies = [
1640
1641
[[package]]
1642
name = "journal-common"
1589
-version = "0.1.0"
1643
+version = "0.1.1"
1644
dependencies = [
1645
"allocative",
1646
"nix",
@@ -1597,7 +1651,7 @@ dependencies = [
1651
1652
[[package]]
1653
name = "journal-core"
1600
-version = "0.1.0"
1654
+version = "0.1.1"
1655
dependencies = [
1656
"allocative",
1657
"anyhow",
@@ -1608,7 +1662,8 @@ dependencies = [
1662
"journal-common",
1663
"journal-registry",
1664
"libc",
1611
- "lz4",
1665
+ "lz4_flex",
1666
+ "lzma-rust2",
1667
"md5",
1668
"memmap2",
1669
"nix",
@@ -1633,12 +1688,11 @@ dependencies = [
1688
"uuid",
1689
"walkdir",
1690
"zerocopy 0.9.0-alpha.0",
1636
- "zstd",
1691
]
1692
1693
[[package]]
1694
name = "journal-engine"
1641
-version = "0.1.0"
1695
+version = "0.1.1"
1696
dependencies = [
1697
"allocative",
1698
"async-stream",
@@ -1667,7 +1721,7 @@ dependencies = [
1721
1722
[[package]]
1723
name = "journal-function"
1670
-version = "0.1.0"
1724
+version = "0.1.1"
1725
dependencies = [
1726
"allocative",
1727
"async-stream",
@@ -1695,7 +1749,7 @@ dependencies = [
1749
1750
[[package]]
1751
name = "journal-index"
1698
-version = "0.1.0"
1752
+version = "0.1.1"
1753
dependencies = [
1754
"allocative",
1755
"journal-common",
@@ -1713,7 +1767,7 @@ dependencies = [
1767
1768
[[package]]
1769
name = "journal-log-writer"
1716
-version = "0.1.0"
1770
+version = "0.1.1"
1771
dependencies = [
1772
"flatten-serde-json",
1773
"journal-common",
@@ -1731,7 +1785,7 @@ dependencies = [
1785
1786
[[package]]
1787
name = "journal-registry"
1734
-version = "0.1.0"
1788
+version = "0.1.1"
1789
dependencies = [
1790
"allocative",
1791
"journal-common",
@@ -1747,7 +1801,7 @@ dependencies = [
1801
1802
[[package]]
1803
name = "journal-viewer-plugin"
1750
-version = "0.1.0"
1804
+version = "0.1.1"
1805
dependencies = [
1806
"anyhow",
1807
"async-trait",
@@ -1878,6 +1932,22 @@ dependencies = [
1932
"libc",
1933
]
1934
1935
+[[package]]
1936
+name = "lz4_flex"
1937
+version = "0.12.0"
1938
+source = "registry+https://github.com/rust-lang/crates.io-index"
1939
+checksum = "ab6473172471198271ff72e9379150e9dfd70d8e533e0752a27e515b48dd375e"
1940
+
1941
+[[package]]
1942
+name = "lzma-rust2"
1943
+version = "0.15.7"
1944
+source = "registry+https://github.com/rust-lang/crates.io-index"
1945
+checksum = "1670343e58806300d87950e3401e820b519b9384281bbabfb15e3636689ffd69"
1946
+dependencies = [
1947
+ "crc",
1948
+ "sha2",
1949
+]
1950
+
1951
[[package]]
1952
name = "madsim"
1953
version = "0.2.34"
@@ -2045,7 +2115,7 @@ dependencies = [
2115
2116
[[package]]
2117
name = "netdata-plugin-charts-derive"
2048
-version = "0.1.0"
2118
+version = "0.1.1"
2119
dependencies = [
2120
"proc-macro2",
2121
"quote",
@@ -2054,14 +2124,14 @@ dependencies = [
2124
2125
[[package]]
2126
name = "netdata-plugin-error"
2057
-version = "0.1.0"
2127
+version = "0.1.1"
2128
dependencies = [
2129
"thiserror 2.0.17",
2130
]
2131
2132
[[package]]
2133
name = "netdata-plugin-protocol"
2064
-version = "0.1.0"
2134
+version = "0.1.1"
2135
dependencies = [
2136
"atoi",
2137
"bytes",
@@ -2080,7 +2150,7 @@ dependencies = [
2150
2151
[[package]]
2152
name = "netdata-plugin-schema"
2083
-version = "0.1.0"
2153
+version = "0.1.1"
2154
dependencies = [
2155
"netdata-plugin-error",
2156
"netdata-plugin-types",
@@ -2091,7 +2161,7 @@ dependencies = [
2161
2162
[[package]]
2163
name = "netdata-plugin-types"
2094
-version = "0.1.0"
2164
+version = "0.1.1"
2165
dependencies = [
2166
"bitflags 2.10.0",
2167
"netdata-plugin-error",
@@ -2447,7 +2517,7 @@ dependencies = [
2517
2518
[[package]]
2519
name = "otel-plugin"
2450
-version = "0.1.0"
2520
+version = "0.1.1"
2521
dependencies = [
2522
"anyhow",
2523
"atty",
@@ -2865,7 +2935,7 @@ dependencies = [
2935
2936
[[package]]
2937
name = "rdp"
2868
-version = "0.1.0"
2938
+version = "0.1.1"
2939
dependencies = [
2940
"md5",
2941
]
@@ -2995,7 +3065,7 @@ dependencies = [
3065
3066
[[package]]
3067
name = "rt"
2998
-version = "0.1.0"
3068
+version = "0.1.1"
3069
dependencies = [
3070
"async-trait",
3071
"bytes",
@@ -3309,6 +3379,17 @@ dependencies = [
3379
"unsafe-libyaml",
3380
]
3381
3382
+[[package]]
3383
+name = "sha2"
3384
+version = "0.10.9"
3385
+source = "registry+https://github.com/rust-lang/crates.io-index"
3386
+checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283"
3387
+dependencies = [
3388
+ "cfg-if",
3389
+ "cpufeatures",
3390
+ "digest",
3391
+]
3392
+
3393
[[package]]
3394
name = "sharded-slab"
3395
version = "0.1.7"
@@ -3935,6 +4016,12 @@ dependencies = [
4016
"rand 0.9.2",
4017
]
4018
4019
+[[package]]
4020
+name = "typenum"
4021
+version = "1.19.0"
4022
+source = "registry+https://github.com/rust-lang/crates.io-index"
4023
+checksum = "562d481066bde0658276a35467c4af00bdc6ee726305698a55b86e61d7ad82bb"
4024
+
4025
[[package]]
4026
name = "uname"
4027
version = "0.1.1"
src/crates/Cargo.toml
+3
-3
@@ -30,7 +30,7 @@ members = [
30
]
31
32
[workspace.package]
33
-version = "0.1.0"
33
+version = "0.1.1"
34
edition = "2024"
35
rust-version = "1.85"
36
@@ -47,12 +47,12 @@ uuid = { version = "1.0"}
47
nix = { version = "0.30" }
48
49
ruzstd = "0.8"
50
-zstd = "0.13"
50
siphasher = "1.0"
51
hashers = "1.0"
52
twox-hash = { version = "2.1", default-features = false }
53
rand = "0.9"
55
-lz4 = "1.25"
54
+lz4_flex = { version = "0.12", default-features = false, features = ["std", "safe-decode", "checked-decode"] }
55
+lzma-rust2 = { version = "0.15", default-features = false, features = ["std", "xz"] }
56
md5 = "0.7"
57
58
roaring = { git = "https://github.com/netdata/roaring-rs.git", branch="allocative" }
src/crates/journal-core/Cargo.toml
+2
-2
@@ -25,7 +25,6 @@ nix = { workspace = true }
25
26
# File format dependencies
27
ruzstd = { workspace = true }
28
-zstd = { workspace = true }
28
siphasher = { workspace = true }
29
md5 = { workspace = true }
30
hashers = { workspace = true }
@@ -33,7 +32,8 @@ hashers = { workspace = true }
32
twox-hash = { workspace = true, features = ["std"] }
33
rand = { workspace = true }
34
36
-lz4 = { workspace = true }
35
+lz4_flex = { workspace = true }
36
+lzma-rust2 = { workspace = true }
37
38
serde = { workspace = true, features = ["derive"] }
39
src/crates/journal-core/src/file/file.rs
+3
-3
@@ -61,7 +61,7 @@ where
61
type Output = NonZeroU64;
62
63
fn visit(&mut self, object: &ValueGuard<'a, Self::Object>) -> Result<Option<Self::Output>> {
64
- if object.hash() == self.hash && object.get_payload() == self.payload {
64
+ if object.hash() == self.hash && object.raw_payload() == self.payload {
65
Ok(Some(object.offset()))
66
} else {
67
Ok(None)
@@ -532,7 +532,7 @@ impl<M: MemoryMap> JournalFile<M> {
532
533
for data_offset in data_offsets.iter().copied() {
534
let data_object = self.data_ref(data_offset)?;
535
- let payload = data_object.payload_bytes();
535
+ let payload = data_object.raw_payload();
536
537
if payload == remapping_payload {
538
continue;
@@ -564,7 +564,7 @@ impl<M: MemoryMap> JournalFile<M> {
564
if field.payload.starts_with(b"ND") {
565
continue;
566
}
567
- let s = String::from_utf8(field.get_payload().to_vec()).expect("utf8 data");
567
+ let s = String::from_utf8(field.raw_payload().to_vec()).expect("utf8 data");
568
field_map.insert(s.clone(), s);
569
}
570
src/crates/journal-core/src/file/object.rs
+59
-6
@@ -11,7 +11,14 @@ pub trait HashableObject {
11
fn hash(&self) -> u64;
12
13
/// Get the payload data for matching
14
- fn get_payload(&self) -> &[u8];
14
+ fn raw_payload(&self) -> &[u8];
15
+
16
+ /// Check if the payload is compressed
17
+ fn is_compressed(&self) -> bool;
18
+
19
+ /// Decompress the payload into the provided buffer.
20
+ /// Returns the number of decompressed bytes.
21
+ fn decompress(&self, buf: &mut Vec<u8>) -> Result<usize>;
22
23
/// Get the offset to the next object in the hash chain
24
fn next_hash_offset(&self) -> Option<NonZeroU64>;
@@ -158,10 +165,20 @@ impl<B: ByteSlice> HashableObject for FieldObject<B> {
165
self.header.hash
166
}
167
161
- fn get_payload(&self) -> &[u8] {
168
+ fn raw_payload(&self) -> &[u8] {
169
&self.payload
170
}
171
172
+ fn is_compressed(&self) -> bool {
173
+ false
174
+ }
175
+
176
+ fn decompress(&self, buf: &mut Vec<u8>) -> Result<usize> {
177
+ buf.clear();
178
+ buf.extend_from_slice(&self.payload);
179
+ Ok(buf.len())
180
+ }
181
+
182
fn next_hash_offset(&self) -> Option<NonZeroU64> {
183
self.header.next_hash_offset
184
}
@@ -186,8 +203,16 @@ impl<B: ByteSlice> HashableObject for DataObject<B> {
203
self.header.hash
204
}
205
189
- fn get_payload(&self) -> &[u8] {
190
- self.payload_bytes()
206
+ fn raw_payload(&self) -> &[u8] {
207
+ self.raw_payload()
208
+ }
209
+
210
+ fn is_compressed(&self) -> bool {
211
+ DataObject::is_compressed(self)
212
+ }
213
+
214
+ fn decompress(&self, buf: &mut Vec<u8>) -> Result<usize> {
215
+ DataObject::decompress(self, buf)
216
}
217
218
fn next_hash_offset(&self) -> Option<NonZeroU64> {
@@ -908,7 +933,7 @@ impl<B: SplitByteSliceMut> JournalObjectMut<B> for DataObject<B> {
933
}
934
935
impl<B: ByteSlice> DataObject<B> {
911
- pub fn payload_bytes(&self) -> &[u8] {
936
+ pub fn raw_payload(&self) -> &[u8] {
937
match &self.payload {
938
DataPayloadType::Regular(payload) => payload,
939
DataPayloadType::Compact { payload, .. } => payload,
@@ -942,10 +967,38 @@ impl<B: ByteSlice> DataObject<B> {
967
use ruzstd::decoding::StreamingDecoder;
968
use ruzstd::io::Read;
969
945
- let payload = self.payload_bytes();
970
+ let payload = self.raw_payload();
971
let mut decoder =
972
StreamingDecoder::new(payload).map_err(|_| JournalError::DecompressorError)?;
973
974
+ buf.clear();
975
+ decoder
976
+ .read_to_end(buf)
977
+ .map_err(|_| JournalError::DecompressorError)
978
+ } else if self.lz4_compressed() {
979
+ let payload = self.raw_payload();
980
+
981
+ // First 8 bytes are the uncompressed size (little-endian u64)
982
+ if payload.len() < 8 {
983
+ return Err(JournalError::DecompressorError);
984
+ }
985
+
986
+ let uncompressed_size =
987
+ u64::from_le_bytes(payload[..8].try_into().unwrap()) as usize;
988
+ let compressed_data = &payload[8..];
989
+
990
+ buf.clear();
991
+ buf.resize(uncompressed_size, 0);
992
+
993
+ lz4_flex::block::decompress_into(compressed_data, buf)
994
+ .map_err(|_| JournalError::DecompressorError)
995
+ } else if self.xz_compressed() {
996
+ use lzma_rust2::XzReader;
997
+ use std::io::Read;
998
+
999
+ let payload = self.raw_payload();
1000
+ let mut decoder = XzReader::new(payload, false);
1001
+
1002
buf.clear();
1003
decoder
1004
.read_to_end(buf)
src/crates/journal-core/src/file/reader.rs
+1
-74
@@ -25,7 +25,6 @@ pub struct JournalReader<'a, M: MemoryMap> {
25
26
// Field name remapping support
27
remapping_registry: FieldMap,
28
- translated_payload: Vec<u8>,
28
}
29
30
impl<M: MemoryMap> std::fmt::Debug for JournalReader<'_, M> {
@@ -49,7 +48,6 @@ impl<M: MemoryMap> Default for JournalReader<'_, M> {
48
field_guard: None,
49
data_guard: None,
50
remapping_registry: FieldMap::new(),
52
- translated_payload: Vec::new(),
51
}
52
}
53
}
@@ -216,63 +214,6 @@ impl<'a, M: MemoryMap> JournalReader<'a, M> {
214
self.entry_data_iterator = None;
215
}
216
219
- pub fn entry_data_enumerate(
220
- &mut self,
221
- journal_file: &'a JournalFile<M>,
222
- ) -> Result<Option<&ValueGuard<'_, DataObject<&'a [u8]>>>> {
223
- self.drop_guards();
224
-
225
- if self.entry_data_iterator.is_none() {
226
- let entry_offset = self.cursor.position()?;
227
- self.entry_data_iterator = Some(journal_file.entry_data_objects(entry_offset)?);
228
- }
229
-
230
- if let Some(iter) = &mut self.entry_data_iterator {
231
- self.data_guard = iter.next().transpose()?;
232
-
233
- // Translate field name if needed
234
- if let Some(data_guard) = &self.data_guard {
235
- let payload = data_guard.get_payload();
236
-
237
- // Check if this field needs translation
238
- if let Some(field_name) = extract_field_name(payload) {
239
- if field_name.starts_with(b"ND_") {
240
- // This looks like a remapped field
241
- if let Ok(systemd_name) = std::str::from_utf8(field_name) {
242
- if let Some(otel_name) =
243
- self.remapping_registry.get_otel_name(systemd_name)
244
- {
245
- // Translate: build new payload with original field name
246
- let eq_pos = payload.iter().position(|&b| b == b'=').unwrap();
247
- let value = &payload[eq_pos..]; // includes '='
248
-
249
- self.translated_payload.clear();
250
- self.translated_payload.extend_from_slice(otel_name);
251
- self.translated_payload.extend_from_slice(value);
252
- } else {
253
- // No mapping found - clear buffer
254
- self.translated_payload.clear();
255
- }
256
- } else {
257
- // Invalid UTF-8 - clear buffer
258
- self.translated_payload.clear();
259
- }
260
- } else {
261
- // Not a remapped field - clear buffer
262
- self.translated_payload.clear();
263
- }
264
- } else {
265
- // No '=' found - clear buffer
266
- self.translated_payload.clear();
267
- }
268
- }
269
-
270
- Ok(self.data_guard.as_ref())
271
- } else {
272
- Ok(None)
273
- }
274
- }
275
-
217
pub fn entry_data_offsets(
218
&self,
219
journal_file: &'a JournalFile<M>,
@@ -391,7 +332,7 @@ impl<'a, M: MemoryMap> JournalReader<'a, M> {
332
let data_iter = journal_file.entry_data_objects(entry_offset)?;
333
for data_result in data_iter {
334
let data_guard = data_result?;
394
- let payload = data_guard.get_payload();
335
+ let payload = data_guard.raw_payload();
336
payloads.push(payload.to_vec());
337
}
338
}
@@ -449,18 +390,4 @@ impl<'a, M: MemoryMap> JournalReader<'a, M> {
390
391
Ok(())
392
}
452
-
453
- /// Gets the current entry data payload, translating remapped field names if applicable.
454
- ///
455
- /// This should be called after `entry_data_enumerate()` to get the translated version
456
- /// of the field name. If the field name doesn't need translation, returns the original.
457
- pub fn get_entry_data_payload(&self) -> &[u8] {
458
- if !self.translated_payload.is_empty() {
459
- &self.translated_payload
460
- } else if let Some(data_guard) = &self.data_guard {
461
- data_guard.get_payload()
462
- } else {
463
- &[]
464
- }
465
- }
393
}
src/crates/journal-core/src/file/value_guard.rs
+10
-2
@@ -75,8 +75,16 @@ impl<T: HashableObject> HashableObject for ValueGuard<'_, T> {
75
self.value.hash()
76
}
77
78
- fn get_payload(&self) -> &[u8] {
79
- self.value.get_payload()
78
+ fn raw_payload(&self) -> &[u8] {
79
+ self.value.raw_payload()
80
+ }
81
+
82
+ fn is_compressed(&self) -> bool {
83
+ self.value.is_compressed()
84
+ }
85
+
86
+ fn decompress(&self, buf: &mut Vec<u8>) -> crate::error::Result<usize> {
87
+ self.value.decompress(buf)
88
}
89
90
fn next_hash_offset(&self) -> Option<NonZeroU64> {
src/crates/journal-engine/src/logs/query.rs
+9
-1
@@ -470,6 +470,9 @@ fn extract_entry_data(log_entries: &[LogEntryId]) -> Result<Vec<LogEntryData>> {
470
// Pre-allocate result vector with exact capacity
471
let mut result = vec![None; log_entries.len()];
472
473
+ // Scratch buffer to keep any decompressed payload of data objects.
474
+ let mut decompress_buf = Vec::new();
475
+
476
// Process each file's entries
477
for (file, file_entries) in entries_by_file {
478
let journal_file = JournalFile::<Mmap>::open(file, 8 * 1024 * 1024)?;
@@ -499,7 +502,12 @@ fn extract_entry_data(log_entries: &[LogEntryId]) -> Result<Vec<LogEntryData>> {
502
let mut fields = Vec::new();
503
for data_offset in data_offsets.iter().copied() {
504
let data_guard = journal_file.data_ref(data_offset)?;
502
- let payload_bytes = data_guard.payload_bytes();
505
+ let payload_bytes = if data_guard.is_compressed() {
506
+ data_guard.decompress(&mut decompress_buf)?;
507
+ &decompress_buf[..]
508
+ } else {
509
+ data_guard.raw_payload()
510
+ };
511
let payload_str = String::from_utf8_lossy(payload_bytes);
512
513
if let Some(mut pair) = FieldValuePair::parse(&payload_str) {
src/crates/journal-index/src/field_types.rs
+1
-1
@@ -230,7 +230,7 @@ pub fn parse_timestamp(
230
field_name: &[u8],
231
data_object: &journal_core::file::DataObject<&[u8]>,
232
) -> crate::Result<u64> {
233
- let payload = data_object.payload_bytes();
233
+ let payload = data_object.raw_payload();
234
235
let value_bytes = FieldValuePair::strip_field_prefix(field_name, payload)
236
.ok_or_else(|| crate::IndexError::InvalidFieldPrefix)?;
src/crates/journal-index/src/file_index.rs
+14
-29
@@ -452,33 +452,6 @@ where
452
Ok(left)
453
}
454
455
-/// Check if an entry matches a regex pattern without using a cache (for benchmarking).
456
-///
457
-/// This is the original implementation that loads and checks each data object
458
-/// for every entry, without caching results.
459
-#[doc(hidden)]
460
-pub fn entry_matches_regex_uncached(
461
- journal_file: &JournalFile<Mmap>,
462
- entry_offset: NonZeroU64,
463
- regex: &Regex,
464
-) -> Result<bool> {
465
- let data_iter = journal_file.entry_data_objects(entry_offset)?;
466
-
467
- for data_result in data_iter {
468
- let data_object = data_result?;
469
- let payload = data_object.payload_bytes();
470
-
471
- // Try to match as UTF-8 string
472
- if let Ok(payload_str) = std::str::from_utf8(payload) {
473
- if regex.is_match(payload_str) {
474
- return Ok(true);
475
- }
476
- }
477
- }
478
-
479
- Ok(false)
480
-}
481
-
455
/// Check if an entry matches a regex pattern
456
fn entry_matches_regex(
457
journal_file: &JournalFile<Mmap>,
@@ -486,6 +459,7 @@ fn entry_matches_regex(
459
regex: &Regex,
460
data_match_cache: &mut HashMap<NonZeroU64, bool>,
461
data_offsets_scratch: &mut Vec<NonZeroU64>,
462
+ scratch_buffer: &mut Vec<u8>,
463
) -> Result<bool> {
464
// Collect all data object offsets for this entry
465
data_offsets_scratch.clear();
@@ -506,9 +480,15 @@ fn entry_matches_regex(
480
481
// Cache miss - load the data object and check if it matches
482
let data_object = journal_file.data_ref(data_offset)?;
509
- let payload = data_object.payload_bytes();
483
511
- let matches = if let Ok(payload_str) = std::str::from_utf8(payload) {
484
+ let payload_bytes = if data_object.is_compressed() {
485
+ data_object.decompress(scratch_buffer)?;
486
+ &scratch_buffer[..]
487
+ } else {
488
+ data_object.raw_payload()
489
+ };
490
+
491
+ let matches = if let Ok(payload_str) = std::str::from_utf8(payload_bytes) {
492
regex.is_match(payload_str)
493
} else {
494
false
@@ -621,6 +601,9 @@ impl FileIndex {
601
);
602
}
603
604
+ // Scratch buffer for compressed payloads of data objects
605
+ let mut scratch_buffer = Vec::new();
606
+
607
match params.direction() {
608
Direction::Forward => {
609
// Determine starting index: use resume_position or binary search
@@ -684,6 +667,7 @@ impl FileIndex {
667
regex,
668
&mut data_match_cache,
669
&mut data_offsets_scratch,
670
+ &mut scratch_buffer,
671
)? {
672
regex_filtered_count += 1;
673
continue;
@@ -783,6 +767,7 @@ impl FileIndex {
767
regex,
768
&mut data_match_cache,
769
&mut data_offsets_scratch,
770
+ &mut scratch_buffer,
771
)? {
772
regex_filtered_count += 1;
773
continue;
src/crates/journal-index/src/file_indexer.rs
+2
-2
@@ -188,12 +188,12 @@ impl FileIndexer {
188
};
189
190
// Skip the remapping value
191
- if data_object.payload_bytes().ends_with(field_name.as_bytes()) {
191
+ if data_object.raw_payload().ends_with(field_name.as_bytes()) {
192
continue;
193
};
194
195
let data_payload =
196
- String::from_utf8_lossy(data_object.payload_bytes()).into_owned();
196
+ String::from_utf8_lossy(data_object.raw_payload()).into_owned();
197
let Some(inlined_cursor) = data_object.inlined_cursor() else {
198
continue;
199
};