Skip to main content

urna_format/encoding/
fsst.rs

1//! fsst (fast static symbol table) text codec over per-chunk canonical
2//! streams (the `fsst` wire codec, encoding id 9). a clean-room 255-entry
3//! static symbol table maps frequent 1-8 byte substrings to single-byte
4//! codes; byte 0xFF is the escape that emits the next raw byte verbatim, so
5//! any input round-trips losslessly. the table is built deterministically
6//! from a single greedy frequency pass and embedded in the payload header.
7//!
8//! this is the `TXT_STREAMS_V3` variant: it reuses the txt_streams container
9//! (kind byte + count + intpack offset table + N frames, O(1) single-chunk
10//! reopen) but each frame is fsst-coded. fsst keeps O(1) single-string decode
11//! and wins on SHORT streams where a zstd frame's overhead dominates.
12//! [`decode`] rebuilds the EXACT canonical payload byte-for-byte, so
13//! `content_hash` is preserved; every read is bounds-checked (typed
14//! `UrnaError`, never a panic on a hostile frame).
15//!
16//! clean-room from the published fsst design (boncz/leis/zukowski, vldb 2020)
17//! as surfaced by duckdb's research (255-entry table, 1-8 byte symbols ->
18//! 1-byte codes, 0xFF escape). NO code is vendored.
19
20use super::fsst_table::{SymbolTable, parse_table, serialize_table};
21use super::intpack::{IntpackReader, pack_u64s};
22use super::txt_streams::{build_canonical, malformed, write_container};
23use crate::bytes::{le_u32, le_u64};
24
25/// kind/version byte for the fsst-framed variant.
26pub const TXT_STREAMS_V3: u8 = 2;
27
28/// escape code: the next byte in the code stream is emitted raw.
29const ESCAPE: u8 = 0xFF;
30
31/// encode one string with `table`: greedy longest-match, escape any byte no
32/// symbol covers. lossless for arbitrary bytes.
33fn encode_one(table: &SymbolTable, input: &[u8]) -> Vec<u8> {
34    let mut out = Vec::with_capacity(input.len());
35    let mut i = 0;
36    while i < input.len() {
37        match table.longest_match(&input[i..]) {
38            Some((code, len)) => {
39                out.push(code);
40                i += len;
41            }
42            None => {
43                out.push(ESCAPE);
44                out.push(input[i]);
45                i += 1;
46            }
47        }
48    }
49    out
50}
51
52/// decode one fsst code stream against the parsed `symbols`. validates that
53/// an escape is never the final byte and that codes are in range; never
54/// panics on a hostile frame.
55fn decode_one(symbols: &[Vec<u8>], codes: &[u8]) -> crate::Result<Vec<u8>> {
56    let mut out = Vec::with_capacity(codes.len() * 2);
57    let mut i = 0;
58    while i < codes.len() {
59        let c = codes[i];
60        if c == ESCAPE {
61            let b = *codes
62                .get(i + 1)
63                .ok_or_else(|| malformed("fsst: trailing escape"))?;
64            out.push(b);
65            i += 2;
66        } else {
67            let sym = symbols
68                .get(c as usize)
69                .ok_or_else(|| malformed("fsst: code out of table range"))?;
70            out.extend_from_slice(sym);
71            i += 1;
72        }
73    }
74    Ok(out)
75}
76
77/// encode `texts` as per-chunk fsst frames behind the shared txt_streams
78/// offset table, with one corpus-wide symbol table embedded after the
79/// container header. a pure function of the inputs, so two builds match.
80pub fn encode(texts: &[String]) -> crate::Result<Vec<u8>> {
81    let corpus: Vec<u8> = texts.iter().flat_map(|t| t.as_bytes().to_vec()).collect();
82    let table = SymbolTable::build(&corpus);
83    let table_blob = serialize_table(&table);
84    let mut streams: Vec<u8> = Vec::new();
85    let mut offsets: Vec<u64> = Vec::with_capacity(texts.len() + 1);
86    offsets.push(0);
87    for t in texts {
88        streams.extend_from_slice(&encode_one(&table, t.as_bytes()));
89        offsets.push(streams.len() as u64);
90    }
91    let off_table = pack_u64s(&offsets);
92    // the container streams region = u32 table_len + symbol table + frames.
93    // the offset table indexes frames relative to the table blob end, so
94    // decode splits the region on the stored table length.
95    let mut framed = Vec::with_capacity(4 + table_blob.len() + streams.len());
96    framed.extend_from_slice(&(table_blob.len() as u32).to_le_bytes());
97    framed.extend_from_slice(&table_blob);
98    framed.extend_from_slice(&streams);
99    Ok(write_container(
100        TXT_STREAMS_V3,
101        texts.len(),
102        &off_table,
103        &framed,
104    ))
105}
106
107/// reconstruct the canonical `chunks_canonical` payload from an fsst-framed
108/// `txt_streams` V3 payload. byte-identical to
109/// `sections::encode_chunks_canonical`, so `content_hash` is preserved.
110pub fn decode(bytes: &[u8]) -> crate::Result<Vec<u8>> {
111    let (count, offsets, framed) = parse_v3(bytes)?;
112    if framed.len() < 4 {
113        return Err(malformed("fsst: truncated region header"));
114    }
115    let table_len = le_u32(&framed[0..4])? as usize;
116    let region = &framed[4..];
117    let (symbols, parsed_len) = parse_table(region)?;
118    if parsed_len != table_len {
119        return Err(malformed("fsst: declared table length mismatch"));
120    }
121    let streams = region
122        .get(table_len..)
123        .ok_or_else(|| malformed("fsst: truncated streams region"))?;
124    if offsets.last().copied() != Some(streams.len() as u64) {
125        return Err(malformed("fsst: final offset != streams length"));
126    }
127    let mut bodies: Vec<Vec<u8>> = Vec::with_capacity(count);
128    for i in 0..count {
129        let start = offsets[i] as usize;
130        let end = offsets[i + 1] as usize;
131        let frame = streams
132            .get(start..end)
133            .ok_or_else(|| malformed("fsst: frame slice out of bounds"))?;
134        let raw = decode_one(&symbols, frame)?;
135        std::str::from_utf8(&raw).map_err(|e| malformed(format!("fsst: invalid utf-8: {}", e)))?;
136        bodies.push(raw);
137    }
138    build_canonical(count, &bodies)
139}
140
141/// parse the V3 container header + intpack offset table, returning the chunk
142/// count, the n+1 byte offsets, and the framed region (table + streams).
143fn parse_v3(bytes: &[u8]) -> crate::Result<(usize, Vec<u64>, &[u8])> {
144    let (kind, rest) = bytes
145        .split_first()
146        .ok_or_else(|| malformed("fsst: empty"))?;
147    if *kind != TXT_STREAMS_V3 {
148        return Err(malformed(format!("fsst: unknown kind {}", *kind)));
149    }
150    if rest.len() < 8 {
151        return Err(malformed("fsst: truncated count"));
152    }
153    let declared = le_u64(&rest[0..8])?;
154    let table_bytes = &rest[8..];
155    let reader = IntpackReader::parse(table_bytes)?;
156    if reader.is_empty() {
157        return Err(malformed("fsst: offset table must hold n+1 >= 1"));
158    }
159    let count = reader.len() - 1;
160    if declared != count as u64 {
161        return Err(malformed("fsst: declared count != offset count - 1"));
162    }
163    let mut offsets = Vec::with_capacity(reader.len());
164    for i in 0..reader.len() {
165        offsets.push(reader.get(i)?);
166    }
167    if offsets[0] != 0 {
168        return Err(malformed("fsst: first offset must be 0"));
169    }
170    let off_len = pack_u64s(&offsets).len();
171    if off_len > table_bytes.len() {
172        return Err(malformed("fsst: truncated offset table"));
173    }
174    let framed = &table_bytes[off_len..];
175    for w in offsets.windows(2) {
176        if w[1] < w[0] {
177            return Err(malformed("fsst: non-monotonic offsets"));
178        }
179    }
180    Ok((count, offsets, framed))
181}
182
183// positive + escape-path coverage lives in tests/fsst_roundtrip.rs and the
184// negative/fuzz coverage in tests/negative_fsst.rs (both exercise the public
185// encode/decode, which drive encode_one/decode_one through every frame).