Skip to main content

urna_format/encoding/
txt_streams.rs

1//! per-chunk independent canonical-text streams (the `txt_streams` wire
2//! codec, encoding id 10). re-layouts the `chunks_canonical` (0x02)
3//! section's COMPRESSED form from one concatenated zstd-19 blob into N
4//! independently zstd-encoded streams (one per canonical string) plus an
5//! intpack offset table that gives O(1) seek to any single chunk.
6//!
7//! this is the named prerequisite for #11's dict/fsst text levers: a
8//! per-chunk frame is where a trained dict or fsst can beat one big
9//! zstd-19 blob, and it is the O(1) single-chunk reopen layout (decode ONE
10//! stream for cite/materialize instead of inflating the whole section).
11//! a per-chunk frame loses cross-chunk LZ context, so expect a small zstd
12//! size INCREASE today; the win is what it unlocks. draws from facebook
13//! zstd (self-describing per-record frames), lancedb/lance (transparent
14//! per-page string compression with a parallel offset table) and
15//! flatbuffers (an offset vector giving O(1) random access off mmap).
16//!
17//! [`decode`] rebuilds the EXACT canonical payload byte-for-byte, so
18//! `content_hash` (hashed over decoded bytes) and `urna://` citations are
19//! unchanged. every read is bounds-checked and returns a typed
20//! `UrnaError`, never a panic on a truncated or hostile payload.
21//!
22//! container layout (mirrors the intpack/spans-repack discipline; shared by
23//! the V1 plain-zstd, V2 dict-framed, and V3 fsst-framed variants):
24//!
25//! ```text
26//! [0]        u8  kind/version  (V1 plain / V2 dict / V3 fsst)
27//! [1..9]     u64 chunk count   (LE)
28//! [9..9+T]   intpack offset table of N+1 byte offsets into the streams
29//!            region (pack_u64s, encoding-id-4 primitive, reused)
30//! [9+T..]    N per-stream frames, stream i = bytes [off[i] .. off[i+1])
31//! ```
32
33use super::intpack::{IntpackReader, pack_u64s};
34use super::zstd_codec::{zstd_decode, zstd_encode};
35use crate::bytes::le_u64;
36use crate::error::UrnaError;
37use crate::layout::{
38    SECTION_CHUNKS_CANONICAL, SECTION_PAYLOAD_PREFIX_SIZE, SECTION_PAYLOAD_VERSION,
39};
40
41/// leading kind/version byte. v1 = plain per-stream zstd. the dict (V2) and
42/// fsst (V3) variants live in `zstd_dict.rs` / `fsst.rs` and claim the next
43/// values here, reusing this container without a new encoding id beyond
44/// their section-entry codec id.
45pub const TXT_STREAMS_V1: u8 = 0;
46
47/// container header: u8 kind + u64 count. shared by all variants.
48pub(super) const HEADER: usize = 9;
49
50pub(super) fn malformed(reason: impl Into<String>) -> UrnaError {
51    UrnaError::MalformedSectionPayload {
52        section_id: SECTION_CHUNKS_CANONICAL,
53        reason: reason.into(),
54    }
55}
56
57/// assemble the shared container: kind byte + count + intpack offset table +
58/// the concatenated per-stream frames. a pure function of the inputs, so two
59/// builds are byte-identical. shared by the V1/V2/V3 encoders.
60pub(super) fn write_container(kind: u8, count: usize, table: &[u8], streams: &[u8]) -> Vec<u8> {
61    let mut out = Vec::with_capacity(HEADER + table.len() + streams.len());
62    out.push(kind);
63    out.extend_from_slice(&(count as u64).to_le_bytes());
64    out.extend_from_slice(table);
65    out.extend_from_slice(streams);
66    out
67}
68
69/// rebuild the canonical (raw-encoding) `chunks_canonical` payload from the
70/// already-decompressed per-chunk bodies (utf-8 validated by the caller).
71/// byte-identical to `sections::encode_chunks_canonical`. shared by all
72/// variants so `content_hash` rebuilds the same way regardless of codec.
73pub(super) fn build_canonical(count: usize, bodies: &[Vec<u8>]) -> crate::Result<Vec<u8>> {
74    let total: usize = bodies.iter().map(|b| b.len()).sum();
75    let mut out = Vec::with_capacity(SECTION_PAYLOAD_PREFIX_SIZE + count * 4 + total);
76    out.extend_from_slice(&SECTION_PAYLOAD_VERSION.to_le_bytes());
77    out.extend_from_slice(&(count as u64).to_le_bytes());
78    for raw in bodies {
79        let len = u32::try_from(raw.len())
80            .map_err(|_| malformed("txt_streams: stream longer than u32"))?;
81        out.extend_from_slice(&len.to_le_bytes());
82        out.extend_from_slice(raw);
83    }
84    Ok(out)
85}
86
87/// encode `texts` (the canonical strings, in chunk order) as per-chunk
88/// independent zstd streams behind an intpack offset table. the layout is
89/// a pure function of the inputs, so two builds are byte-identical.
90pub fn encode_txt_streams(texts: &[String]) -> crate::Result<Vec<u8>> {
91    let mut streams: Vec<u8> = Vec::new();
92    // n+1 offsets so stream i is [off[i] .. off[i+1]); off[0] == 0 and the
93    // last is the total streams length, giving O(1) seek and exact bounds.
94    let mut offsets: Vec<u64> = Vec::with_capacity(texts.len() + 1);
95    offsets.push(0);
96    for t in texts {
97        streams.extend_from_slice(&zstd_encode(t.as_bytes())?);
98        offsets.push(streams.len() as u64);
99    }
100    let table = pack_u64s(&offsets);
101    Ok(write_container(
102        TXT_STREAMS_V1,
103        texts.len(),
104        &table,
105        &streams,
106    ))
107}
108
109/// a parsed `txt_streams` payload. `parse` validates the header and offset
110/// table once; `stream`/`text` then reach any single chunk in O(1) without
111/// touching the others (the O(1)-reopen enabler). `decode_payload` rebuilds
112/// the full canonical section by concatenating all of them.
113pub struct TxtStreams<'a> {
114    streams: &'a [u8],
115    offsets: Vec<u64>,
116    count: usize,
117}
118
119impl<'a> TxtStreams<'a> {
120    pub fn parse(bytes: &'a [u8]) -> crate::Result<Self> {
121        let (kind, rest) = bytes
122            .split_first()
123            .ok_or_else(|| malformed("txt_streams: empty"))?;
124        if *kind != TXT_STREAMS_V1 {
125            return Err(malformed(format!("txt_streams: unknown kind {}", *kind)));
126        }
127        if rest.len() < 8 {
128            return Err(malformed("txt_streams: truncated count"));
129        }
130        let declared = le_u64(&rest[0..8])?;
131        let table_bytes = &rest[8..];
132        // the offset table is itself an intpack payload; IntpackReader
133        // parses its header/directory and refuses an oversized claim, so the
134        // n+1 offsets are the bounded source of truth (never the raw count).
135        let reader = IntpackReader::parse(table_bytes)?;
136        if reader.is_empty() {
137            return Err(malformed("txt_streams: offset table must hold n+1 >= 1"));
138        }
139        let count = reader.len() - 1;
140        // n streams require exactly n+1 offsets; the declared count (checked
141        // against the table) must agree without overflowing on a hostile u64.
142        if declared != count as u64 {
143            return Err(malformed("txt_streams: declared count != offset count - 1"));
144        }
145        let mut offsets = Vec::with_capacity(reader.len());
146        for i in 0..reader.len() {
147            offsets.push(reader.get(i)?);
148        }
149        if offsets[0] != 0 {
150            return Err(malformed("txt_streams: first offset must be 0"));
151        }
152        // the offset table sits before the streams region. pack_u64s output
153        // length is a deterministic function of the offsets, so re-pack the
154        // parsed offsets and measure to locate the streams start without
155        // storing a separate table length.
156        let table_len = pack_u64s(&offsets).len();
157        if table_len > table_bytes.len() {
158            return Err(malformed("txt_streams: truncated offset table"));
159        }
160        let streams = &table_bytes[table_len..];
161        // offsets must be monotonic non-decreasing and end at the streams
162        // length, so every stream slice is in-bounds and the layout is exact.
163        for w in offsets.windows(2) {
164            if w[1] < w[0] {
165                return Err(malformed("txt_streams: non-monotonic offsets"));
166            }
167        }
168        if offsets.last().copied() != Some(streams.len() as u64) {
169            return Err(malformed("txt_streams: final offset != streams length"));
170        }
171        Ok(Self {
172            streams,
173            offsets,
174            count,
175        })
176    }
177
178    #[inline]
179    pub fn len(&self) -> usize {
180        self.count
181    }
182
183    #[inline]
184    pub fn is_empty(&self) -> bool {
185        self.count == 0
186    }
187
188    /// the raw zstd bytes of stream `i`, bounds-checked. O(1) seek.
189    fn stream(&self, i: usize) -> crate::Result<&'a [u8]> {
190        if i >= self.count {
191            return Err(malformed("txt_streams: index out of range"));
192        }
193        let start = self.offsets[i] as usize;
194        let end = self.offsets[i + 1] as usize;
195        self.streams
196            .get(start..end)
197            .ok_or_else(|| malformed("txt_streams: stream slice out of bounds"))
198    }
199
200    /// decode a SINGLE chunk's canonical text in O(1) (the runtime reopen
201    /// path decodes one stream for cite/materialize instead of the whole
202    /// section). validates utf-8 so a hostile stream cannot smuggle non-utf8.
203    pub fn text(&self, i: usize) -> crate::Result<String> {
204        let raw = zstd_decode(self.stream(i)?).map_err(|e| match e {
205            UrnaError::MalformedSectionPayload { reason, .. } => malformed(reason),
206            other => other,
207        })?;
208        String::from_utf8(raw).map_err(|e| malformed(format!("txt_streams: invalid utf-8: {}", e)))
209    }
210}
211
212/// reconstruct the canonical (raw-encoding) `chunks_canonical` payload from
213/// a full `txt_streams` payload (including the leading kind byte). the
214/// output is byte-identical to `sections::encode_chunks_canonical`, so
215/// `content_hash` is preserved. dispatched by `encoding::decode_payload`.
216pub fn decode(bytes: &[u8]) -> crate::Result<Vec<u8>> {
217    let parsed = TxtStreams::parse(bytes)?;
218    let mut bodies: Vec<Vec<u8>> = Vec::with_capacity(parsed.count);
219    for i in 0..parsed.count {
220        let raw = zstd_decode(parsed.stream(i)?).map_err(|e| match e {
221            UrnaError::MalformedSectionPayload { reason, .. } => malformed(reason),
222            other => other,
223        })?;
224        // validate utf-8 (the raw canonical encoder only ever wrote utf-8;
225        // a tampered stream must be rejected, not silently passed through).
226        std::str::from_utf8(&raw)
227            .map_err(|e| malformed(format!("txt_streams: invalid utf-8: {}", e)))?;
228        bodies.push(raw);
229    }
230    build_canonical(parsed.count, &bodies)
231}
232
233#[cfg(test)]
234mod tests {
235    use super::*;
236    use crate::sections::encode_chunks_canonical;
237
238    fn texts(items: &[&str]) -> Vec<String> {
239        items.iter().map(|s| s.to_string()).collect()
240    }
241
242    fn assert_byte_identical(items: &[&str]) {
243        let t = texts(items);
244        let packed = encode_txt_streams(&t).unwrap();
245        assert_eq!(packed[0], TXT_STREAMS_V1);
246        let raw = encode_chunks_canonical(&t).unwrap();
247        assert_eq!(decode(&packed).unwrap(), raw, "decode must rebuild raw");
248    }
249
250    #[test]
251    #[cfg_attr(miri, ignore)] // zstd is c code, miri cannot call it
252    fn byte_identical_across_corpora() {
253        assert_byte_identical(&[]);
254        assert_byte_identical(&["only one"]);
255        assert_byte_identical(&["primeiro", "segundo", "terceiro"]);
256        // multibyte utf-8 (pt-br accents) must round-trip exactly.
257        assert_byte_identical(&["coração", "informação", "açaí é ótimo", ""]);
258    }
259
260    #[test]
261    #[cfg_attr(miri, ignore)] // zstd is c code, miri cannot call it
262    fn o1_seek_returns_the_right_stream() {
263        let t = texts(&["alpha", "beta", "gama", "delta"]);
264        let packed = encode_txt_streams(&t).unwrap();
265        let parsed = TxtStreams::parse(&packed).unwrap();
266        assert_eq!(parsed.len(), 4);
267        for (i, s) in t.iter().enumerate() {
268            assert_eq!(&parsed.text(i).unwrap(), s, "text({}) mismatch", i);
269        }
270        assert!(parsed.text(4).is_err(), "oob index must error");
271    }
272
273    #[test]
274    #[cfg_attr(miri, ignore)] // zstd is c code, miri cannot call it
275    fn determinism_two_encodes_byte_identical() {
276        let t = texts(&["a", "bb", "ccc", "coração"]);
277        assert_eq!(
278            encode_txt_streams(&t).unwrap(),
279            encode_txt_streams(&t).unwrap()
280        );
281    }
282}