1use rudb_common::{Error, Result};
42
43use crate::fsst::SymbolTable;
44use crate::integer;
45use crate::reader::Reader;
46
47const MAX_DEPTH: u8 = 2;
50
51pub(crate) const SAMPLE_BYTES: usize = 64 * 1024;
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq)]
60pub enum Kind {
61 Constant = 0,
63 Plain = 1,
65 Fsst = 2,
67 Dict = 3,
69}
70
71impl Kind {
72 fn tag(self) -> u8 {
73 self as u8
74 }
75
76 fn from_tag(tag: u8) -> Result<Self> {
77 match tag {
78 0 => Ok(Self::Constant),
79 1 => Ok(Self::Plain),
80 2 => Ok(Self::Fsst),
81 3 => Ok(Self::Dict),
82 other => Err(Error::internal(format!("unknown string encoding tag {other}"))),
83 }
84 }
85
86 #[must_use]
88 pub fn name(self) -> &'static str {
89 match self {
90 Self::Constant => "CONSTANT",
91 Self::Plain => "PLAIN",
92 Self::Fsst => "FSST",
93 Self::Dict => "DICT",
94 }
95 }
96}
97
98pub fn encode(values: &[&[u8]]) -> Result<Vec<u8>> {
105 encode_at(values, 0)
106}
107
108pub fn decode_prefix(bytes: &[u8]) -> Result<(Vec<Vec<u8>>, usize)> {
117 let mut reader = Reader::new(bytes);
118 let values = decode_chunk(&mut reader)?;
119 Ok((values, reader.used()))
120}
121
122pub fn describe_prefix(bytes: &[u8]) -> Result<(String, usize)> {
128 let mut reader = Reader::new(bytes);
129 let text = describe_chunk(&mut reader)?;
130 Ok((text, reader.used()))
131}
132
133pub fn decode(bytes: &[u8]) -> Result<Vec<Vec<u8>>> {
139 let mut reader = Reader::new(bytes);
140 let values = decode_chunk(&mut reader)?;
141 if reader.remaining() != 0 {
142 return Err(Error::internal(format!(
143 "{} bytes left over after decoding a string chunk",
144 reader.remaining()
145 )));
146 }
147 Ok(values)
148}
149
150pub fn candidate_sizes(values: &[&[u8]]) -> Result<Vec<(Kind, usize)>> {
157 let mut sizes = Vec::new();
158 for kind in candidates(values, 0) {
159 if let Some(bytes) = encode_as(kind, values, 0)? {
160 sizes.push((kind, bytes.len()));
161 }
162 }
163 Ok(sizes)
164}
165
166pub fn describe(bytes: &[u8]) -> Result<String> {
172 let mut reader = Reader::new(bytes);
173 describe_chunk(&mut reader)
174}
175
176fn encode_at(values: &[&[u8]], depth: u8) -> Result<Vec<u8>> {
177 let mut best: Option<Vec<u8>> = None;
178 for kind in candidates(values, depth) {
179 let Some(bytes) = encode_as(kind, values, depth)? else {
180 continue;
181 };
182 if best.as_ref().is_none_or(|current| bytes.len() < current.len()) {
183 best = Some(bytes);
184 }
185 }
186 best.ok_or_else(|| Error::internal("no string encoding applied to the chunk"))
187}
188
189fn candidates(values: &[&[u8]], depth: u8) -> Vec<Kind> {
190 let mut kinds = vec![Kind::Plain];
191 if values.is_empty() {
192 return kinds;
193 }
194 if values.iter().all(|value| *value == values[0]) {
195 return vec![Kind::Constant];
196 }
197 kinds.push(Kind::Fsst);
198 if depth < MAX_DEPTH && distinct_values(values).len() < values.len() {
199 kinds.push(Kind::Dict);
200 }
201 kinds
202}
203
204fn encode_as(kind: Kind, values: &[&[u8]], depth: u8) -> Result<Option<Vec<u8>>> {
205 let mut out = vec![kind.tag()];
206 put_u32(&mut out, u32::try_from(values.len()).map_err(|_| too_long(values.len()))?);
207 match kind {
208 Kind::Constant => {
209 let Some(first) = values.first() else {
210 return Ok(None);
211 };
212 if values.iter().any(|value| value != first) {
213 return Ok(None);
214 }
215 put_u32(&mut out, u32::try_from(first.len()).map_err(|_| too_long(first.len()))?);
216 out.extend_from_slice(first);
217 }
218 Kind::Plain => {
219 out.extend_from_slice(&encode_lengths(values)?);
220 for value in values {
221 out.extend_from_slice(value);
222 }
223 }
224 Kind::Fsst => {
225 let sample = sample_of(values);
226 let table = SymbolTable::train(&sample);
227 if table.is_empty() {
228 return Ok(None);
229 }
230 let mut compressed = Vec::new();
231 let mut lengths = Vec::with_capacity(values.len());
232 for value in values {
233 let before = compressed.len();
234 table.compress(value, &mut compressed);
235 lengths.push((compressed.len() - before) as i64);
236 }
237 table.serialize(&mut out);
238 out.extend_from_slice(&integer::encode(&lengths)?);
239 out.extend_from_slice(&compressed);
240 }
241 Kind::Dict => {
242 let dictionary = distinct_values(values);
243 if dictionary.is_empty() {
244 return Ok(None);
245 }
246 let codes = codes_over(values, &dictionary);
247 let entries: Vec<&[u8]> = dictionary.iter().map(Vec::as_slice).collect();
248 out.extend_from_slice(&encode_at(&entries, depth + 1)?);
249 out.extend_from_slice(&integer::encode(&codes)?);
250 }
251 }
252 Ok(Some(out))
253}
254
255fn decode_chunk(reader: &mut Reader<'_>) -> Result<Vec<Vec<u8>>> {
256 let kind = Kind::from_tag(reader.u8()?)?;
257 let count = reader.u32()? as usize;
258 match kind {
259 Kind::Constant => {
260 let len = reader.u32()? as usize;
261 let value = reader.bytes(len)?.to_vec();
262 Ok(vec![value; count])
263 }
264 Kind::Plain => {
265 let lengths = decode_lengths(reader, count)?;
266 let mut values = Vec::with_capacity(count);
267 for length in lengths {
268 values.push(reader.bytes(length)?.to_vec());
269 }
270 Ok(values)
271 }
272 Kind::Fsst => {
273 let (table, used) = SymbolTable::deserialize(reader.rest())?;
274 reader.skip(used)?;
275 let lengths = decode_lengths(reader, count)?;
276 let mut values = Vec::with_capacity(count);
277 for length in lengths {
278 let compressed = reader.bytes(length)?;
279 let mut value = Vec::new();
280 table.decompress(compressed, &mut value)?;
281 values.push(value);
282 }
283 Ok(values)
284 }
285 Kind::Dict => {
286 let dictionary = decode_chunk(reader)?;
287 let codes = decode_integers(reader)?;
288 if codes.len() != count {
289 return Err(Error::internal(format!(
290 "a dictionary chunk says it holds {count} values and has {} codes",
291 codes.len()
292 )));
293 }
294 let mut values = Vec::with_capacity(count);
295 for code in codes {
296 let entry =
297 usize::try_from(code).ok().and_then(|index| dictionary.get(index)).ok_or_else(
298 || Error::internal(format!("code {code} is not in the dictionary")),
299 )?;
300 values.push(entry.clone());
301 }
302 Ok(values)
303 }
304 }
305}
306
307fn describe_chunk(reader: &mut Reader<'_>) -> Result<String> {
308 let kind = Kind::from_tag(reader.u8()?)?;
309 let count = reader.u32()? as usize;
310 Ok(match kind {
311 Kind::Constant => {
312 let len = reader.u32()? as usize;
313 reader.bytes(len)?;
314 "CONSTANT".to_string()
315 }
316 Kind::Plain => {
317 let (shape, lengths) = describe_lengths(reader, count)?;
318 reader.skip(lengths.iter().sum())?;
319 format!("PLAIN({shape})")
320 }
321 Kind::Fsst => {
322 let (table, used) = SymbolTable::deserialize(reader.rest())?;
323 reader.skip(used)?;
324 let (shape, lengths) = describe_lengths(reader, count)?;
325 reader.skip(lengths.iter().sum())?;
326 format!("FSST[{}]({shape})", table.len())
327 }
328 Kind::Dict => {
329 let entries = describe_chunk(reader)?;
330 let codes = describe_integers(reader)?;
331 format!("DICT({entries}, {codes})")
332 }
333 })
334}
335
336fn describe_lengths(reader: &mut Reader<'_>, count: usize) -> Result<(String, Vec<usize>)> {
340 let (shape, _) = integer::describe_prefix(reader.rest())?;
341 let lengths = decode_lengths(reader, count)?;
342 Ok((shape, lengths))
343}
344
345fn encode_lengths(values: &[&[u8]]) -> Result<Vec<u8>> {
346 let lengths: Vec<i64> = values.iter().map(|value| value.len() as i64).collect();
347 integer::encode(&lengths)
348}
349
350fn decode_lengths(reader: &mut Reader<'_>, count: usize) -> Result<Vec<usize>> {
351 let lengths = decode_integers(reader)?;
352 if lengths.len() != count {
353 return Err(Error::internal(format!(
354 "a string chunk says it holds {count} values and has {} lengths",
355 lengths.len()
356 )));
357 }
358 lengths
359 .into_iter()
360 .map(|length| {
361 usize::try_from(length).map_err(|_| Error::internal("a negative string length"))
362 })
363 .collect()
364}
365
366fn decode_integers(reader: &mut Reader<'_>) -> Result<Vec<i64>> {
370 let (values, used) = integer::decode_prefix(reader.rest())?;
371 reader.skip(used)?;
372 Ok(values)
373}
374
375fn describe_integers(reader: &mut Reader<'_>) -> Result<String> {
376 let (text, used) = integer::describe_prefix(reader.rest())?;
377 reader.skip(used)?;
378 Ok(text)
379}
380
381pub(crate) fn sample_of<'a>(values: &[&'a [u8]]) -> Vec<&'a [u8]> {
400 sample_bytes_of(values, SAMPLE_BYTES)
401}
402
403pub(crate) fn sample_bytes_of<'a>(values: &[&'a [u8]], budget: usize) -> Vec<&'a [u8]> {
406 let budget = budget.max(1);
407 let total: usize = values.iter().map(|value| value.len()).sum();
408 if total <= budget {
409 return values.to_vec();
410 }
411 let stride = total.div_ceil(budget).max(1);
412 let span = (stride * 2 - 1).max(1) as u64;
413 let mut state = 0x2545_f491_4f6c_dd1du64;
414 let mut sample = Vec::with_capacity(values.len() / stride + 1);
415 let mut at = 0usize;
416 while at < values.len() {
417 sample.push(values[at]);
418 state ^= state << 13;
419 state ^= state >> 7;
420 state ^= state << 17;
421 at += 1 + (state % span) as usize;
422 }
423 sample
424}
425
426fn distinct_values(values: &[&[u8]]) -> Vec<Vec<u8>> {
429 let mut distinct: Vec<Vec<u8>> = values.iter().map(|value| value.to_vec()).collect();
430 distinct.sort_unstable();
431 distinct.dedup();
432 distinct
433}
434
435fn codes_over(values: &[&[u8]], dictionary: &[Vec<u8>]) -> Vec<i64> {
436 values
437 .iter()
438 .map(|value| {
439 dictionary
440 .binary_search_by(|entry| entry.as_slice().cmp(value))
441 .expect("the dictionary is the distinct values of this chunk") as i64
442 })
443 .collect()
444}
445
446fn too_long(len: usize) -> Error {
447 Error::internal(format!("a string chunk of {len} is longer than the format allows"))
448}
449
450fn put_u32(out: &mut Vec<u8>, value: u32) {
451 out.extend_from_slice(&value.to_le_bytes());
452}
453
454#[cfg(test)]
455mod tests {
456 use super::*;
457
458 fn urls(count: usize) -> Vec<Vec<u8>> {
459 let hosts = ["www.example.com", "shop.example.com", "news.other.example.org"];
460 let paths = ["/index.html", "/catalog/item", "/search", "/user/profile/settings"];
461 (0..count)
462 .map(|index| {
463 let host = hosts[index % hosts.len()];
464 let path = paths[(index / 3) % paths.len()];
465 format!("http://{host}{path}?session={}&ref=google", index * 7).into_bytes()
466 })
467 .collect()
468 }
469
470 fn borrow(values: &[Vec<u8>]) -> Vec<&[u8]> {
471 values.iter().map(Vec::as_slice).collect()
472 }
473
474 fn round_trip(values: &[Vec<u8>]) -> Vec<u8> {
475 let borrowed = borrow(values);
476 let bytes = encode(&borrowed).unwrap();
477 let back = decode(&bytes).unwrap();
478 assert_eq!(back, values, "{}", describe(&bytes).unwrap());
479 bytes
480 }
481
482 fn kind_of(bytes: &[u8]) -> Kind {
483 Kind::from_tag(bytes[0]).unwrap()
484 }
485
486 fn raw_size(values: &[Vec<u8>]) -> usize {
487 values.iter().map(Vec::len).sum::<usize>() + values.len() * 4
488 }
489
490 #[test]
491 fn an_empty_chunk_round_trips() {
492 let bytes = round_trip(&[]);
493 assert_eq!(kind_of(&bytes), Kind::Plain);
494 }
495
496 #[test]
497 fn a_constant_column_costs_what_one_value_costs() {
498 let values = vec![b"https://www.example.com/".to_vec(); 100_000];
499 let bytes = round_trip(&values);
500 assert_eq!(kind_of(&bytes), Kind::Constant);
501 assert_eq!(bytes.len(), 9 + 24);
502 }
503
504 #[test]
505 fn a_url_column_of_unique_values_uses_fsst() {
506 let values = urls(20_000);
510 let bytes = round_trip(&values);
511 assert_eq!(kind_of(&bytes), Kind::Fsst);
512 let ratio = raw_size(&values) as f64 / bytes.len() as f64;
513 assert!(ratio > 5.0, "{ratio:.2}x");
514 }
515
516 #[test]
517 fn a_sample_of_a_periodic_column_learns_every_phase_of_it() {
518 let values = urls(20_000);
523 let borrowed = borrow(&values);
524 let sample = sample_of(&borrowed);
525 let mut phases: Vec<&[u8]> = sample
526 .iter()
527 .map(|value| {
528 let query =
529 value.iter().position(|byte| *byte == b'?').expect("every value has a query");
530 &value[..query]
531 })
532 .collect();
533 phases.sort_unstable();
534 phases.dedup();
535 assert_eq!(phases.len(), 12);
537 let whole = SymbolTable::train(&borrowed);
538 let sampled = SymbolTable::train(&sample);
539 let mut on_whole = Vec::new();
540 let mut on_sample = Vec::new();
541 for value in &borrowed {
542 whole.compress(value, &mut on_whole);
543 sampled.compress(value, &mut on_sample);
544 }
545 assert!(
548 on_sample.len() < on_whole.len() * 5 / 4,
549 "{} against {}",
550 on_sample.len(),
551 on_whole.len()
552 );
553 }
554
555 #[test]
556 fn a_repeating_column_becomes_a_dictionary_of_compressed_entries() {
557 let distinct = urls(500);
560 let values: Vec<Vec<u8>> =
561 (0..50_000).map(|index| distinct[index * 7919 % distinct.len()].clone()).collect();
562 let bytes = round_trip(&values);
563 assert_eq!(kind_of(&bytes), Kind::Dict);
564 let shape = describe(&bytes).unwrap();
565 assert!(shape.starts_with("DICT(FSST"), "{shape}");
566 let ratio = raw_size(&values) as f64 / bytes.len() as f64;
567 assert!(ratio > 20.0, "{ratio:.2}x, {shape}");
568 }
569
570 #[test]
571 fn a_column_of_long_runs_costs_almost_nothing() {
572 let distinct = urls(50);
575 let mut values = Vec::new();
576 for entry in &distinct {
577 values.extend(std::iter::repeat_n(entry.clone(), 1000));
578 }
579 let bytes = round_trip(&values);
580 let shape = describe(&bytes).unwrap();
581 assert!(shape.contains("RLE"), "{shape}");
582 assert!(bytes.len() < 2000, "{} bytes: {shape}", bytes.len());
583 }
584
585 #[test]
586 fn incompressible_strings_stay_close_to_their_own_size() {
587 let mut state = 0x2545_f491_4f6c_dd1du64;
590 let values: Vec<Vec<u8>> = (0..2000)
591 .map(|_| {
592 (0..32)
593 .map(|_| {
594 state ^= state << 13;
595 state ^= state >> 7;
596 state ^= state << 17;
597 state as u8
598 })
599 .collect()
600 })
601 .collect();
602 let bytes = round_trip(&values);
603 assert!(bytes.len() < 2000 * 32 + 3000, "{} bytes", bytes.len());
604 }
605
606 #[test]
607 fn lengths_are_stored_rather_than_offsets() {
608 let values: Vec<Vec<u8>> =
611 (0..100_000).map(|index| format!("{index:024}").into_bytes()).collect();
612 let borrowed = borrow(&values);
613 let bytes = encode_as(Kind::Plain, &borrowed, 0).unwrap().unwrap();
614 assert_eq!(bytes.len(), 5 + 13 + 100_000 * 24);
615 }
616
617 #[test]
618 fn empty_strings_are_values_and_not_nulls() {
619 let values = vec![Vec::new(), b"a".to_vec(), Vec::new(), b"bb".to_vec()];
620 round_trip(&values);
621 }
622
623 #[test]
624 fn a_chunk_with_one_value_round_trips() {
625 round_trip(&[b"only".to_vec()]);
626 }
627
628 #[test]
629 fn every_candidate_that_applies_decodes_to_the_input() {
630 let values = urls(3000);
631 let borrowed = borrow(&values);
632 let applicable = candidates(&borrowed, 0);
633 assert!(applicable.len() >= 2, "{applicable:?}");
634 for kind in applicable {
635 let bytes = encode_as(kind, &borrowed, 0).unwrap().unwrap();
636 assert_eq!(decode(&bytes).unwrap(), values, "{}", kind.name());
637 }
638 }
639
640 #[test]
641 fn the_chooser_picks_the_smallest_candidate() {
642 let values = urls(2000);
643 let borrowed = borrow(&values);
644 let chosen = encode(&borrowed).unwrap();
645 for (_, size) in candidate_sizes(&borrowed).unwrap() {
646 assert!(chosen.len() <= size);
647 }
648 }
649
650 #[test]
651 fn a_truncated_chunk_is_an_error_and_not_a_panic() {
652 let values = urls(40);
653 let bytes = encode(&borrow(&values)).unwrap();
654 for len in 0..bytes.len() {
655 assert!(decode(&bytes[..len]).is_err(), "{len} bytes decoded");
656 }
657 }
658
659 #[test]
660 fn trailing_bytes_are_an_error() {
661 let mut bytes = encode(&borrow(&urls(10))).unwrap();
662 bytes.push(0);
663 let error = decode(&bytes).unwrap_err();
664 assert!(error.message().contains("left over"), "{error}");
665 }
666
667 #[test]
668 fn an_unknown_tag_is_an_error() {
669 let error = decode(&[99, 0, 0, 0, 0]).unwrap_err();
670 assert!(error.message().contains("unknown string encoding tag"), "{error}");
671 }
672
673 #[test]
674 fn a_dictionary_code_outside_the_dictionary_is_an_error() {
675 let mut bytes = vec![Kind::Dict.tag()];
676 put_u32(&mut bytes, 1);
677 bytes.extend_from_slice(&encode(&[b"one".as_slice()]).unwrap());
678 bytes.extend_from_slice(&integer::encode(&[9]).unwrap());
679 let error = decode(&bytes).unwrap_err();
680 assert!(error.message().contains("not in the dictionary"), "{error}");
681 }
682
683 #[test]
684 fn a_negative_length_is_an_error() {
685 let mut bytes = vec![Kind::Plain.tag()];
686 put_u32(&mut bytes, 1);
687 bytes.extend_from_slice(&integer::encode(&[-1]).unwrap());
688 let error = decode(&bytes).unwrap_err();
689 assert!(error.message().contains("negative string length"), "{error}");
690 }
691
692 #[test]
693 fn the_sample_is_spread_across_the_chunk_and_not_taken_from_the_front() {
694 let mut values: Vec<Vec<u8>> = Vec::new();
697 for index in 0..20_000 {
698 let head = if index < 10_000 { "aaaaaaaaaaaaaaaa" } else { "zzzzzzzzzzzzzzzz" };
699 values.push(format!("{head}/{index:08}").into_bytes());
700 }
701 let borrowed = borrow(&values);
702 let sample = sample_of(&borrowed);
703 let first_half = sample.iter().filter(|value| value.starts_with(b"aaaa")).count();
704 let second_half = sample.len() - first_half;
705 assert!(first_half > 0 && second_half > 0, "{first_half} and {second_half}");
706 let bytes = round_trip(&values);
707 let ratio = raw_size(&values) as f64 / bytes.len() as f64;
708 assert!(ratio > 4.0, "{ratio:.2}x");
709 }
710}