Skip to main content

vole_document/materialize/
seek.rs

1//! Seek-based partial materialization (Phase 8.2).
2//!
3//! [`materialize_observation_seeked`] serves the same narrow byte range as
4//! [`crate::materialize::observation::materialize_observation`], but from a
5//! `Read + Seek` source and by reading **only the records the query needs**. It
6//! does this by reading the optional `DIRECTORY` record (written as the first
7//! record at offset 64), which locates every record class, and then seeking
8//! directly to the `GRAPH`, `OBSERVATION_INDEX`, `INTEGRITY`, and the referenced
9//! `OBJECT`/`ENTROPY_CHANNEL`/`MODEL` records.
10//!
11//! ## Decline, never guess
12//!
13//! A descriptor whose header does not advertise the seek feature, or whose first
14//! record is not a `DIRECTORY`, is declined with
15//! [`crate::ErrorClass::UnsupportedFeature`]; the reader never silently falls
16//! back to reading the whole file. Op selection, lazy channel decoding, and the
17//! evaluated slice are shared verbatim with the Phase-7 in-memory path
18//! ([`crate::materialize::observation::select_ops`] / `serve_selection`), so the
19//! two readers cannot diverge.
20//!
21//! ## The directory is advisory, never authority
22//!
23//! Every locator is cross-checked against the record it points at (tag byte and
24//! payload length), every record's own CRC32C is verified when it is read, the
25//! directory's internal geometry is validated by
26//! [`SeekDirectory::validate_structural`], and the `ObservationIndex` is
27//! re-derived against the directory-derived object/channel lengths with
28//! [`ObservationIndex::validate`]. A lying directory is rejected, never trusted.
29//!
30//! ## Integrity is honest
31//!
32//! A partial read cannot recompute the whole-source SHA-256, so a served slice is
33//! an *observation* consistent with the descriptor's own validated
34//! program/index/directory — **not** a verified archival read. The returned
35//! [`ObservationStats`] carries `integrity_verified == false` and a real
36//! `bytes_read` measured by the internal [`CountingReader`]. Only
37//! `materialize`/`decode`/`verify` check `INTEGRITY` and are the archival
38//! authority.
39
40use std::io::{self, Read, Seek, SeekFrom};
41
42use crate::container::directory::{
43    DirectoryEntry, SECTION_CHANNEL_LENGTHS, SECTION_LOCATORS, SeekDirectory,
44};
45use crate::container::header::{FEATURE_SEEK_DIRECTORY, HEADER_LEN, Header};
46use crate::container::observation::ObservationIndex;
47use crate::container::record::{FLAG_OPTIONAL, RECORD_OVERHEAD, Record, RecordTag, read_record_at};
48use crate::dra::Program;
49use crate::entropy::codec::EntropyChannelDescriptor;
50use crate::entropy::model::EntropyModel;
51use crate::entropy::{CODER_ORDER0_BYTE_RANS, CODER_VERSION_1};
52use crate::error::{Error, Result};
53use crate::limits::Limits;
54use crate::materialize::observation::{
55    ObservationReport, ObservationSelector, ObservationStats, resolve_selector, select_ops,
56    selection_references, serve_selection,
57};
58
59/// A `Read + Seek` wrapper that counts the bytes actually returned by `read` and
60/// the number of `read`/`seek` calls.
61///
62/// This is the primary instrument for the Phase-8 bytes-read claim: it is exact,
63/// deterministic, attributes bytes to VOLE (not to a loader or pipe read-ahead),
64/// and is immune to page-cache effects. The library wraps one internally in
65/// [`materialize_observation_seeked`] and reports `bytes_read`; it is exposed so
66/// callers and tests can measure I/O directly.
67#[derive(Debug)]
68pub struct CountingReader<R> {
69    inner: R,
70    bytes_read: u64,
71    read_calls: u32,
72    seeks: u32,
73}
74
75impl<R> CountingReader<R> {
76    /// Wrap `inner`, starting all counters at zero.
77    pub fn new(inner: R) -> Self {
78        CountingReader {
79            inner,
80            bytes_read: 0,
81            read_calls: 0,
82            seeks: 0,
83        }
84    }
85
86    /// Total payload bytes returned by `read` so far.
87    pub fn bytes_read(&self) -> u64 {
88        self.bytes_read
89    }
90
91    /// Number of `read` calls (each may return fewer bytes than requested).
92    pub fn read_calls(&self) -> u32 {
93        self.read_calls
94    }
95
96    /// Number of `seek` calls.
97    pub fn seeks(&self) -> u32 {
98        self.seeks
99    }
100
101    /// Unwrap the inner reader.
102    pub fn into_inner(self) -> R {
103        self.inner
104    }
105}
106
107impl<R: Read> Read for CountingReader<R> {
108    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
109        let n = self.inner.read(buf)?;
110        self.bytes_read += n as u64;
111        self.read_calls += 1;
112        Ok(n)
113    }
114}
115
116impl<R: Seek> Seek for CountingReader<R> {
117    fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
118        self.seeks += 1;
119        self.inner.seek(pos)
120    }
121}
122
123/// Serve one observation from a seekable `.voldoc` source by reading only the
124/// records the query needs.
125///
126/// The returned bytes equal `materialize(parsed)[a..b]` for the resolved range.
127/// A source without a seek directory is declined (never silently fully read).
128/// See the module documentation for the validation and integrity rules.
129pub fn materialize_observation_seeked<R: Read + Seek>(
130    reader: R,
131    selector: ObservationSelector,
132    limits: Limits,
133) -> Result<ObservationReport> {
134    let mut reader = CountingReader::new(reader);
135    let file_len = reader
136        .seek(SeekFrom::End(0))
137        .map_err(|e| Error::io(format!("seek to end failed: {e}")))?;
138
139    // (a) Header, then the fixed-offset DIRECTORY record.
140    let header = read_header(&mut reader)?;
141    if header.optional_features & FEATURE_SEEK_DIRECTORY == 0 {
142        return Err(Error::unsupported_feature(
143            "descriptor has no seek directory (the header does not declare the seek feature)",
144        ));
145    }
146    let dir_rec = read_record_at(&mut reader, HEADER_LEN as u64, limits)?;
147    if dir_rec.tag != RecordTag::Directory as u8 {
148        return Err(Error::unsupported_feature(
149            "descriptor has no seek directory record at offset 64",
150        ));
151    }
152    if dir_rec.flags & FLAG_OPTIONAL == 0 {
153        return Err(Error::invalid_container(
154            "DIRECTORY record must carry FLAG_OPTIONAL",
155        ));
156    }
157    let dir = SeekDirectory::decode(&dir_rec.payload, limits)?;
158    dir.validate_structural(file_len, limits)?;
159    if dir.section_flags & SECTION_LOCATORS == 0 {
160        return Err(Error::unsupported_feature(
161            "seek directory omits the locator section",
162        ));
163    }
164    if dir.entries[0].payload_len != dir_rec.payload.len() as u32 {
165        return Err(Error::invalid_container(
166            "seek directory locator 0 payload length disagrees with the record framing",
167        ));
168    }
169
170    let object_entries = class_entries(&dir, RecordTag::Object)?;
171    let channel_entries = class_entries(&dir, RecordTag::EntropyChannel)?;
172    let model_entries = class_entries(&dir, RecordTag::Model)?;
173    if !class_entries(&dir, RecordTag::ExternalRef)?.is_empty() {
174        return Err(Error::unsupported_feature(
175            "seek-based partial read cannot resolve external objects",
176        ));
177    }
178    if !channel_entries.is_empty() && dir.section_flags & SECTION_CHANNEL_LENGTHS == 0 {
179        return Err(Error::unsupported_feature(
180            "seek directory omits the channel-lengths section",
181        ));
182    }
183    if dir.channel_lengths.len() != channel_entries.len() {
184        return Err(Error::invalid_container(
185            "seek directory channel-length table disagrees with the channel locators",
186        ));
187    }
188    let object_lens: Vec<u64> = object_entries
189        .iter()
190        .map(|e| u64::from(e.payload_len))
191        .collect();
192    let channel_lens: Vec<u64> = dir.channel_lengths.clone();
193
194    // (b) GRAPH and OBSERVATION_INDEX (the reader learns op->output and deps
195    // only from these), plus INTEGRITY for the declared source length.
196    let graph_site = class_entries(&dir, RecordTag::Graph)?
197        .first()
198        .ok_or_else(|| {
199            Error::unsupported_feature("seek directory does not locate a GRAPH record")
200        })?;
201    let index_site = class_entries(&dir, RecordTag::ObservationIndex)?
202        .first()
203        .ok_or_else(|| {
204            Error::unsupported_feature("seek directory does not locate an OBSERVATION_INDEX record")
205        })?;
206    let integrity_site = class_entries(&dir, RecordTag::Integrity)?
207        .first()
208        .ok_or_else(|| {
209            Error::unsupported_feature("seek directory does not locate an INTEGRITY record")
210        })?;
211
212    let graph_rec = read_checked(&mut reader, graph_site, limits)?;
213    let program = Program::decode(&graph_rec.payload, limits)?;
214    let index_rec = read_checked(&mut reader, index_site, limits)?;
215    let index = ObservationIndex::decode(&index_rec.payload, limits)?;
216    let integrity_rec = read_checked(&mut reader, integrity_site, limits)?;
217    if integrity_rec.payload.len() != 40 {
218        return Err(Error::invalid_container(
219            "INTEGRITY payload must be 40 bytes",
220        ));
221    }
222    let declared_len = u64::from_le_bytes([
223        integrity_rec.payload[32],
224        integrity_rec.payload[33],
225        integrity_rec.payload[34],
226        integrity_rec.payload[35],
227        integrity_rec.payload[36],
228        integrity_rec.payload[37],
229        integrity_rec.payload[38],
230        integrity_rec.payload[39],
231    ]);
232    if declared_len != header.declared_source_len {
233        return Err(Error::integrity_mismatch(format!(
234            "INTEGRITY length {declared_len} disagrees with header {}",
235            header.declared_source_len
236        )));
237    }
238
239    // The decisive cross-check: the index, re-derived over the *directory-derived*
240    // lengths, must reproduce the program's op table and the CRC-framed index.
241    index.validate(&program, &object_lens, &channel_lens, limits)?;
242
243    // (c) Resolve the selector and compute the minimal op set/prefix.
244    let (a, b) = resolve_selector(&index, selector, declared_len)?;
245    let window = select_ops(&program, &object_lens, &channel_lens, a, b, limits)?;
246    let (objects_used, channels_used) =
247        selection_references(&window.ops, object_entries.len(), channel_entries.len());
248
249    // (d) Read ONLY the referenced OBJECT records.
250    let mut objects: Vec<Vec<u8>> = vec![Vec::new(); object_entries.len()];
251    for (id, used) in objects_used.iter().enumerate() {
252        if *used {
253            objects[id] = read_checked(&mut reader, &object_entries[id], limits)?.payload;
254        }
255    }
256
257    // (d) Read ONLY the referenced ENTROPY_CHANNEL records, then the MODEL records
258    // they name.
259    let mut channels: Vec<EntropyChannelDescriptor> =
260        vec![placeholder_channel(); channel_entries.len()];
261    let mut models_needed = vec![false; model_entries.len()];
262    for (id, used) in channels_used.iter().enumerate() {
263        if *used {
264            let rec = read_checked(&mut reader, &channel_entries[id], limits)?;
265            let channel = EntropyChannelDescriptor::decode(&rec.payload, limits)?;
266            if channel.decoded_length != channel_lens[id] {
267                return Err(Error::invalid_container(format!(
268                    "seek directory channel-length {id} disagrees with the channel record"
269                )));
270            }
271            if channel.model_id as usize >= model_entries.len() {
272                return Err(Error::invalid_model(format!(
273                    "entropy channel {id} references missing model {}",
274                    channel.model_id
275                )));
276            }
277            models_needed[channel.model_id as usize] = true;
278            channels[id] = channel;
279        }
280    }
281    let mut models: Vec<EntropyModel> = vec![placeholder_model(); model_entries.len()];
282    for (id, used) in models_needed.iter().enumerate() {
283        if *used {
284            let rec = read_checked(&mut reader, &model_entries[id], limits)?;
285            models[id] = EntropyModel::decode(&rec.payload)?;
286        }
287    }
288
289    // (e) Evaluate the selected ops and slice `[a, b)` exactly as Phase 7 does.
290    let served = serve_selection(
291        &objects,
292        &channels,
293        &models,
294        window,
295        &objects_used,
296        &channels_used,
297        a,
298        b,
299        limits,
300    )?;
301
302    // `descriptor_bytes_traversed` stays comparable with the Phase-7 path: graph
303    // payload + index payload (+ framing) + referenced object/channel payloads.
304    let descriptor_bytes_traversed = graph_rec.payload.len() as u64
305        + index_rec.payload.len() as u64
306        + RECORD_OVERHEAD as u64
307        + served.referenced_object_bytes
308        + served.referenced_channel_bytes;
309
310    let stats = ObservationStats {
311        ops_evaluated: served.ops_evaluated,
312        ops_total: served.ops_total,
313        objects_fetched: served.objects_fetched,
314        objects_total: object_entries.len(),
315        channels_decoded: served.channels_decoded,
316        channels_total: channel_entries.len(),
317        entropy_bytes_decoded: served.entropy_bytes_decoded,
318        descriptor_bytes_traversed,
319        output_bytes: served.bytes.len() as u64,
320        // (f) Real I/O, distinct from the Phase-7 CPU-side approximation.
321        bytes_read: reader.bytes_read(),
322        integrity_verified: false,
323    };
324
325    Ok(ObservationReport {
326        range: (a, b),
327        bytes: served.bytes,
328        stats,
329    })
330}
331
332/// Read and decode the fixed 64-byte header from the start of the source.
333fn read_header<R: Read + Seek>(reader: &mut R) -> Result<Header> {
334    reader
335        .seek(SeekFrom::Start(0))
336        .map_err(|e| Error::io(format!("seek to header failed: {e}")))?;
337    let mut buf = [0u8; HEADER_LEN];
338    reader
339        .read_exact(&mut buf)
340        .map_err(|_| Error::invalid_container("truncated header"))?;
341    Header::decode(&buf)
342}
343
344/// Read the record at `site` and require its framing to match the locator.
345fn read_checked<R: Read + Seek>(
346    reader: &mut R,
347    site: &DirectoryEntry,
348    limits: Limits,
349) -> Result<Record> {
350    let rec = read_record_at(reader, site.offset, limits)?;
351    if rec.tag != site.tag || rec.payload.len() as u64 != u64::from(site.payload_len) {
352        return Err(Error::invalid_container(format!(
353            "DIRECTORY locator for tag {:#04x} disagrees with the record framing at offset {}",
354            site.tag, site.offset
355        )));
356    }
357    Ok(rec)
358}
359
360/// The locators of `tag`, using the directory's `CLASS_INDEX` and re-checking it
361/// against a linear scan of the locator table (defence in depth;
362/// [`SeekDirectory::validate_structural`] has already required agreement).
363fn class_entries(dir: &SeekDirectory, tag: RecordTag) -> Result<&[DirectoryEntry]> {
364    let (first, count) = match dir.classes.iter().find(|c| c.tag == tag as u8) {
365        Some(c) => (c.first as usize, c.count as usize),
366        None => (0, 0),
367    };
368    let scan_first = dir.entries.iter().position(|e| e.tag == tag as u8);
369    let scan_count = dir.entries.iter().filter(|e| e.tag == tag as u8).count();
370    let expected_first = if count == 0 { None } else { Some(first) };
371    if scan_first != expected_first || scan_count != count {
372        return Err(Error::invalid_container(
373            "seek directory class index disagrees with the locator table",
374        ));
375    }
376    if count == 0 {
377        return Ok(&[]);
378    }
379    let end = first
380        .checked_add(count)
381        .ok_or_else(|| Error::invalid_container("seek directory class range overflow"))?;
382    dir.entries
383        .get(first..end)
384        .ok_or_else(|| Error::invalid_container("seek directory class range out of bounds"))
385}
386
387/// A never-dereferenced placeholder occupying an unreferenced channel slot, so
388/// index positions in the partial vectors stay aligned with the descriptor's
389/// own tables.
390fn placeholder_channel() -> EntropyChannelDescriptor {
391    EntropyChannelDescriptor {
392        coder: CODER_ORDER0_BYTE_RANS,
393        coder_version: CODER_VERSION_1,
394        scale_bits: 0,
395        lane_count: 1,
396        model_id: 0,
397        symbol_count: 0,
398        decoded_length: 0,
399        initial_state: 0,
400        payload: Vec::new(),
401    }
402}
403
404/// A never-dereferenced placeholder occupying an unreferenced model slot.
405fn placeholder_model() -> EntropyModel {
406    EntropyModel {
407        scale_bits: 0,
408        frequencies: Vec::new(),
409    }
410}