vole-document 0.1.0-alpha.19

Persistent procedural document runtime: byte-exact reconstruction plus a content-addressed procedural seed DAG, queryable observations with provenance, and selective late materialization.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
//! Seek-based partial materialization (Phase 8.2).
//!
//! [`materialize_observation_seeked`] serves the same narrow byte range as
//! [`crate::materialize::observation::materialize_observation`], but from a
//! `Read + Seek` source and by reading **only the records the query needs**. It
//! does this by reading the optional `DIRECTORY` record (written as the first
//! record at offset 64), which locates every record class, and then seeking
//! directly to the `GRAPH`, `OBSERVATION_INDEX`, `INTEGRITY`, and the referenced
//! `OBJECT`/`ENTROPY_CHANNEL`/`MODEL` records.
//!
//! ## Decline, never guess
//!
//! A descriptor whose header does not advertise the seek feature, or whose first
//! record is not a `DIRECTORY`, is declined with
//! [`crate::ErrorClass::UnsupportedFeature`]; the reader never silently falls
//! back to reading the whole file. Op selection, lazy channel decoding, and the
//! evaluated slice are shared verbatim with the Phase-7 in-memory path
//! ([`crate::materialize::observation::select_ops`] / `serve_selection`), so the
//! two readers cannot diverge.
//!
//! ## The directory is advisory, never authority
//!
//! Every locator is cross-checked against the record it points at (tag byte and
//! payload length), every record's own CRC32C is verified when it is read, the
//! directory's internal geometry is validated by
//! [`SeekDirectory::validate_structural`], and the `ObservationIndex` is
//! re-derived against the directory-derived object/channel lengths with
//! [`ObservationIndex::validate`]. A lying directory is rejected, never trusted.
//!
//! ## Integrity is honest
//!
//! A partial read cannot recompute the whole-source SHA-256, so a served slice is
//! an *observation* consistent with the descriptor's own validated
//! program/index/directory — **not** a verified archival read. The returned
//! [`ObservationStats`] carries `integrity_verified == false` and a real
//! `bytes_read` measured by the internal [`CountingReader`]. Only
//! `materialize`/`decode`/`verify` check `INTEGRITY` and are the archival
//! authority.

use std::io::{self, Read, Seek, SeekFrom};

use crate::container::checkpoint::CheckpointTable;
use crate::container::directory::{
    DirectoryEntry, SECTION_CHANNEL_LENGTHS, SECTION_LOCATORS, SeekDirectory,
};
use crate::container::header::{FEATURE_SEEK_DIRECTORY, HEADER_LEN, Header};
use crate::container::observation::ObservationIndex;
use crate::container::record::{FLAG_OPTIONAL, RECORD_OVERHEAD, Record, RecordTag, read_record_at};
use crate::dra::Program;
use crate::entropy::codec::EntropyChannelDescriptor;
use crate::entropy::model::EntropyModel;
use crate::entropy::{CODER_ORDER0_BYTE_RANS, CODER_VERSION_1};
use crate::error::{Error, Result};
use crate::limits::Limits;
use crate::materialize::observation::{
    ObservationReport, ObservationSelector, ObservationStats, OpWindow, resolve_byte_range,
    resolve_selector, select_ops, select_ops_from_lengths, selection_references, serve_selection,
};

/// A `Read + Seek` wrapper that counts the bytes actually returned by `read` and
/// the number of `read`/`seek` calls.
///
/// This is the primary instrument for the Phase-8 bytes-read claim: it is exact,
/// deterministic, attributes bytes to VOLE (not to a loader or pipe read-ahead),
/// and is immune to page-cache effects. The library wraps one internally in
/// [`materialize_observation_seeked`] and reports `bytes_read`; it is exposed so
/// callers and tests can measure I/O directly.
#[derive(Debug)]
pub struct CountingReader<R> {
    inner: R,
    bytes_read: u64,
    read_calls: u32,
    seeks: u32,
}

impl<R> CountingReader<R> {
    /// Wrap `inner`, starting all counters at zero.
    pub fn new(inner: R) -> Self {
        CountingReader {
            inner,
            bytes_read: 0,
            read_calls: 0,
            seeks: 0,
        }
    }

    /// Total payload bytes returned by `read` so far.
    pub fn bytes_read(&self) -> u64 {
        self.bytes_read
    }

    /// Number of `read` calls (each may return fewer bytes than requested).
    pub fn read_calls(&self) -> u32 {
        self.read_calls
    }

    /// Number of `seek` calls.
    pub fn seeks(&self) -> u32 {
        self.seeks
    }

    /// Unwrap the inner reader.
    pub fn into_inner(self) -> R {
        self.inner
    }
}

impl<R: Read> Read for CountingReader<R> {
    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
        let n = self.inner.read(buf)?;
        self.bytes_read += n as u64;
        self.read_calls += 1;
        Ok(n)
    }
}

impl<R: Seek> Seek for CountingReader<R> {
    fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
        self.seeks += 1;
        self.inner.seek(pos)
    }
}

/// The selection a lane yields: the resolved output range `[a, b)`, the op
/// window to evaluate, and the referenced-object/channel flags.
type Selection = (u64, u64, OpWindow, Vec<bool>, Vec<bool>);

/// Serve one observation from a seekable `.voldoc` source by reading only the
/// records the query needs.
///
/// The returned bytes equal `materialize(parsed)[a..b]` for the resolved range.
/// A source without a seek directory is declined (never silently fully read).
/// See the module documentation for the validation and integrity rules.
pub fn materialize_observation_seeked<R: Read + Seek>(
    reader: R,
    selector: ObservationSelector,
    limits: Limits,
) -> Result<ObservationReport> {
    let mut reader = CountingReader::new(reader);
    let file_len = reader
        .seek(SeekFrom::End(0))
        .map_err(|e| Error::io(format!("seek to end failed: {e}")))?;

    // (a) Header, then the fixed-offset DIRECTORY record.
    let header = read_header(&mut reader)?;
    if header.optional_features & FEATURE_SEEK_DIRECTORY == 0 {
        return Err(Error::unsupported_feature(
            "descriptor has no seek directory (the header does not declare the seek feature)",
        ));
    }
    let dir_rec = read_record_at(&mut reader, HEADER_LEN as u64, limits)?;
    if dir_rec.tag != RecordTag::Directory as u8 {
        return Err(Error::unsupported_feature(
            "descriptor has no seek directory record at offset 64",
        ));
    }
    if dir_rec.flags & FLAG_OPTIONAL == 0 {
        return Err(Error::invalid_container(
            "DIRECTORY record must carry FLAG_OPTIONAL",
        ));
    }
    let dir = SeekDirectory::decode(&dir_rec.payload, limits)?;
    dir.validate_structural(file_len, limits)?;
    if dir.section_flags & SECTION_LOCATORS == 0 {
        return Err(Error::unsupported_feature(
            "seek directory omits the locator section",
        ));
    }
    if dir.entries[0].payload_len != dir_rec.payload.len() as u32 {
        return Err(Error::invalid_container(
            "seek directory locator 0 payload length disagrees with the record framing",
        ));
    }

    let object_entries = class_entries(&dir, RecordTag::Object)?;
    let channel_entries = class_entries(&dir, RecordTag::EntropyChannel)?;
    let model_entries = class_entries(&dir, RecordTag::Model)?;
    if !class_entries(&dir, RecordTag::ExternalRef)?.is_empty() {
        return Err(Error::unsupported_feature(
            "seek-based partial read cannot resolve external objects",
        ));
    }
    if !channel_entries.is_empty() && dir.section_flags & SECTION_CHANNEL_LENGTHS == 0 {
        return Err(Error::unsupported_feature(
            "seek directory omits the channel-lengths section",
        ));
    }
    if dir.channel_lengths.len() != channel_entries.len() {
        return Err(Error::invalid_container(
            "seek directory channel-length table disagrees with the channel locators",
        ));
    }
    let object_lens: Vec<u64> = object_entries
        .iter()
        .map(|e| u64::from(e.payload_len))
        .collect();
    let channel_lens: Vec<u64> = dir.channel_lengths.clone();

    // (b) GRAPH and INTEGRITY are needed by every lane: the program is the
    // authority for op selection and for validating any checkpoint, and INTEGRITY
    // supplies the declared source length.
    let graph_site = class_entries(&dir, RecordTag::Graph)?
        .first()
        .ok_or_else(|| {
            Error::unsupported_feature("seek directory does not locate a GRAPH record")
        })?;
    let integrity_site = class_entries(&dir, RecordTag::Integrity)?
        .first()
        .ok_or_else(|| {
            Error::unsupported_feature("seek directory does not locate an INTEGRITY record")
        })?;
    let checkpoint_sites = class_entries(&dir, RecordTag::Checkpoint)?;

    let graph_rec = read_checked(&mut reader, graph_site, limits)?;
    let program = Program::decode(&graph_rec.payload, limits)?;
    let integrity_rec = read_checked(&mut reader, integrity_site, limits)?;
    if integrity_rec.payload.len() != 40 {
        return Err(Error::invalid_container(
            "INTEGRITY payload must be 40 bytes",
        ));
    }
    let declared_len = u64::from_le_bytes([
        integrity_rec.payload[32],
        integrity_rec.payload[33],
        integrity_rec.payload[34],
        integrity_rec.payload[35],
        integrity_rec.payload[36],
        integrity_rec.payload[37],
        integrity_rec.payload[38],
        integrity_rec.payload[39],
    ]);
    if declared_len != header.declared_source_len {
        return Err(Error::integrity_mismatch(format!(
            "INTEGRITY length {declared_len} disagrees with header {}",
            header.declared_source_len
        )));
    }

    // (c) Selection. Two lanes produce the same `(range, op window, dependency
    // set)`:
    //   * the Phase-8 index lane reads and re-derives the OBSERVATION_INDEX;
    //   * the checkpoint lane (Phase 13.4), usable only for a raw byte range
    //     (the one selector that needs no index), consumes a *validated*
    //     CHECKPOINT and reads no index record.
    // The checkpoint is advisory: a missing, corrupt, non-optional, or lying one
    // is ignored and the reader falls back to the index lane -- never to a
    // guess -- so a checkpoint can neither change a served byte nor deny service.
    let mut checkpoint_bytes: u64 = 0;
    let mut index_bytes: u64 = 0;
    let (a, b, window, objects_used, channels_used) = {
        let mut via_checkpoint: Option<Selection> = None;
        if matches!(selector, ObservationSelector::ByteRange { .. })
            && let Some(site) = checkpoint_sites.first()
            && let Ok(cp_rec) = read_checked(&mut reader, site, limits)
            && cp_rec.is_optional()
            && let Ok(cp) = CheckpointTable::decode(&cp_rec.payload, limits)
            && cp
                .validate(
                    &program,
                    &graph_rec.payload,
                    declared_len,
                    &object_lens,
                    &channel_lens,
                    limits,
                )
                .is_ok()
        {
            let (ra, rb) = resolve_byte_range(selector, declared_len)?;
            let w = select_ops_from_lengths(&program, &cp.lengths(), ra, rb)?;
            let (ou, cu) =
                selection_references(&w.ops, object_entries.len(), channel_entries.len());
            checkpoint_bytes = cp_rec.payload.len() as u64 + RECORD_OVERHEAD as u64;
            via_checkpoint = Some((ra, rb, w, ou, cu));
        }
        match via_checkpoint {
            Some((ra, rb, w, ou, cu)) => (ra, rb, w, ou, cu),
            None => {
                let index_site = class_entries(&dir, RecordTag::ObservationIndex)?
                    .first()
                    .ok_or_else(|| {
                        Error::unsupported_feature(
                            "seek directory does not locate an OBSERVATION_INDEX record",
                        )
                    })?;
                let index_rec = read_checked(&mut reader, index_site, limits)?;
                let index = ObservationIndex::decode(&index_rec.payload, limits)?;
                index.validate(&program, &object_lens, &channel_lens, limits)?;
                index_bytes = index_rec.payload.len() as u64 + RECORD_OVERHEAD as u64;
                let (ra, rb) = resolve_selector(&index, selector, declared_len)?;
                let w = select_ops(&program, &object_lens, &channel_lens, ra, rb, limits)?;
                let (ou, cu) =
                    selection_references(&w.ops, object_entries.len(), channel_entries.len());
                (ra, rb, w, ou, cu)
            }
        }
    };

    // (d) Read ONLY the referenced OBJECT records.
    let mut objects: Vec<Vec<u8>> = vec![Vec::new(); object_entries.len()];
    for (id, used) in objects_used.iter().enumerate() {
        if *used {
            objects[id] = read_checked(&mut reader, &object_entries[id], limits)?.payload;
        }
    }

    // (d) Read ONLY the referenced ENTROPY_CHANNEL records, then the MODEL records
    // they name.
    let mut channels: Vec<EntropyChannelDescriptor> =
        vec![placeholder_channel(); channel_entries.len()];
    let mut models_needed = vec![false; model_entries.len()];
    for (id, used) in channels_used.iter().enumerate() {
        if *used {
            let rec = read_checked(&mut reader, &channel_entries[id], limits)?;
            let channel = EntropyChannelDescriptor::decode(&rec.payload, limits)?;
            if channel.decoded_length != channel_lens[id] {
                return Err(Error::invalid_container(format!(
                    "seek directory channel-length {id} disagrees with the channel record"
                )));
            }
            if channel.model_id as usize >= model_entries.len() {
                return Err(Error::invalid_model(format!(
                    "entropy channel {id} references missing model {}",
                    channel.model_id
                )));
            }
            models_needed[channel.model_id as usize] = true;
            channels[id] = channel;
        }
    }
    let mut models: Vec<EntropyModel> = vec![placeholder_model(); model_entries.len()];
    for (id, used) in models_needed.iter().enumerate() {
        if *used {
            let rec = read_checked(&mut reader, &model_entries[id], limits)?;
            models[id] = EntropyModel::decode(&rec.payload)?;
        }
    }

    // (e) Evaluate the selected ops and slice `[a, b)` exactly as Phase 7 does.
    let served = serve_selection(
        &objects,
        &channels,
        &models,
        window,
        &objects_used,
        &channels_used,
        a,
        b,
        limits,
    )?;

    // `descriptor_bytes_traversed` stays comparable with the Phase-7 path: graph
    // payload + the selection lane's own record (index **or** checkpoint) + the
    // referenced object/channel payloads.
    let descriptor_bytes_traversed = graph_rec.payload.len() as u64
        + index_bytes
        + checkpoint_bytes
        + served.referenced_object_bytes
        + served.referenced_channel_bytes;

    let stats = ObservationStats {
        ops_evaluated: served.ops_evaluated,
        ops_total: served.ops_total,
        objects_fetched: served.objects_fetched,
        objects_total: object_entries.len(),
        channels_decoded: served.channels_decoded,
        channels_total: channel_entries.len(),
        entropy_bytes_decoded: served.entropy_bytes_decoded,
        descriptor_bytes_traversed,
        output_bytes: served.bytes.len() as u64,
        // (f) Real I/O, distinct from the Phase-7 CPU-side approximation.
        bytes_read: reader.bytes_read(),
        integrity_verified: false,
    };

    Ok(ObservationReport {
        range: (a, b),
        bytes: served.bytes,
        stats,
    })
}

/// Read and decode the fixed 64-byte header from the start of the source.
fn read_header<R: Read + Seek>(reader: &mut R) -> Result<Header> {
    reader
        .seek(SeekFrom::Start(0))
        .map_err(|e| Error::io(format!("seek to header failed: {e}")))?;
    let mut buf = [0u8; HEADER_LEN];
    reader
        .read_exact(&mut buf)
        .map_err(|_| Error::invalid_container("truncated header"))?;
    Header::decode(&buf)
}

/// Read the record at `site` and require its framing to match the locator.
fn read_checked<R: Read + Seek>(
    reader: &mut R,
    site: &DirectoryEntry,
    limits: Limits,
) -> Result<Record> {
    let rec = read_record_at(reader, site.offset, limits)?;
    if rec.tag != site.tag || rec.payload.len() as u64 != u64::from(site.payload_len) {
        return Err(Error::invalid_container(format!(
            "DIRECTORY locator for tag {:#04x} disagrees with the record framing at offset {}",
            site.tag, site.offset
        )));
    }
    Ok(rec)
}

/// The locators of `tag`, using the directory's `CLASS_INDEX` and re-checking it
/// against a linear scan of the locator table (defence in depth;
/// [`SeekDirectory::validate_structural`] has already required agreement).
fn class_entries(dir: &SeekDirectory, tag: RecordTag) -> Result<&[DirectoryEntry]> {
    let (first, count) = match dir.classes.iter().find(|c| c.tag == tag as u8) {
        Some(c) => (c.first as usize, c.count as usize),
        None => (0, 0),
    };
    let scan_first = dir.entries.iter().position(|e| e.tag == tag as u8);
    let scan_count = dir.entries.iter().filter(|e| e.tag == tag as u8).count();
    let expected_first = if count == 0 { None } else { Some(first) };
    if scan_first != expected_first || scan_count != count {
        return Err(Error::invalid_container(
            "seek directory class index disagrees with the locator table",
        ));
    }
    if count == 0 {
        return Ok(&[]);
    }
    let end = first
        .checked_add(count)
        .ok_or_else(|| Error::invalid_container("seek directory class range overflow"))?;
    dir.entries
        .get(first..end)
        .ok_or_else(|| Error::invalid_container("seek directory class range out of bounds"))
}

/// A never-dereferenced placeholder occupying an unreferenced channel slot, so
/// index positions in the partial vectors stay aligned with the descriptor's
/// own tables.
fn placeholder_channel() -> EntropyChannelDescriptor {
    EntropyChannelDescriptor {
        coder: CODER_ORDER0_BYTE_RANS,
        coder_version: CODER_VERSION_1,
        scale_bits: 0,
        lane_count: 1,
        model_id: 0,
        symbol_count: 0,
        decoded_length: 0,
        initial_state: 0,
        payload: Vec::new(),
    }
}

/// A never-dereferenced placeholder occupying an unreferenced model slot.
fn placeholder_model() -> EntropyModel {
    EntropyModel {
        scale_bits: 0,
        frequencies: Vec::new(),
    }
}