Skip to main content

kernel/
text.rs

1//! 2h: full-text search as key discipline (FACT-03).
2//!
3//! A newcomer's map: an inverted index answers "which documents contain
4//! this word?" -- it inverts documents into (word -> documents). Here the
5//! btree IS that index: the sorted term keys are the dictionary (prefix
6//! search = a range scan), and postings live in two shapes:
7//!
8//!  - the HEAD (segment 0): one row per (term, doc), written BLIND at
9//!    index time -- searchable the instant it lands, no build step;
10//!  - FOLDED segments (1..): one value per term with packed
11//!    (docid-delta, tf) varints -- compact, immutable, produced by
12//!    folding the head (2h fold, bulk-dock shape).
13//!
14//! BM25 in one breath: score = idf(term) * tf / (tf + k1*(1-b+b*|d|/avg)).
15//! idf rewards rare words, the denominator saturates repeated words and
16//! normalises by document length. Everything it needs lives in keyed rows:
17//! per-(field,doc) token counts (0x0D) and per-(field,seg) totals (0x0E).
18
19use crate::graph::Graph;
20use crate::keys;
21use crate::{Error, Result};
22use std::collections::{BTreeSet, BinaryHeap, HashMap, HashSet, VecDeque};
23
24pub const BM25_K1: f32 = 1.2;
25pub const BM25_B: f32 = 0.75;
26
27/// Lowercase alphanumeric tokens; everything else separates. Deliberately
28/// simple (no stemming, no stopwords) -- linguistic layers stack on later
29/// without changing the keyspace.
30pub fn tokenize(text: &str) -> Vec<String> {
31    let mut out = Vec::new();
32    let mut cur = String::new();
33    for ch in text.chars() {
34        if ch.is_alphanumeric() {
35            for lc in ch.to_lowercase() { cur.push(lc); }
36        } else if !cur.is_empty() {
37            out.push(std::mem::take(&mut cur));
38        }
39    }
40    if !cur.is_empty() { out.push(cur); }
41    out
42}
43
44/// Varint (LEB128) helpers for packed postings.
45pub fn write_varint(v: &mut Vec<u8>, mut x: u64) {
46    loop {
47        let b = (x & 0x7F) as u8;
48        x >>= 7;
49        if x == 0 { v.push(b); break; }
50        v.push(b | 0x80);
51    }
52}
53pub fn read_varint(b: &[u8], pos: &mut usize) -> Option<u64> {
54    let mut x = 0u64; let mut shift = 0;
55    loop {
56        let byte = *b.get(*pos)?;
57        *pos += 1;
58        x |= ((byte & 0x7F) as u64) << shift;
59        if byte & 0x80 == 0 { return Some(x); }
60        shift += 7;
61        if shift > 63 { return None; }
62    }
63}
64
65const SEG_MAGIC: &[u8] = b"TSEG2";
66const FIELD_MAGIC: &[u8] = b"TFM3";
67const NORM_MAGIC: &[u8] = b"TN3";
68const POSTING_BLOCK: usize = 128;
69const TEXT_BUILD_BLOCK_MAGIC: &[u8] = b"TB1";
70const TEXT_BUILD_OPEN_BUDGET: usize = 32 << 20;
71const TEXT_BUILD_OUTPUT_BUDGET: usize = 16 << 20;
72const MERGE_FANOUT: usize = 8;
73const BUILD_TERM_CACHE: usize = 4096;
74
75struct PostingSource<'a> {
76    scan: Option<crate::btree::RangeIter<'a>>,
77    prefix: Vec<u8>,
78    pending: VecDeque<(u64, u64, u64)>,
79    dead: HashSet<u64>,
80    head: bool,
81    blocks_read: u64,
82    postings_decoded: u64,
83}
84
85impl PostingSource<'_> {
86    fn next_live(&mut self) -> Result<Option<(u64, u64, u64)>> {
87        loop {
88            while let Some(posting) = self.pending.pop_front() {
89                if !self.dead.contains(&posting.0) { return Ok(Some(posting)); }
90            }
91            let Some(scan) = self.scan.as_mut() else { return Ok(None) };
92            let Some(item) = scan.next() else { self.scan = None; return Ok(None) };
93            let (key, value) = item?;
94            if !key.starts_with(&self.prefix) || key.len() != self.prefix.len() + 8 {
95                self.scan = None;
96                return Ok(None);
97            }
98            self.blocks_read += 1;
99            if self.head {
100                let doc = u64::from_be_bytes(key[self.prefix.len()..].try_into().unwrap());
101                let mut pos = 0;
102                let tf = required_varint(&value, &mut pos, "head posting frequency is truncated")?;
103                let dl = required_varint(&value, &mut pos, "head posting length is truncated")?;
104                if pos != value.len() || tf == 0 {
105                    return Err(corrupt("head posting has invalid trailing bytes or frequency"));
106                }
107                self.postings_decoded += 1;
108                if !self.dead.contains(&doc) { return Ok(Some((doc, tf, dl))); }
109            } else {
110                let posts = decode_postings(&value)?;
111                if posts.len() > POSTING_BLOCK {
112                    return Err(corrupt("text posting block exceeds its bound"));
113                }
114                self.postings_decoded += posts.len() as u64;
115                self.pending.extend(posts);
116            }
117        }
118    }
119}
120
121struct TextPostingCursor<'a> {
122    sources: Vec<PostingSource<'a>>,
123    heads: Vec<Option<(u64, u64, u64)>>,
124}
125
126impl TextPostingCursor<'_> {
127    fn next(&mut self) -> Result<Option<(u64, u64, u64)>> {
128        let Some(doc) = self.heads.iter().flatten().map(|posting| posting.0).min() else {
129            return Ok(None);
130        };
131        let mut chosen = None;
132        for i in 0..self.sources.len() {
133            if self.heads[i].is_some_and(|posting| posting.0 == doc) {
134                // Head is appended last and therefore wins the only legitimate
135                // cross-source duplicate: a replacement published beside its
136                // not-yet-retired old segment.
137                chosen = self.heads[i];
138                self.heads[i] = self.sources[i].next_live()?;
139            }
140        }
141        Ok(chosen)
142    }
143
144    fn counters(&self) -> (u64, u64) {
145        self.sources.iter().fold((0, 0), |(blocks, postings), source| {
146            (blocks + source.blocks_read, postings + source.postings_decoded)
147        })
148    }
149}
150
151struct RankedText { id: u64, score: f64 }
152impl PartialEq for RankedText {
153    fn eq(&self, other: &Self) -> bool { self.id == other.id && self.score == other.score }
154}
155impl Eq for RankedText {}
156impl Ord for RankedText {
157    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
158        // The worst retained hit is the max-heap root: lower score, then larger id.
159        other.score.total_cmp(&self.score).then(self.id.cmp(&other.id))
160    }
161}
162impl PartialOrd for RankedText {
163    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> { Some(self.cmp(other)) }
164}
165
166/// A fully sorted and logically verified text segment that has not touched the
167/// shared tree. It is Send because it owns only run-file descriptors and
168/// scalar manifests; D26 publication still requires a caller-owned Graph.
169pub struct PackedTextCandidate {
170    field: u64,
171    posting_runs: crate::bulk::SortedRuns,
172    norm_runs: crate::bulk::SortedRuns,
173    membership_runs: crate::bulk::SortedRuns,
174    dictionary_runs: crate::bulk::SortedRuns,
175    expected: SegMeta,
176    doc_count: u64,
177    total_tokens: u64,
178    posting_rows: u64,
179    emissions: u64,
180    partials: u64,
181    block_rows: u64,
182    block_min: Option<Vec<u8>>,
183    block_max: Option<Vec<u8>>,
184    norm_rows: u64,
185    norm_min: Option<Vec<u8>>,
186    norm_max: Option<Vec<u8>>,
187    member_rows: u64,
188    member_min: Option<Vec<u8>>,
189    member_max: Option<Vec<u8>>,
190    dict_rows: u64,
191    dict_min: Option<Vec<u8>>,
192    dict_max: Option<Vec<u8>>,
193    scan_stage: std::time::Duration,
194    source_finish: std::time::Duration,
195    posting_merge: std::time::Duration,
196    block_term_finish: std::time::Duration,
197    dictionary_stage: std::time::Duration,
198    dictionary_finish: std::time::Duration,
199    validate: std::time::Duration,
200    prepare_total: std::time::Duration,
201    posting_scratch: u64,
202    norm_scratch: u64,
203    membership_scratch: u64,
204    term_scratch: u64,
205    dictionary_scratch: u64,
206    partial_block_scratch: u64,
207    posting_input_runs: usize,
208    norm_input_runs: usize,
209    membership_input_runs: usize,
210    term_input_runs: usize,
211    dictionary_input_runs: usize,
212}
213
214#[derive(Debug, PartialEq, Eq)]
215struct PreparedPackedTextDoc {
216    doc: u64,
217    dl: u64,
218    terms: Vec<(String, u64)>,
219    norm_terms: Vec<(u64, u64)>,
220}
221
222fn prepare_packed_text_doc(doc: u64, text: &str) -> Result<PreparedPackedTextDoc> {
223    let tokens = tokenize(text);
224    let dl = tokens.len() as u64;
225    let mut tf = std::collections::BTreeMap::<String, u64>::new();
226    for term in tokens { *tf.entry(term).or_insert(0) += 1; }
227
228    let terms: Vec<_> = tf.into_iter().collect();
229    let mut norm_terms: Vec<_> = terms.iter()
230        .map(|(term, count)| (term_number(term.as_bytes()), *count))
231        .collect();
232    norm_terms.sort_unstable_by_key(|&(id, _)| id);
233    if norm_terms.windows(2).any(|pair| pair[0].0 == pair[1].0) {
234        return Err(corrupt("packed text document contains colliding term numbers"));
235    }
236    Ok(PreparedPackedTextDoc { doc, dl, terms, norm_terms })
237}
238
239fn text_prepare_worker_limit() -> usize {
240    std::env::var("SEKEJAP_INDEX_BUILD_WORKERS")
241        .ok()
242        .and_then(|value| value.parse::<usize>().ok())
243        .unwrap_or_else(|| std::thread::available_parallelism().map_or(1, usize::from))
244        .clamp(1, 8)
245}
246
247#[cfg(test)]
248mod parallel_prepare_tests {
249    use super::*;
250
251    #[test]
252    fn packed_document_transform_is_deterministic() {
253        let text = "Railway railway junction café";
254        assert_eq!(
255            prepare_packed_text_doc(17, text).unwrap(),
256            prepare_packed_text_doc(17, text).unwrap(),
257        );
258    }
259
260    #[test]
261    fn packed_document_transform_preserves_document_identity_and_length() {
262        let prepared = prepare_packed_text_doc(9_223_372_036_854_775_000, "one two two").unwrap();
263        assert_eq!(prepared.doc, 9_223_372_036_854_775_000);
264        assert_eq!(prepared.dl, 3);
265        assert_eq!(prepared.terms, vec![("one".into(), 1), ("two".into(), 2)]);
266    }
267
268    #[test]
269    fn packed_document_transform_canonicalizes_norm_term_order() {
270        let prepared = prepare_packed_text_doc(1, "zulu alpha beta alpha").unwrap();
271        assert!(prepared.norm_terms.windows(2).all(|pair| pair[0].0 < pair[1].0));
272        assert_eq!(prepared.norm_terms.iter().map(|(_, count)| count).sum::<u64>(), 4);
273    }
274}
275
276/// A fixed-size, direct-mapped accelerator used only while initially
277/// backfilling an index.  It never grows with the corpus; long words bypass it
278/// so its retained memory is below one MiB.
279pub struct TextBuildCache {
280    slots: Vec<Option<(String, u64)>>,
281}
282
283/// Disk-first accumulator for an initial index build.  A fixed 4 MiB table
284/// combines repeated `(raw term number, exact word)` records before a fixed
285/// 4 MiB external sort. Finishing resolves the vanishingly rare number
286/// collision exactly and writes dictionary rows in key order.
287pub struct TextBuildAccumulator {
288    sort: Option<crate::bulk::ExternalSort>,
289    scratch: std::path::PathBuf,
290    grouped: HashMap<Vec<u8>, u64>,
291    grouped_bytes: usize,
292    docs: u64,
293    tokens: u64,
294}
295
296impl TextBuildAccumulator {
297    pub fn new() -> Result<Self> {
298        use std::sync::atomic::{AtomicU64, Ordering};
299        static SEQ: AtomicU64 = AtomicU64::new(0);
300        let scratch = std::env::temp_dir().join(format!("text-build-{}-{}",
301            std::process::id(), SEQ.fetch_add(1, Ordering::Relaxed)));
302        Ok(Self { sort: None, scratch, grouped: HashMap::new(), grouped_bytes: 0,
303            docs: 0, tokens: 0 })
304    }
305
306    fn document(&mut self, total: u64) -> Result<()> {
307        self.docs = self.docs.checked_add(1).ok_or(Error::TooLarge)?;
308        self.tokens = self.tokens.checked_add(total).ok_or(Error::TooLarge)?;
309        Ok(())
310    }
311
312    fn term(&mut self, term: &str, _docid: u64) -> Result<u64> {
313        let id = term_number(term.as_bytes());
314        let mut key = Vec::with_capacity(term.len() + 8);
315        key.extend_from_slice(&id.to_be_bytes()); key.extend_from_slice(term.as_bytes());
316        if let Some(count) = self.grouped.get_mut(&key) {
317            *count = count.checked_add(1).ok_or(Error::TooLarge)?;
318            return Ok(id);
319        }
320        if self.grouped_bytes.saturating_add(key.len() + 64) > (4 << 20) {
321            self.flush()?;
322        }
323        self.grouped_bytes = self.grouped_bytes.saturating_add(key.len() + 64);
324        self.grouped.insert(key, 1);
325        Ok(id)
326    }
327
328    fn flush(&mut self) -> Result<()> {
329        if self.sort.is_none() {
330            self.sort = Some(crate::bulk::ExternalSort::new(&self.scratch, 4 << 20)?);
331        }
332        let sort = self.sort.as_mut().unwrap();
333        for (key, count) in self.grouped.drain() {
334            sort.push(key, count.to_be_bytes().to_vec())?;
335        }
336        self.grouped_bytes = 0;
337        Ok(())
338    }
339}
340
341impl Default for TextBuildCache {
342    fn default() -> Self { Self { slots: (0..BUILD_TERM_CACHE).map(|_| None).collect() } }
343}
344
345impl TextBuildCache {
346    fn slot(term: &str) -> usize { term_number(term.as_bytes()) as usize & (BUILD_TERM_CACHE - 1) }
347    fn get(&self, term: &str) -> Option<u64> {
348        self.slots[Self::slot(term)].as_ref()
349            .and_then(|(stored, id)| (stored == term).then_some(*id))
350    }
351    fn insert(&mut self, term: &str, id: u64) {
352        if term.len() <= 128 { self.slots[Self::slot(term)] = Some((term.to_owned(), id)); }
353    }
354}
355
356fn corrupt(why: &'static str) -> Error { Error::Corrupt { page_no: 0, why } }
357
358fn required_varint(v: &[u8], pos: &mut usize, why: &'static str) -> Result<u64> {
359    read_varint(v, pos).ok_or_else(|| corrupt(why))
360}
361
362/// Per-(field, segment) manifest.  The logical hashes are computed from the
363/// source postings before the candidate is written and recomputed by the
364/// independently reopened verifier.  Counts catch truncation; the two
365/// differently seeded hashes catch changed/reordered logical records.
366#[derive(Clone, Debug, PartialEq, Eq)]
367pub struct SegMeta {
368    pub doc_count: u64,
369    pub total_tokens: u64,
370    /// Token occurrences belonging to dead docs, captured from the norm
371    /// row at delete time -- so avg document length is identical before
372    /// and after a fold physically drops the dead (score-drift bug).
373    pub dead_tokens: u64,
374    pub dead: Vec<u64>,
375    pub level: u32,
376    pub term_rows: u64,
377    pub posting_count: u64,
378    pub logical_xor: u64,
379    pub logical_sum: u64,
380}
381impl SegMeta {
382    pub fn decode(v: &[u8]) -> Result<SegMeta> {
383        if !v.starts_with(SEG_MAGIC) {
384            // Read compatibility for pre-manifest stores.  New writes always
385            // use the exact format below; malformed mandatory legacy fields
386            // are refused rather than silently becoming zero.
387            let mut pos = 0usize;
388            let doc_count = required_varint(v, &mut pos, "text segment has no document count")?;
389            let total_tokens = required_varint(v, &mut pos, "text segment has no token total")?;
390            let dead_tokens = required_varint(v, &mut pos, "text segment has no dead-token total")?;
391            let mut dead = Vec::new();
392            let mut last = 0u64;
393            while pos < v.len() {
394                let delta = required_varint(v, &mut pos, "text segment dead list is truncated")?;
395                last = last.checked_add(delta).ok_or_else(|| corrupt("text segment dead id overflows"))?;
396                dead.push(last);
397            }
398            return Ok(SegMeta { doc_count, total_tokens, dead_tokens, dead, level: 0,
399                term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0 });
400        }
401        let mut pos = SEG_MAGIC.len();
402        let doc_count = required_varint(v, &mut pos, "text segment has no document count")?;
403        let total_tokens = required_varint(v, &mut pos, "text segment has no token total")?;
404        let dead_tokens = required_varint(v, &mut pos, "text segment has no dead-token total")?;
405        let level = required_varint(v, &mut pos, "text segment has no level")?;
406        let term_rows = required_varint(v, &mut pos, "text segment has no term-row count")?;
407        let posting_count = required_varint(v, &mut pos, "text segment has no posting count")?;
408        let logical_xor = required_varint(v, &mut pos, "text segment has no xor manifest")?;
409        let logical_sum = required_varint(v, &mut pos, "text segment has no sum manifest")?;
410        let dead_count = required_varint(v, &mut pos, "text segment has no dead count")?;
411        if level > u32::MAX as u64 || dead_count > usize::MAX as u64 {
412            return Err(corrupt("text segment metadata exceeds format bounds"));
413        }
414        let mut dead = Vec::with_capacity(dead_count as usize);
415        let mut last = 0u64;
416        for _ in 0..dead_count {
417            let delta = required_varint(v, &mut pos, "text segment dead list is truncated")?;
418            last = last.checked_add(delta).ok_or_else(|| corrupt("text segment dead id overflows"))?;
419            dead.push(last);
420        }
421        if pos != v.len() { return Err(corrupt("text segment metadata has trailing bytes")); }
422        Ok(SegMeta { doc_count, total_tokens, dead_tokens, dead, level: level as u32,
423            term_rows, posting_count, logical_xor, logical_sum })
424    }
425    pub fn encode(&self) -> Vec<u8> {
426        let mut v = SEG_MAGIC.to_vec();
427        write_varint(&mut v, self.doc_count);
428        write_varint(&mut v, self.total_tokens);
429        write_varint(&mut v, self.dead_tokens);
430        write_varint(&mut v, self.level as u64);
431        write_varint(&mut v, self.term_rows);
432        write_varint(&mut v, self.posting_count);
433        write_varint(&mut v, self.logical_xor);
434        write_varint(&mut v, self.logical_sum);
435        let mut sorted = self.dead.clone();
436        sorted.sort_unstable(); sorted.dedup();
437        write_varint(&mut v, sorted.len() as u64);
438        let mut last = 0u64;
439        for d in sorted { write_varint(&mut v, d - last); last = d; }
440        v
441    }
442}
443
444#[derive(Clone, Debug)]
445struct FieldMeta {
446    live_docs: u64,
447    total_tokens: u64,
448    next_seg: u32,
449    generation: u64,
450    active: Vec<(u32, u32)>, // (segment, level)
451    head_redirect: Option<u32>,
452    redirects: Vec<(u32, u32)>,
453    term_stats_ready: bool,
454}
455
456impl Default for FieldMeta {
457    fn default() -> Self {
458        Self { live_docs: 0, total_tokens: 0, next_seg: 1, generation: 0,
459            active: Vec::new(), head_redirect: None, redirects: Vec::new(), term_stats_ready: true }
460    }
461}
462
463impl FieldMeta {
464    fn decode(v: &[u8]) -> Result<Self> {
465        if !v.starts_with(FIELD_MAGIC) { return Err(corrupt("text field manifest has an unknown format")); }
466        let mut pos = FIELD_MAGIC.len();
467        let live_docs = required_varint(v, &mut pos, "text field manifest has no document count")?;
468        let total_tokens = required_varint(v, &mut pos, "text field manifest has no token total")?;
469        let next_seg = required_varint(v, &mut pos, "text field manifest has no next segment")?;
470        let generation = required_varint(v, &mut pos, "text field manifest has no generation")?;
471        let term_stats_ready = required_varint(v, &mut pos, "text field manifest has no term-stat state")?;
472        if term_stats_ready > 1 { return Err(corrupt("text field term-stat state is invalid")); }
473        let head = required_varint(v, &mut pos, "text field manifest has no head state")?;
474        let nactive = required_varint(v, &mut pos, "text field manifest has no active count")?;
475        if next_seg >= keys::TEXT_TERM_STATS_SEG as u64 || nactive > 1024 {
476            return Err(corrupt("text field manifest exceeds format bounds"));
477        }
478        let mut active = Vec::with_capacity(nactive as usize);
479        for _ in 0..nactive {
480            let seg = required_varint(v, &mut pos, "text field active segment is truncated")?;
481            let level = required_varint(v, &mut pos, "text field active level is truncated")?;
482            if seg == 0 || seg >= keys::TEXT_TERM_STATS_SEG as u64 || level > u32::MAX as u64 {
483                return Err(corrupt("text field active segment is out of bounds"));
484            }
485            active.push((seg as u32, level as u32));
486        }
487        let nr = required_varint(v, &mut pos, "text field manifest has no redirect count")?;
488        if nr > 1024 { return Err(corrupt("text field redirect count exceeds its bound")); }
489        let mut redirects = Vec::with_capacity(nr as usize);
490        for _ in 0..nr {
491            let old = required_varint(v, &mut pos, "text field redirect is truncated")?;
492            let new = required_varint(v, &mut pos, "text field redirect is truncated")?;
493            if old >= keys::TEXT_TERM_STATS_SEG as u64 || new >= keys::TEXT_TERM_STATS_SEG as u64 {
494                return Err(corrupt("text field redirect is out of bounds"));
495            }
496            redirects.push((old as u32, new as u32));
497        }
498        if pos != v.len() { return Err(corrupt("text field manifest has trailing bytes")); }
499        Ok(FieldMeta { live_docs, total_tokens, next_seg: next_seg as u32, generation,
500            active, head_redirect: if head == 0 { None } else { Some((head - 1) as u32) }, redirects,
501            term_stats_ready: term_stats_ready == 1 })
502    }
503
504    fn encode(&self) -> Vec<u8> {
505        let mut v = FIELD_MAGIC.to_vec();
506        write_varint(&mut v, self.live_docs);
507        write_varint(&mut v, self.total_tokens);
508        write_varint(&mut v, self.next_seg as u64);
509        write_varint(&mut v, self.generation);
510        write_varint(&mut v, self.term_stats_ready as u64);
511        write_varint(&mut v, self.head_redirect.map_or(0, |s| s as u64 + 1));
512        write_varint(&mut v, self.active.len() as u64);
513        for &(seg, level) in &self.active { write_varint(&mut v, seg as u64); write_varint(&mut v, level as u64); }
514        write_varint(&mut v, self.redirects.len() as u64);
515        for &(old, new) in &self.redirects { write_varint(&mut v, old as u64); write_varint(&mut v, new as u64); }
516        v
517    }
518
519    fn resolve_owner(&self, owner: u32) -> u32 {
520        if owner == 0 {
521            if let Some(seg) = self.head_redirect { return seg; }
522        }
523        self.redirects.iter().find_map(|&(old, new)| (old == owner).then_some(new)).unwrap_or(owner)
524    }
525}
526
527#[derive(Clone, Debug)]
528struct Norm {
529    total: u64,
530    owner: u32,
531    /// Sorted `(field-local term number, frequency)` pairs.  The word itself
532    /// is interned once per field, rather than repeated in every document.
533    terms: Vec<(u64, u64)>,
534}
535
536impl Norm {
537    fn decode(v: &[u8]) -> Result<Self> {
538        if !v.starts_with(NORM_MAGIC) {
539            let mut pos = 0;
540            let total = required_varint(v, &mut pos, "text norm has no token count")?;
541            let owner = required_varint(v, &mut pos, "text norm has no owner")?;
542            if pos != v.len() || owner >= keys::TEXT_TERM_STATS_SEG as u64 {
543                return Err(corrupt("legacy text norm is malformed"));
544            }
545            return Ok(Norm { total, owner: owner as u32, terms: Vec::new() });
546        }
547        let mut pos = NORM_MAGIC.len();
548        let total = required_varint(v, &mut pos, "text norm has no token count")?;
549        let owner = required_varint(v, &mut pos, "text norm has no owner")?;
550        let n = required_varint(v, &mut pos, "text norm has no term count")?;
551        if owner >= keys::TEXT_TERM_STATS_SEG as u64 || n > u32::MAX as u64 {
552            return Err(corrupt("text norm exceeds format bounds"));
553        }
554        let mut terms = Vec::with_capacity(n as usize);
555        for _ in 0..n {
556            let term = required_varint(v, &mut pos, "text norm term number is truncated")?;
557            let tf = required_varint(v, &mut pos, "text norm term frequency is truncated")?;
558            if term == 0 || tf == 0 || terms.last().is_some_and(|(previous, _)| *previous >= term) {
559                return Err(corrupt("text norm terms are invalid or out of order"));
560            }
561            terms.push((term, tf));
562        }
563        if pos != v.len() { return Err(corrupt("text norm has trailing bytes")); }
564        Ok(Norm { total, owner: owner as u32, terms })
565    }
566
567    fn encode(&self) -> Vec<u8> {
568        let mut v = NORM_MAGIC.to_vec();
569        write_varint(&mut v, self.total);
570        write_varint(&mut v, self.owner as u64);
571        write_varint(&mut v, self.terms.len() as u64);
572        for (term, tf) in &self.terms {
573            write_varint(&mut v, *term); write_varint(&mut v, *tf);
574        }
575        v
576    }
577}
578
579fn posting_hash(term: &[u8], docid: u64, tf: u64, dl: u64) -> u64 {
580    let mut v = Vec::with_capacity(term.len() + 24);
581    v.extend_from_slice(term); v.extend_from_slice(&docid.to_be_bytes());
582    v.extend_from_slice(&tf.to_be_bytes()); v.extend_from_slice(&dl.to_be_bytes());
583    let lo = crc32c::crc32c(&v) as u64;
584    v.push(0xA5);
585    lo | ((crc32c::crc32c(&v) as u64) << 32)
586}
587
588fn decode_postings(v: &[u8]) -> Result<Vec<(u64, u64, u64)>> {
589    let mut out = Vec::new();
590    let (mut pos, mut last) = (0usize, 0u64);
591    while pos < v.len() {
592        let delta = required_varint(v, &mut pos, "text posting doc delta is truncated")?;
593        let tf = required_varint(v, &mut pos, "text posting term frequency is truncated")?;
594        let dl = required_varint(v, &mut pos, "text posting document length is truncated")?;
595        if tf == 0 { return Err(corrupt("text posting has zero term frequency")); }
596        last = last.checked_add(delta).ok_or_else(|| corrupt("text posting document id overflows"))?;
597        if out.last().is_some_and(|(previous, _, _)| *previous >= last) {
598            return Err(corrupt("text postings are not strictly ordered"));
599        }
600        out.push((last, tf, dl));
601    }
602    Ok(out)
603}
604
605#[derive(Default)]
606struct TextOpenBlock {
607    rows: Vec<(u64, u64, u64)>,
608    last_seen: u64,
609}
610
611impl TextOpenBlock {
612    fn new(doc: u64, tf: u64, dl: u64) -> Self {
613        let mut rows = Vec::with_capacity(POSTING_BLOCK);
614        rows.push((doc, tf, dl));
615        Self { rows, last_seen: doc }
616    }
617}
618
619fn encode_text_build_block(block: &TextOpenBlock) -> Vec<u8> {
620    debug_assert!(!block.rows.is_empty() && block.rows.len() <= POSTING_BLOCK);
621    let mut out = Vec::with_capacity(TEXT_BUILD_BLOCK_MAGIC.len() + 1 + block.rows.len() * 24);
622    out.extend_from_slice(TEXT_BUILD_BLOCK_MAGIC);
623    out.push(block.rows.len() as u8);
624    for &(doc, tf, dl) in &block.rows {
625        out.extend_from_slice(&doc.to_le_bytes());
626        write_varint(&mut out, tf);
627        write_varint(&mut out, dl);
628    }
629    out
630}
631
632fn decode_text_build_block(bytes: &[u8]) -> Result<Vec<(u64, u64, u64)>> {
633    if !bytes.starts_with(TEXT_BUILD_BLOCK_MAGIC) || bytes.len() < 4 {
634        return Err(corrupt("text accumulator block has an invalid header"));
635    }
636    let count = bytes[TEXT_BUILD_BLOCK_MAGIC.len()] as usize;
637    if count == 0 || count > POSTING_BLOCK {
638        return Err(corrupt("text accumulator block has an invalid posting count"));
639    }
640    let mut rows = Vec::with_capacity(count);
641    let mut pos = 4usize;
642    for _ in 0..count {
643        let end = pos.checked_add(8).ok_or(Error::TooLarge)?;
644        let raw = bytes.get(pos..end)
645            .ok_or_else(|| corrupt("text accumulator document id is truncated"))?;
646        let doc = u64::from_le_bytes(raw.try_into().unwrap());
647        pos = end;
648        let tf = required_varint(bytes, &mut pos, "text accumulator frequency is truncated")?;
649        let dl = required_varint(bytes, &mut pos, "text accumulator length is truncated")?;
650        if tf == 0 { return Err(corrupt("text accumulator has zero term frequency")); }
651        rows.push((doc, tf, dl));
652    }
653    if pos != bytes.len() || rows.windows(2).any(|pair| pair[0].0 >= pair[1].0) {
654        return Err(corrupt("text accumulator block is malformed or unordered"));
655    }
656    Ok(rows)
657}
658
659/// Fixed-memory posting accumulator for initial BM25/SEARCH builds. Documents
660/// arrive in increasing id order, so only the term streams interleave. One
661/// open 128-document block per hot term removes the row-per-posting scratch
662/// materialisation. If the vocabulary exceeds the arena, the coldest partial
663/// block is emitted and the final merge re-forms canonical blocks.
664struct TextPostingAccumulator {
665    blocks: HashMap<Vec<u8>, TextOpenBlock>,
666    recency: BTreeSet<(u64, Vec<u8>)>,
667    tracking_recency: bool,
668    used: usize,
669    open_budget: usize,
670    output: crate::bulk::ExternalSort,
671    postings: u64,
672    emissions: u64,
673    partials: u64,
674    partial_framed_bytes: u64,
675    /// Diagnostic ablation: emit one scratch fragment per posting, reproducing
676    /// the pre-accumulator materialisation without changing final bytes.
677    direct_fragments: bool,
678}
679
680impl TextPostingAccumulator {
681    fn new(path: &std::path::Path, open_budget: usize, output_budget: usize) -> Result<Self> {
682        Ok(Self {
683            blocks: HashMap::new(),
684            recency: BTreeSet::new(),
685            tracking_recency: false,
686            used: 0,
687            open_budget,
688            output: crate::bulk::ExternalSort::new(path, output_budget)?,
689            postings: 0,
690            emissions: 0,
691            partials: 0,
692            partial_framed_bytes: 0,
693            direct_fragments: std::env::var_os("SEKEJAP_ABLATE_BM25_ACCUMULATOR").is_some(),
694        })
695    }
696
697    fn entry_bytes(term: &[u8]) -> usize {
698        term.len() * 2 + POSTING_BLOCK * std::mem::size_of::<(u64, u64, u64)>() + 192
699    }
700
701    fn emit(&mut self, term: Vec<u8>, block: TextOpenBlock, partial: bool) -> Result<()> {
702        let mut key = Vec::with_capacity(term.len() + 9);
703        key.extend_from_slice(&term);
704        key.push(0);
705        key.extend_from_slice(&block.rows[0].0.to_be_bytes());
706        let value = encode_text_build_block(&block);
707        if partial {
708            self.partial_framed_bytes = self.partial_framed_bytes
709                .checked_add((12 + key.len() + value.len()) as u64)
710                .ok_or(Error::TooLarge)?;
711        }
712        self.output.push(key, value)?;
713        self.emissions = self.emissions.checked_add(1).ok_or(Error::TooLarge)?;
714        if partial { self.partials = self.partials.checked_add(1).ok_or(Error::TooLarge)?; }
715        Ok(())
716    }
717
718    fn evict_coldest(&mut self) -> Result<()> {
719        if !self.tracking_recency {
720            self.recency.extend(self.blocks.iter()
721                .map(|(term, block)| (block.last_seen, term.clone())));
722            self.tracking_recency = true;
723        }
724        let Some((last_seen, term)) = self.recency.iter().next().cloned() else {
725            return Err(Error::TooLarge);
726        };
727        self.recency.remove(&(last_seen, term.clone()));
728        let block = self.blocks.remove(&term).ok_or(Error::DuplicateKey)?;
729        self.used = self.used.saturating_sub(Self::entry_bytes(&term));
730        self.emit(term, block, true)
731    }
732
733    fn push(&mut self, term: Vec<u8>, doc: u64, tf: u64, dl: u64) -> Result<()> {
734        self.postings = self.postings.checked_add(1).ok_or(Error::TooLarge)?;
735        if self.direct_fragments {
736            return self.emit(term, TextOpenBlock::new(doc, tf, dl), false);
737        }
738        if let Some(block) = self.blocks.get_mut(&term) {
739            if doc <= block.last_seen || tf == 0 { return Err(Error::DuplicateKey); }
740            if self.tracking_recency { self.recency.remove(&(block.last_seen, term.clone())); }
741            block.rows.push((doc, tf, dl));
742            block.last_seen = doc;
743            if self.tracking_recency { self.recency.insert((doc, term.clone())); }
744            if block.rows.len() == POSTING_BLOCK {
745                let block = self.blocks.remove(&term).unwrap();
746                if self.tracking_recency { self.recency.remove(&(doc, term.clone())); }
747                self.used = self.used.saturating_sub(Self::entry_bytes(&term));
748                self.emit(term, block, false)?;
749            }
750            return Ok(());
751        }
752        if tf == 0 { return Err(corrupt("text accumulator has zero term frequency")); }
753        let bytes = Self::entry_bytes(&term);
754        if bytes > self.open_budget { return Err(Error::TooLarge); }
755        while self.used.saturating_add(bytes) > self.open_budget { self.evict_coldest()?; }
756        self.used += bytes;
757        if self.tracking_recency { self.recency.insert((doc, term.clone())); }
758        self.blocks.insert(term, TextOpenBlock::new(doc, tf, dl));
759        Ok(())
760    }
761
762    fn profile(&self) -> (u64, u64, u64, (u64, u64, usize)) {
763        (self.postings, self.emissions + self.blocks.len() as u64, self.partials,
764         self.output.profile())
765    }
766
767    fn flush_run(&mut self) -> Result<()> { self.output.flush_run() }
768
769    fn finish(mut self) -> Result<(crate::bulk::SortedRuns, u64, u64, u64, u64)> {
770        let blocks = std::mem::take(&mut self.blocks);
771        for (term, block) in blocks { self.emit(term, block, false)?; }
772        let partial_scratch = self.partial_framed_bytes;
773        Ok((self.output.finish()?, self.postings, self.emissions, self.partials,
774            partial_scratch))
775    }
776}
777
778struct TextBlockIter {
779    input: crate::bulk::MergeIter,
780    field: u64,
781    pending: Option<(Vec<u8>, Vec<(u64, u64, u64)>, usize)>,
782    previous_term: Option<Vec<u8>>,
783    done: bool,
784}
785
786impl TextBlockIter {
787    fn new(input: crate::bulk::MergeIter, field: u64) -> Self {
788        Self { input, field, pending: None, previous_term: None, done: false }
789    }
790
791    fn parse(key: Vec<u8>, value: Vec<u8>, marker: bool)
792        -> Result<(Vec<u8>, Vec<(u64, u64, u64)>, usize)>
793    {
794        if marker || key.len() < 10 || key[key.len() - 9] != 0 {
795            return Err(corrupt("text accumulator sort emitted a malformed key"));
796        }
797        let term = key[..key.len() - 9].to_vec();
798        let first = u64::from_be_bytes(key[key.len() - 8..].try_into().unwrap());
799        let rows = decode_text_build_block(&value)?;
800        if rows.first().map(|row| row.0) != Some(first) {
801            return Err(corrupt("text accumulator key disagrees with its block"));
802        }
803        Ok((term, rows, 0))
804    }
805}
806
807impl Iterator for TextBlockIter {
808    type Item = Result<(Vec<u8>, Vec<u8>, bool)>;
809
810    fn next(&mut self) -> Option<Self::Item> {
811        if self.done { return None; }
812        let (term, chunk, mut at) = match self.pending.take() {
813            Some(item) => item,
814            None => match self.input.next()? {
815                Ok((key, value, marker)) => match Self::parse(key, value, marker) {
816                    Ok(item) => item,
817                    Err(error) => { self.done = true; return Some(Err(error)); }
818                },
819                Err(error) => { self.done = true; return Some(Err(error)); }
820            },
821        };
822        let first_for_term = self.previous_term.as_deref() != Some(term.as_slice());
823        let mut rows = Vec::with_capacity(POSTING_BLOCK);
824        while at < chunk.len() && rows.len() < POSTING_BLOCK {
825            rows.push(chunk[at]); at += 1;
826        }
827        if at < chunk.len() { self.pending = Some((term.clone(), chunk, at)); }
828        while rows.len() < POSTING_BLOCK && self.pending.is_none() {
829            let Some(item) = self.input.next() else { break };
830            let (next_term, next_rows, mut next_at) = match item {
831                Ok((key, value, marker)) => match Self::parse(key, value, marker) {
832                    Ok(item) => item,
833                    Err(error) => { self.done = true; return Some(Err(error)); }
834                },
835                Err(error) => { self.done = true; return Some(Err(error)); }
836            };
837            if next_term != term {
838                self.pending = Some((next_term, next_rows, next_at));
839                break;
840            }
841            while next_at < next_rows.len() && rows.len() < POSTING_BLOCK {
842                if next_rows[next_at].0 <= rows.last().unwrap().0 {
843                    self.done = true;
844                    return Some(Err(corrupt("text accumulator merge regressed")));
845                }
846                rows.push(next_rows[next_at]); next_at += 1;
847            }
848            if next_at < next_rows.len() {
849                self.pending = Some((next_term, next_rows, next_at));
850            }
851        }
852        let mut value = Vec::new();
853        let mut previous = 0u64;
854        for &(doc, tf, dl) in &rows {
855            write_varint(&mut value, doc - previous);
856            write_varint(&mut value, tf);
857            write_varint(&mut value, dl);
858            previous = doc;
859        }
860        let key = if first_for_term { keys::text_seg_key(self.field, 1, &term) }
861            else { keys::text_seg_block_key(self.field, 1, &term, rows[0].0) };
862        self.previous_term = Some(term);
863        Some(Ok((key, value, false)))
864    }
865}
866
867fn term_number(term: &[u8]) -> u64 {
868    let lo = crc32c::crc32c(term) as u64;
869    let mut salted = Vec::with_capacity(term.len() + 1);
870    salted.extend_from_slice(term); salted.push(0x5D);
871    let id = lo | ((crc32c::crc32c(&salted) as u64) << 32);
872    id.max(1)
873}
874
875fn next_term_probe(id: u64) -> u64 {
876    id.wrapping_add(0x9E37_79B9_7F4A_7C15).max(1)
877}
878
879fn decode_term_row(v: &[u8]) -> Result<(u64, &[u8])> {
880    let mut pos = 0usize;
881    let n = required_varint(v, &mut pos, "text term count is truncated")?;
882    let len = required_varint(v, &mut pos, "text term length is truncated")? as usize;
883    let end = pos.checked_add(len).ok_or_else(|| corrupt("text term boundary overflows"))?;
884    let term = v.get(pos..end).ok_or_else(|| corrupt("text term crosses its row"))?;
885    if term.is_empty() || end != v.len() || std::str::from_utf8(term).is_err() {
886        return Err(corrupt("text term dictionary row is malformed"));
887    }
888    Ok((n, term))
889}
890
891fn encode_term_row(term: &[u8], count: u64) -> Vec<u8> {
892    let mut v = Vec::with_capacity(term.len() + 16);
893    write_varint(&mut v, count); write_varint(&mut v, term.len() as u64); v.extend_from_slice(term);
894    v
895}
896
897impl Graph {
898    fn field_meta(&self, field: u64) -> Result<Option<FieldMeta>> {
899        self.store_ref().get(&keys::text_field_meta_key(field))?
900            .map(|v| FieldMeta::decode(&v)).transpose()
901    }
902
903    fn legacy_segment_metas(&self, field: u64) -> Result<Vec<(u32, SegMeta)>> {
904        let mut out = Vec::new();
905        let mut failure = None;
906        let from = keys::text_meta_key(field, 0);
907        self.store_ref().scan(&from)?.for_each_ref(|key, val| {
908            if key.first() != Some(&keys::TAG_TEXTMETA) || key.len() != 13 { return false; }
909            if u64::from_be_bytes(key[1..9].try_into().unwrap()) != field { return false; }
910            let seg = u32::from_be_bytes(key[9..13].try_into().unwrap());
911            if seg < keys::TEXT_TERM_STATS_SEG {
912                match SegMeta::decode(val) {
913                    Ok(meta) => out.push((seg, meta)),
914                    Err(e) => { failure = Some(e); return false; }
915                }
916            }
917            true
918        })?;
919        if let Some(e) = failure { return Err(e); }
920        Ok(out)
921    }
922
923    fn ensure_field_meta(&mut self, field: u64) -> Result<FieldMeta> {
924        if let Some(mut meta) = self.field_meta(field)? {
925            // Resume a publication that crashed after the one-row manifest
926            // flip but before retirement. Reads are already correct through
927            // redirects; the next writer finishes the idempotent cleanup
928            // before admitting fresh head rows.
929            if let Some(new_seg) = meta.head_redirect {
930                self.rewrite_norm_owners(field, &[0], new_seg)?;
931                self.store().delete_prefix(&keys::text_seg_key(field, 0, b""))?;
932                self.store().delete(&keys::text_meta_key(field, 0))?;
933                self.commit()?; self.checkpoint()?;
934                meta.head_redirect = None;
935                self.put_field_meta(field, &meta)?;
936                self.commit()?; self.checkpoint()?;
937            }
938            if !meta.redirects.is_empty() {
939                let new_seg = meta.redirects[0].1;
940                if meta.redirects.iter().any(|&(_, new)| new != new_seg) {
941                    return Err(corrupt("text manifest contains redirects to multiple pending merges"));
942                }
943                let old: Vec<u32> = meta.redirects.iter().map(|&(old, _)| old).collect();
944                self.rewrite_norm_owners(field, &old, new_seg)?;
945                for &source in &old {
946                    self.store().delete_prefix(&keys::text_seg_key(field, source, b""))?;
947                    self.store().delete_prefix(&keys::text_seg_doc_prefix(field, source))?;
948                    self.store().delete(&keys::text_meta_key(field, source))?;
949                }
950                self.commit()?; self.checkpoint()?;
951                meta.redirects.clear();
952                self.put_field_meta(field, &meta)?;
953                self.commit()?; self.checkpoint()?;
954            }
955            return Ok(meta);
956        }
957        let legacy = self.legacy_segment_metas(field)?;
958        let mut meta = FieldMeta { next_seg: 1, ..FieldMeta::default() };
959        if !legacy.is_empty() { meta.term_stats_ready = false; }
960        for (seg, m) in legacy {
961            let live = m.doc_count.checked_sub(m.dead.len() as u64)
962                .ok_or_else(|| corrupt("text segment has more dead documents than documents"))?;
963            meta.live_docs = meta.live_docs.checked_add(live).ok_or(Error::TooLarge)?;
964            meta.total_tokens = meta.total_tokens
965                .checked_add(m.total_tokens.checked_sub(m.dead_tokens)
966                    .ok_or_else(|| corrupt("text segment dead tokens exceed its token total"))?)
967                .ok_or(Error::TooLarge)?;
968            if seg == 0 { continue; }
969            meta.next_seg = meta.next_seg.max(seg.checked_add(1).ok_or(Error::TooLarge)?);
970            meta.active.push((seg, m.level));
971        }
972        if meta.next_seg >= keys::TEXT_TERM_STATS_SEG { return Err(Error::TooLarge); }
973        self.store().put(&keys::text_field_meta_key(field), &meta.encode())?;
974        Ok(meta)
975    }
976
977    fn put_field_meta(&mut self, field: u64, meta: &FieldMeta) -> Result<()> {
978        self.store().put(&keys::text_field_meta_key(field), &meta.encode())
979    }
980
981    fn term_info(&self, field: u64, term: &str) -> Result<Option<(u64, u64)>> {
982        let mut id = term_number(term.as_bytes());
983        for _ in 0..1024 {
984            let Some(v) = self.store_ref().get(&keys::text_term_id_key(field, id))? else {
985                return Ok(None);
986            };
987            let (count, stored) = decode_term_row(&v)?;
988            if stored == term.as_bytes() { return Ok(Some((count, id))); }
989            id = next_term_probe(id);
990        }
991        Err(corrupt("text term collision chain exceeds its bound"))
992    }
993
994    fn term_by_id(&self, field: u64, id: u64) -> Result<String> {
995        let v = self.store_ref().get(&keys::text_term_id_key(field, id))?
996            .ok_or_else(|| corrupt("text term number has no dictionary entry"))?;
997        let (_, term) = decode_term_row(&v)?;
998        Ok(std::str::from_utf8(term).unwrap().to_owned())
999    }
1000
1001    /// Adjust one word's live-document count and return its compact,
1002    /// field-local number.  A new word is interned exactly once; document
1003    /// norms then store only the number and frequency.
1004    fn adjust_term_df(&mut self, field: u64, term: &str, delta: i8) -> Result<u64> {
1005        let mut id = term_number(term.as_bytes());
1006        for _ in 0..1024 {
1007            let key = keys::text_term_id_key(field, id);
1008            match self.store_ref().get(&key)? {
1009                Some(v) => {
1010                    let (old, stored) = decode_term_row(&v)?;
1011                    if stored != term.as_bytes() { id = next_term_probe(id); continue; }
1012                    let new = match delta {
1013                        1 => old.checked_add(1).ok_or(Error::TooLarge)?,
1014                        -1 => old.checked_sub(1)
1015                            .ok_or_else(|| corrupt("text term document count underflows"))?,
1016                        _ => return Err(Error::TooLarge),
1017                    };
1018                    self.store().put(&key, &encode_term_row(term.as_bytes(), new))?;
1019                    let lex = keys::text_term_lex_key(field, term.as_bytes());
1020                    if new == 0 { self.store().delete(&lex)?; }
1021                    else if old == 0 { self.store().put(&lex, &[])?; }
1022                    return Ok(id);
1023                }
1024                None if delta == 1 => {
1025                    self.store().put(&key, &encode_term_row(term.as_bytes(), 1))?;
1026                    self.store().put(&keys::text_term_lex_key(field, term.as_bytes()), &[])?;
1027                    return Ok(id);
1028                }
1029                None => return Err(corrupt("text term document count is missing")),
1030            }
1031        }
1032        Err(corrupt("text term collision chain exceeds its bound"))
1033    }
1034
1035    /// Initial backfills write a new word's row once with count one, but do
1036    /// not rewrite it for every later document.  `finish_text_build` obtains
1037    /// the exact repeated counts by an external sort over compact norm ids.
1038    fn intern_term_for_build(&mut self, field: u64, term: &str) -> Result<u64> {
1039        let mut id = term_number(term.as_bytes());
1040        for _ in 0..1024 {
1041            let key = keys::text_term_id_key(field, id);
1042            match self.store_ref().get(&key)? {
1043                Some(v) => {
1044                    let (_, stored) = decode_term_row(&v)?;
1045                    if stored == term.as_bytes() { return Ok(id); }
1046                    id = next_term_probe(id);
1047                }
1048                None => {
1049                    self.store().put(&key, &encode_term_row(term.as_bytes(), 1))?;
1050                    self.store().put(&keys::text_term_lex_key(field, term.as_bytes()), &[])?;
1051                    return Ok(id);
1052                }
1053            }
1054        }
1055        Err(corrupt("text term collision chain exceeds its bound"))
1056    }
1057
1058    fn segment_meta(&self, field: u64, seg: u32) -> Result<SegMeta> {
1059        let v = self.store_ref().get(&keys::text_meta_key(field, seg))?
1060            .ok_or_else(|| corrupt("active text segment has no metadata"))?;
1061        SegMeta::decode(&v)
1062    }
1063
1064    /// Remove every posting, norm and statistics row owned by one text
1065    /// field. A rebuild starts from an empty corpus; otherwise its document
1066    /// counters and folded terms describe both the old and new contents.
1067    pub fn clear_text(&mut self, field: u64) -> Result<()> {
1068        self.store().delete_prefix(&keys::text_prefix(field))?;
1069        self.store().delete_prefix(&keys::text_norm_prefix(field))?;
1070        self.store().delete_prefix(&keys::text_meta_prefix(field))?;
1071        Ok(())
1072    }
1073
1074    /// Index `text` for `(field, docid)`, blind writes only: one head row
1075    /// per distinct term, the norm row, and the head meta counters. Cost
1076    /// O(distinct terms) -- never touches other documents (Law 2).
1077    /// Re-indexing the same (field, doc) must be preceded by delete_text.
1078    pub fn index_text(&mut self, field: u64, docid: u64, text: &str) -> Result<()> {
1079        self.index_text_inner(field, docid, text, true, None, None)
1080    }
1081
1082    /// Backfill variant: corpus/document metadata and postings are still
1083    /// complete after every row, but repeated term-count rewrites are deferred
1084    /// to one bounded external aggregation in `finish_text_build`.
1085    pub fn index_text_build(&mut self, field: u64, docid: u64, text: &str) -> Result<()> {
1086        self.index_text_inner(field, docid, text, false, None, None)
1087    }
1088
1089    pub fn index_text_build_cached(&mut self, field: u64, docid: u64, text: &str,
1090                                   cache: &mut TextBuildCache) -> Result<()> {
1091        self.index_text_inner(field, docid, text, false, Some(cache), None)
1092    }
1093
1094    pub fn index_text_build_accum(&mut self, field: u64, docid: u64, text: &str,
1095                                  accumulator: &mut TextBuildAccumulator) -> Result<()> {
1096        self.index_text_inner(field, docid, text, false, None, Some(accumulator))
1097    }
1098
1099    pub fn begin_text_build(&mut self, field: u64) -> Result<()> {
1100        let mut meta = self.ensure_field_meta(field)?;
1101        meta.live_docs = 0; meta.total_tokens = 0; meta.term_stats_ready = false;
1102        self.put_field_meta(field, &meta)?;
1103        self.put_head_build_meta(field, 0, 0)
1104    }
1105
1106    /// Build a fresh text field directly as one immutable segment. The input
1107    /// is a replayable external sort of `(docid_be, utf8_text)` records. No
1108    /// row-per-posting head is ever installed: term/doc records are grouped
1109    /// into 128-document values before their range reaches `graft_sorted_range`.
1110    /// Live writes after publication continue to use segment zero unchanged.
1111    pub fn prepare_text_packed(
1112        field: u64,
1113        docs: &mut crate::bulk::SortedRuns,
1114        scratch: &std::path::Path,
1115    ) -> Result<PackedTextCandidate> {
1116        Self::prepare_text_packed_with_workers(
1117            field, docs, scratch, text_prepare_worker_limit(),
1118        )
1119    }
1120
1121    pub fn prepare_text_packed_with_workers(
1122        field: u64,
1123        docs: &mut crate::bulk::SortedRuns,
1124        scratch: &std::path::Path,
1125        worker_limit: usize,
1126    ) -> Result<PackedTextCandidate> {
1127        let trace = std::env::var_os("SEKEJAP_LOAD_BREAKDOWN").is_some();
1128        let total_started = trace.then(std::time::Instant::now);
1129        let scan_started = trace.then(std::time::Instant::now);
1130        let mut postings = TextPostingAccumulator::new(
1131            &scratch.join(format!("text-{field}-posting-blocks")),
1132            TEXT_BUILD_OPEN_BUDGET, TEXT_BUILD_OUTPUT_BUDGET)?;
1133        let mut norms = crate::bulk::ExternalSort::new(
1134            &scratch.join(format!("text-{field}-norms")), 32 << 20)?;
1135        let mut memberships = crate::bulk::ExternalSort::new(
1136            &scratch.join(format!("text-{field}-members")), 32 << 20)?;
1137        let mut doc_count = 0u64;
1138        let mut total_tokens = 0u64;
1139        let mut norm_min = None;
1140        let mut norm_max = None;
1141        let mut previous_doc = None;
1142
1143        let mut stage = |prepared: PreparedPackedTextDoc| -> Result<()> {
1144            for (term, count) in prepared.terms {
1145                postings.push(term.into_bytes(), prepared.doc, count, prepared.dl)?;
1146            }
1147            let norm_key = keys::text_norm_key(field, prepared.doc);
1148            if norm_min.is_none() { norm_min = Some(norm_key.clone()); }
1149            norm_max = Some(norm_key.clone());
1150            norms.push(norm_key, Norm {
1151                total: prepared.dl, owner: 1, terms: prepared.norm_terms,
1152            }.encode())?;
1153            let mut length = Vec::new(); write_varint(&mut length, prepared.dl);
1154            memberships.push(keys::text_seg_doc_key(field, 1, prepared.doc), length)?;
1155            doc_count = doc_count.checked_add(1).ok_or(Error::TooLarge)?;
1156            total_tokens = total_tokens.checked_add(prepared.dl).ok_or(Error::TooLarge)?;
1157            if doc_count % 8192 == 0 {
1158                postings.flush_run()?;
1159                norms.flush_run()?;
1160                memberships.flush_run()?;
1161            }
1162            Ok(())
1163        };
1164
1165        let worker_limit = worker_limit.clamp(1, 8);
1166        let parallel_tokenize = std::env::var_os("SEKEJAP_INDEX_BUILD_SERIAL").is_none()
1167            && std::env::var_os("SEKEJAP_NO_PARALLEL_TEXT_PREP").is_none()
1168            && worker_limit > 1;
1169        if parallel_tokenize {
1170            const DOC_CHUNK: usize = 512;
1171            let worker_count = worker_limit;
1172            std::thread::scope(|scope| -> Result<()> {
1173                let mut inputs = Vec::with_capacity(worker_count);
1174                let mut outputs = Vec::with_capacity(worker_count);
1175                let mut handles = Vec::with_capacity(worker_count);
1176                for _ in 0..worker_count {
1177                    let (input_tx, input_rx) = std::sync::mpsc::sync_channel::<
1178                        Option<Vec<(u64, String)>>
1179                    >(1);
1180                    let (output_tx, output_rx) = std::sync::mpsc::channel();
1181                    inputs.push(input_tx);
1182                    outputs.push(output_rx);
1183                    handles.push(scope.spawn(move || {
1184                        while let Ok(Some(chunk)) = input_rx.recv() {
1185                            let result = chunk.into_iter().map(|(doc, text)|
1186                                prepare_packed_text_doc(doc, &text)
1187                            ).collect::<Result<Vec<_>>>();
1188                            if output_tx.send(result).is_err() { break; }
1189                        }
1190                    }));
1191                }
1192
1193                let mut input = docs.iter()?;
1194                loop {
1195                    let mut active = 0usize;
1196                    for worker in 0..worker_count {
1197                        let mut chunk = Vec::with_capacity(DOC_CHUNK);
1198                        while chunk.len() < DOC_CHUNK {
1199                            let Some(item) = input.next() else { break };
1200                            let (doc_key, text, marker) = item?;
1201                            if marker || doc_key.len() != 8 {
1202                                return Err(corrupt("packed text source has a malformed document key"));
1203                            }
1204                            let doc = u64::from_be_bytes(doc_key.try_into().unwrap());
1205                            if previous_doc.is_some_and(|previous| previous >= doc) {
1206                                return Err(corrupt("packed text documents are not strictly ordered"));
1207                            }
1208                            previous_doc = Some(doc);
1209                            let text = String::from_utf8(text)
1210                                .map_err(|_| corrupt("packed text source is not UTF-8"))?;
1211                            chunk.push((doc, text));
1212                        }
1213                        if chunk.is_empty() { break; }
1214                        inputs[worker].send(Some(chunk))
1215                            .map_err(|_| corrupt("packed text worker stopped before input"))?;
1216                        active += 1;
1217                    }
1218                    if active == 0 { break; }
1219                    for output in outputs.iter().take(active) {
1220                        let prepared = output.recv()
1221                            .map_err(|_| corrupt("packed text worker stopped before output"))??;
1222                        for document in prepared { stage(document)?; }
1223                    }
1224                }
1225                for input in &inputs { let _ = input.send(None); }
1226                for handle in handles {
1227                    handle.join().map_err(|_| corrupt("packed text worker panicked"))?;
1228                }
1229                Ok(())
1230            })?;
1231        } else {
1232            for item in docs.iter()? {
1233                let (doc_key, text, marker) = item?;
1234                if marker || doc_key.len() != 8 {
1235                    return Err(corrupt("packed text source has a malformed document key"));
1236                }
1237                let doc = u64::from_be_bytes(doc_key.try_into().unwrap());
1238                if previous_doc.is_some_and(|previous| previous >= doc) {
1239                    return Err(corrupt("packed text documents are not strictly ordered"));
1240                }
1241                previous_doc = Some(doc);
1242                let text = std::str::from_utf8(&text)
1243                    .map_err(|_| corrupt("packed text source is not UTF-8"))?;
1244                stage(prepare_packed_text_doc(doc, text)?)?;
1245            }
1246        }
1247
1248        let scan_stage =
1249            scan_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1250        let posting_profile = postings.profile();
1251        let norm_profile = norms.profile();
1252        let membership_profile = memberships.profile();
1253        let finish_sources_started = trace.then(std::time::Instant::now);
1254        let (mut posting_runs, posting_rows, emissions, partials, partial_block_scratch) =
1255            postings.finish()?;
1256        let mut norm_runs = norms.finish()?;
1257        let mut membership_runs = memberships.finish()?;
1258        let source_finish =
1259            finish_sources_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1260        let posting_input_runs = posting_runs.run_count();
1261        let norm_input_runs = norm_runs.run_count();
1262        let membership_input_runs = membership_runs.run_count();
1263
1264        let mut terms = crate::bulk::ExternalSort::new(
1265            &scratch.join(format!("text-{field}-terms")), 16 << 20)?;
1266        let mut expected = SegMeta {
1267            doc_count, total_tokens, dead_tokens: 0, dead: Vec::new(), level: 0,
1268            term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0,
1269        };
1270        let mut block_min: Option<Vec<u8>> = None;
1271        let mut block_max: Option<Vec<u8>> = None;
1272        let mut previous_key: Option<Vec<u8>> = None;
1273        let mut current_term: Option<Vec<u8>> = None;
1274        let mut term_docs = 0u64;
1275        let posting_merge_started = trace.then(std::time::Instant::now);
1276        for item in TextBlockIter::new(posting_runs.iter()?, field) {
1277            let (key, value, marker) = item?;
1278            if marker || previous_key.as_ref().is_some_and(|old| old >= &key) {
1279                return Err(corrupt("packed text block stream is not strictly ordered"));
1280            }
1281            let prefix_len = 1 + 8 + 4;
1282            if key.len() <= prefix_len { return Err(corrupt("packed text block key has no term")); }
1283            let body = &key[prefix_len..];
1284            let term_end = body.iter().position(|byte| *byte == 0).unwrap_or(body.len());
1285            let term = &body[..term_end];
1286            if term.is_empty() { return Err(corrupt("packed text block has an empty term")); }
1287            if current_term.as_deref().is_some_and(|old| old != term) {
1288                let old = current_term.take().unwrap();
1289                let mut term_key = term_number(&old).to_be_bytes().to_vec();
1290                term_key.extend_from_slice(&old);
1291                terms.push(term_key, term_docs.to_be_bytes().to_vec())?;
1292                term_docs = 0;
1293            }
1294            if current_term.is_none() { current_term = Some(term.to_vec()); }
1295            for (doc, tf, dl) in decode_postings(&value)? {
1296                expected.posting_count = expected.posting_count.checked_add(1).ok_or(Error::TooLarge)?;
1297                let hash = posting_hash(term, doc, tf, dl);
1298                expected.logical_xor ^= hash;
1299                expected.logical_sum = expected.logical_sum.wrapping_add(hash);
1300                term_docs = term_docs.checked_add(1).ok_or(Error::TooLarge)?;
1301            }
1302            expected.term_rows = expected.term_rows.checked_add(1).ok_or(Error::TooLarge)?;
1303            if block_min.is_none() { block_min = Some(key.clone()); }
1304            block_max = Some(key.clone());
1305            previous_key = Some(key);
1306        }
1307        if let Some(term) = current_term.take() {
1308            let mut term_key = term_number(&term).to_be_bytes().to_vec();
1309            term_key.extend_from_slice(&term);
1310            terms.push(term_key, term_docs.to_be_bytes().to_vec())?;
1311        }
1312        if expected.posting_count != posting_rows {
1313            return Err(corrupt("text accumulator lost or duplicated a posting"));
1314        }
1315        let block_rows = expected.term_rows;
1316        let posting_merge =
1317            posting_merge_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1318        let term_profile = terms.profile();
1319        let block_term_finish_started = trace.then(std::time::Instant::now);
1320        let mut term_runs = terms.finish()?;
1321        let block_term_finish = block_term_finish_started
1322            .map_or(std::time::Duration::ZERO, |started| started.elapsed());
1323        let term_input_runs = term_runs.run_count();
1324        let mut dictionary = crate::bulk::ExternalSort::new(
1325            &scratch.join(format!("text-{field}-dictionary")), 16 << 20)?;
1326        let mut dict_rows = 0u64;
1327        let mut previous_actual = None;
1328        let mut current_raw = None;
1329        let mut probe_offset = 0u32;
1330        let dictionary_started = trace.then(std::time::Instant::now);
1331        for item in term_runs.iter()? {
1332            let (key, count, marker) = item?;
1333            if marker || key.len() < 9 || count.len() != 8 {
1334                return Err(corrupt("packed text term sort emitted a malformed row"));
1335            }
1336            let raw = u64::from_be_bytes(key[..8].try_into().unwrap());
1337            let term = &key[8..];
1338            if current_raw != Some(raw) { current_raw = Some(raw); probe_offset = 0; }
1339            let mut actual = raw;
1340            for _ in 0..probe_offset { actual = next_term_probe(actual); }
1341            probe_offset = probe_offset.checked_add(1).ok_or(Error::TooLarge)?;
1342            if probe_offset > 1024 { return Err(corrupt("packed text term collision chain exceeds its bound")); }
1343            if previous_actual == Some(actual) {
1344                return Err(corrupt("packed text term assignment produced a duplicate id"));
1345            }
1346            previous_actual = Some(actual);
1347            let df = u64::from_be_bytes(count.try_into().unwrap());
1348            dictionary.push(keys::text_term_id_key(field, actual), encode_term_row(term, df))?;
1349            dictionary.push(keys::text_term_lex_key(field, term), Vec::new())?;
1350            dict_rows = dict_rows.checked_add(2).ok_or(Error::TooLarge)?;
1351            if actual != raw {
1352                // The double-CRC term number makes this astronomically rare.
1353                // Refuse rather than publish norms whose ids need a corpus
1354                // rewrite; a future collision fixture can justify the extra
1355                // external join without charging every ordinary build today.
1356                return Err(corrupt("packed text build encountered a term-number collision"));
1357            }
1358        }
1359        let dictionary_stage =
1360            dictionary_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1361        let dictionary_profile = dictionary.profile();
1362        let dictionary_finish_started = trace.then(std::time::Instant::now);
1363        let mut dictionary_runs = dictionary.finish()?;
1364        let dictionary_finish = dictionary_finish_started
1365            .map_or(std::time::Duration::ZERO, |started| started.elapsed());
1366        let dictionary_input_runs = dictionary_runs.run_count();
1367
1368        // Validate every replayable stream, including duplicate keys, before
1369        // the first graft can publish any of the three physical ranges.
1370        fn validate_sorted(
1371            iter: impl Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1372        ) -> Result<(u64, Option<Vec<u8>>, Option<Vec<u8>>)> {
1373            let mut rows = 0u64; let mut min = None; let mut max = None;
1374            for item in iter {
1375                let (key, _, _) = item?;
1376                if max.as_ref().is_some_and(|old: &Vec<u8>| old >= &key) {
1377                    return Err(Error::DuplicateKey);
1378                }
1379                if min.is_none() { min = Some(key.clone()); }
1380                max = Some(key); rows = rows.checked_add(1).ok_or(Error::TooLarge)?;
1381            }
1382            Ok((rows, min, max))
1383        }
1384        let validate_started = trace.then(std::time::Instant::now);
1385        let (norm_rows, checked_norm_min, checked_norm_max) = validate_sorted(norm_runs.iter()?)?;
1386        if norm_rows != doc_count || checked_norm_min != norm_min || checked_norm_max != norm_max {
1387            return Err(corrupt("packed text norm manifest disagrees with its stream"));
1388        }
1389        let (member_rows, member_min, member_max) = validate_sorted(membership_runs.iter()?)?;
1390        let (checked_dict_rows, dict_min, dict_max) = validate_sorted(dictionary_runs.iter()?)?;
1391        if member_rows != doc_count || checked_dict_rows != dict_rows {
1392            return Err(corrupt("packed text metadata manifest disagrees with its streams"));
1393        }
1394        let validate =
1395            validate_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1396
1397        Ok(PackedTextCandidate {
1398            field,
1399            posting_runs,
1400            norm_runs,
1401            membership_runs,
1402            dictionary_runs,
1403            expected,
1404            doc_count,
1405            total_tokens,
1406            posting_rows,
1407            emissions,
1408            partials,
1409            block_rows,
1410            block_min,
1411            block_max,
1412            norm_rows,
1413            norm_min: checked_norm_min,
1414            norm_max: checked_norm_max,
1415            member_rows,
1416            member_min,
1417            member_max,
1418            dict_rows,
1419            dict_min,
1420            dict_max,
1421            scan_stage,
1422            source_finish,
1423            posting_merge,
1424            block_term_finish,
1425            dictionary_stage,
1426            dictionary_finish,
1427            validate,
1428            prepare_total: total_started
1429                .map_or(std::time::Duration::ZERO, |started| started.elapsed()),
1430            posting_scratch: posting_profile.3.1,
1431            norm_scratch: norm_profile.1,
1432            membership_scratch: membership_profile.1,
1433            term_scratch: term_profile.1,
1434            dictionary_scratch: dictionary_profile.1,
1435            partial_block_scratch,
1436            posting_input_runs,
1437            norm_input_runs,
1438            membership_input_runs,
1439            term_input_runs,
1440            dictionary_input_runs,
1441        })
1442    }
1443
1444    /// Graft one privately built candidate in the caller's deterministic
1445    /// single-writer order, then publish and independently verify it.
1446    pub fn publish_text_packed(
1447        &mut self,
1448        mut candidate: PackedTextCandidate,
1449        scratch: &std::path::Path,
1450    ) -> Result<()> {
1451        let trace = std::env::var_os("SEKEJAP_LOAD_BREAKDOWN").is_some();
1452        let graft_started = trace.then(std::time::Instant::now);
1453        if candidate.block_rows != 0 {
1454            if trace { eprintln!("consumer BM25 postings graft phases:"); }
1455            self.store().graft_sorted_range(
1456                TextBlockIter::new(candidate.posting_runs.iter()?, candidate.field),
1457                candidate.block_rows,
1458                candidate.block_min.take().unwrap(),
1459                candidate.block_max.take().unwrap(),
1460                scratch,
1461            )?;
1462        }
1463        if candidate.norm_rows != 0 {
1464            if trace { eprintln!("consumer BM25 norms graft phases:"); }
1465            self.store().graft_sorted_range(
1466                candidate.norm_runs.iter()?,
1467                candidate.norm_rows,
1468                candidate.norm_min.take().unwrap(),
1469                candidate.norm_max.take().unwrap(),
1470                scratch,
1471            )?;
1472        }
1473        if candidate.member_rows != 0 {
1474            if trace { eprintln!("consumer BM25 memberships graft phases:"); }
1475            self.store().graft_sorted_range(
1476                candidate.membership_runs.iter()?,
1477                candidate.member_rows,
1478                candidate.member_min.take().unwrap(),
1479                candidate.member_max.take().unwrap(),
1480                scratch,
1481            )?;
1482        }
1483        if candidate.dict_rows != 0 {
1484            if trace { eprintln!("consumer BM25 dictionary graft phases:"); }
1485            self.store().graft_sorted_range(
1486                candidate.dictionary_runs.iter()?,
1487                candidate.dict_rows,
1488                candidate.dict_min.take().unwrap(),
1489                candidate.dict_max.take().unwrap(),
1490                scratch,
1491            )?;
1492        }
1493        let graft = graft_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1494
1495        let publish_started = trace.then(std::time::Instant::now);
1496        let field_meta = FieldMeta {
1497            live_docs: candidate.doc_count,
1498            total_tokens: candidate.total_tokens,
1499            next_seg: 2,
1500            generation: 1,
1501            active: vec![(1, 0)],
1502            head_redirect: None,
1503            redirects: Vec::new(),
1504            term_stats_ready: true,
1505        };
1506        self.store().put(&keys::text_meta_key(candidate.field, 1), &candidate.expected.encode())?;
1507        self.store().put(&keys::text_field_meta_key(candidate.field), &field_meta.encode())?;
1508        self.commit()?;
1509        self.checkpoint()?;
1510
1511        let cfg = crate::store::Config {
1512            budget_bytes: 16 * crate::page::PAGE_SIZE,
1513            io: self.store_ref().io_mode(),
1514            sync: crate::store::SyncMode::Off,
1515        };
1516        let snapshot = crate::store::Store::open_snapshot(self.store_ref().dir(), cfg)?;
1517        let result = Graph::new(snapshot)?
1518            .verify_text_segment(candidate.field, 1, &candidate.expected);
1519        let publish_verify = publish_started
1520            .map_or(std::time::Duration::ZERO, |started| started.elapsed());
1521        if trace {
1522            let mb = |bytes: u64| bytes as f64 / (1024.0 * 1024.0);
1523            let doc_count = candidate.doc_count;
1524            eprintln!(
1525                "\nBM25 packed build detail ({} documents, {} term/doc postings)",
1526                doc_count, candidate.posting_rows,
1527            );
1528            eprintln!("{:<34} {:>11} {:>12}", "item", "seconds", "ns/doc");
1529            eprintln!("{}", "-".repeat(61));
1530            let print = |name: &str, elapsed: std::time::Duration| eprintln!(
1531                "{name:<34} {:>11.6} {:>12.1}", elapsed.as_secs_f64(),
1532                elapsed.as_secs_f64() * 1e9 / doc_count.max(1) as f64);
1533            print("tokenize + stage raw rows", candidate.scan_stage);
1534            print("finish raw sorters", candidate.source_finish);
1535            print("merge postings + form blocks", candidate.posting_merge);
1536            print("finish block + term sorters", candidate.block_term_finish);
1537            print("build dictionary rows", candidate.dictionary_stage);
1538            print("finish dictionary sorter", candidate.dictionary_finish);
1539            print("validate all replay streams", candidate.validate);
1540            print("pack + verify + graft 4 ranges", graft);
1541            print("publish + logical reopen verify", publish_verify);
1542            print("BM25 packed total", candidate.prepare_total + graft + publish_verify);
1543            eprintln!(
1544                "postings: {} rows, {:.2}/doc -> {} packed rows ({:.2} postings/block)",
1545                candidate.posting_rows,
1546                candidate.posting_rows as f64 / doc_count.max(1) as f64,
1547                candidate.block_rows,
1548                candidate.posting_rows as f64 / candidate.block_rows.max(1) as f64,
1549            );
1550            eprintln!(
1551                "accumulator emissions: {} ({} partial spills), open arena {} MiB, finished-block sorter {} MiB",
1552                candidate.emissions, candidate.partials,
1553                TEXT_BUILD_OPEN_BUDGET >> 20, TEXT_BUILD_OUTPUT_BUDGET >> 20,
1554            );
1555            eprintln!(
1556                "accumulator churn: evictions={} partial-spills={} merge-fragments={} (emitted fragments minus canonical blocks)",
1557                candidate.partials,
1558                candidate.partials,
1559                candidate.emissions.saturating_sub(candidate.block_rows),
1560            );
1561            eprintln!(
1562                "dictionary: {} distinct terms, {} rows; norms: {}; memberships: {}",
1563                candidate.dict_rows / 2, candidate.dict_rows,
1564                candidate.norm_rows, candidate.member_rows,
1565            );
1566            eprintln!(
1567                "scratch first-pass: posting blocks {:.2} MiB/{} runs; norms {:.2} MiB/{}; members {:.2} MiB/{}; terms {:.2} MiB/{}; dictionary {:.2} MiB/{}",
1568                mb(candidate.posting_scratch), candidate.posting_input_runs,
1569                mb(candidate.norm_scratch), candidate.norm_input_runs,
1570                mb(candidate.membership_scratch), candidate.membership_input_runs,
1571                mb(candidate.term_scratch), candidate.term_input_runs,
1572                mb(candidate.dictionary_scratch), candidate.dictionary_input_runs,
1573            );
1574            eprintln!(
1575                "partial-block scratch subset: {:.2} MiB ({} evicted fragments; included in posting blocks)",
1576                mb(candidate.partial_block_scratch), candidate.partials,
1577            );
1578        }
1579        result
1580    }
1581
1582    pub fn build_text_packed(
1583        &mut self,
1584        field: u64,
1585        docs: &mut crate::bulk::SortedRuns,
1586        scratch: &std::path::Path,
1587    ) -> Result<()> {
1588        let candidate = Self::prepare_text_packed(field, docs, scratch)?;
1589        self.publish_text_packed(candidate, scratch)
1590    }
1591
1592    /// Copy a quiescent packed initial segment to another physical field id.
1593    /// This is how a BM25 build and a later SEARCH build over the same source
1594    /// text share tokenisation without sharing physical indexes. A field with
1595    /// any live head or merge history is refused so CREATE INDEX falls back to
1596    /// its ordinary source scan rather than copying a mutable shape.
1597    pub fn clone_packed_text(
1598        &mut self,
1599        source: u64,
1600        target: u64,
1601        scratch: &std::path::Path,
1602    ) -> Result<bool> {
1603        let trace = std::env::var_os("SEKEJAP_LOAD_BREAKDOWN").is_some();
1604        let total_started = trace.then(std::time::Instant::now);
1605        let mut scan_elapsed = std::time::Duration::ZERO;
1606        let mut graft_elapsed = std::time::Duration::ZERO;
1607        let Some(meta) = self.field_meta(source)? else {
1608            return Ok(false);
1609        };
1610        if meta.active.as_slice() != [(1, 0)]
1611            || meta.head_redirect.is_some()
1612            || !meta.redirects.is_empty()
1613            || !meta.term_stats_ready
1614        {
1615            return Ok(false);
1616        }
1617        let head_prefix = keys::text_seg_key(source, 0, b"");
1618        if let Some(row) = self.store_ref().scan(&head_prefix)?.next() {
1619            let (key, _) = row?;
1620            if key.starts_with(&head_prefix) { return Ok(false); }
1621        }
1622
1623        if std::env::var_os("SEKEJAP_ABLATE_SEARCH_CLONE").is_some() {
1624            let ablate_started = trace.then(std::time::Instant::now);
1625            let stage = |prefix: Vec<u8>,
1626                         name: &str|
1627             -> Result<(
1628                crate::bulk::ExternalSort,
1629                u64,
1630                Option<Vec<u8>>,
1631                Option<Vec<u8>>,
1632            )> {
1633                let mut sort = crate::bulk::ExternalSort::new(
1634                    &scratch.join(format!("clone-{target}-{name}")),
1635                    32 << 20,
1636                )?;
1637                let mut rows = 0u64;
1638                let mut min = None;
1639                let mut max = None;
1640                let mut failure = None;
1641                self.store_ref().scan(&prefix)?.for_each_ref(|key, value| {
1642                    if !key.starts_with(&prefix) {
1643                        return false;
1644                    }
1645                    let mut rewritten = key.to_vec();
1646                    rewritten[1..9].copy_from_slice(&target.to_be_bytes());
1647                    if min.is_none() {
1648                        min = Some(rewritten.clone());
1649                    }
1650                    max = Some(rewritten.clone());
1651                    rows += 1;
1652                    match sort.push(rewritten, value.to_vec()) {
1653                        Ok(()) => true,
1654                        Err(error) => {
1655                            failure = Some(error);
1656                            false
1657                        }
1658                    }
1659                })?;
1660                if let Some(error) = failure {
1661                    return Err(error);
1662                }
1663                Ok((sort, rows, min, max))
1664            };
1665            let (text, text_rows, text_min, text_max) =
1666                stage(keys::text_prefix(source), "postings")?;
1667            let (norms, norm_rows, norm_min, norm_max) =
1668                stage(keys::text_norm_prefix(source), "norms")?;
1669            let (metadata, meta_rows, meta_min, meta_max) =
1670                stage(keys::text_meta_prefix(source), "metadata")?;
1671            let mut text = text.finish()?;
1672            let mut norms = norms.finish()?;
1673            let mut metadata = metadata.finish()?;
1674            if text_rows != 0 {
1675                self.store().graft_sorted_range(
1676                    text.iter()?,
1677                    text_rows,
1678                    text_min.unwrap(),
1679                    text_max.unwrap(),
1680                    scratch,
1681                )?;
1682            }
1683            if norm_rows != 0 {
1684                self.store().graft_sorted_range(
1685                    norms.iter()?,
1686                    norm_rows,
1687                    norm_min.unwrap(),
1688                    norm_max.unwrap(),
1689                    scratch,
1690                )?;
1691            }
1692            if meta_rows != 0 {
1693                self.store().graft_sorted_range(
1694                    metadata.iter()?,
1695                    meta_rows,
1696                    meta_min.unwrap(),
1697                    meta_max.unwrap(),
1698                    scratch,
1699                )?;
1700            }
1701            let expected = self.segment_meta(target, 1)?;
1702            let cfg = crate::store::Config {
1703                budget_bytes: 16 * crate::page::PAGE_SIZE,
1704                io: self.store_ref().io_mode(),
1705                sync: crate::store::SyncMode::Off,
1706            };
1707            let snapshot = crate::store::Store::open_snapshot(self.store_ref().dir(), cfg)?;
1708            Graph::new(snapshot)?.verify_text_segment(target, 1, &expected)?;
1709            if trace {
1710                eprintln!(
1711                    "consumer SEARCH re-sort ablation total: {:.6}s",
1712                    ablate_started.unwrap().elapsed().as_secs_f64()
1713                );
1714            }
1715            return Ok(true);
1716        }
1717
1718        let cfg = crate::store::Config {
1719            budget_bytes: 16 * crate::page::PAGE_SIZE,
1720            io: self.store_ref().io_mode(),
1721            sync: crate::store::SyncMode::Off,
1722        };
1723        // The source was just logically verified and published. Read it from
1724        // an independent snapshot so its range iterator can remain live while
1725        // this writer packs the target. Replacing the fixed-width field id in
1726        // every key preserves key order exactly, so sorting the clone again
1727        // only wrote and reread the whole segment as scratch.
1728        let source_snapshot = crate::store::Store::open_snapshot(self.store_ref().dir(), cfg)?;
1729        let source_graph = Graph::new(source_snapshot)?;
1730        let mut graft_prefix = |name: &str, prefix: Vec<u8>| -> Result<()> {
1731            let scan_started = trace.then(std::time::Instant::now);
1732            let mut rows = 0u64;
1733            let mut min = None;
1734            let mut max = None;
1735            let mut count_error = None;
1736            source_graph
1737                .store_ref()
1738                .scan(&prefix)?
1739                .for_each_ref(|key, _| {
1740                    if !key.starts_with(&prefix) {
1741                        return false;
1742                    }
1743                    let mut rewritten = key.to_vec();
1744                    rewritten[1..9].copy_from_slice(&target.to_be_bytes());
1745                    if min.is_none() {
1746                        min = Some(rewritten.clone());
1747                    }
1748                    max = Some(rewritten);
1749                    match rows.checked_add(1) {
1750                        Some(next) => rows = next,
1751                        None => {
1752                            count_error = Some(Error::TooLarge);
1753                            return false;
1754                        }
1755                    }
1756                    true
1757                })?;
1758            if let Some(started) = scan_started {
1759                scan_elapsed += started.elapsed();
1760            }
1761            if let Some(error) = count_error {
1762                return Err(error);
1763            }
1764            if rows == 0 {
1765                return Ok(());
1766            }
1767
1768            let mut source_rows = source_graph.store_ref().scan(&prefix)?;
1769            let iter_prefix = prefix.clone();
1770            let iter = std::iter::from_fn(move || match source_rows.next() {
1771                Some(Ok((mut key, value))) if key.starts_with(&iter_prefix) => {
1772                    key[1..9].copy_from_slice(&target.to_be_bytes());
1773                    Some(Ok((key, value, false)))
1774                }
1775                Some(Ok(_)) | None => None,
1776                Some(Err(error)) => Some(Err(error)),
1777            });
1778            if trace {
1779                eprintln!("consumer SEARCH clone {name} graft phases:");
1780            }
1781            let graft_started = trace.then(std::time::Instant::now);
1782            let result =
1783                self.store()
1784                    .graft_sorted_range(iter, rows, min.unwrap(), max.unwrap(), scratch);
1785            if let Some(started) = graft_started {
1786                graft_elapsed += started.elapsed();
1787            }
1788            result
1789        };
1790        graft_prefix("postings", keys::text_prefix(source))?;
1791        graft_prefix("norms", keys::text_norm_prefix(source))?;
1792        graft_prefix("metadata", keys::text_meta_prefix(source))?;
1793        drop(graft_prefix);
1794        drop(source_graph);
1795
1796        let logical_verify_started = trace.then(std::time::Instant::now);
1797        let expected = self.segment_meta(target, 1)?;
1798        let snapshot = crate::store::Store::open_snapshot(self.store_ref().dir(), cfg)?;
1799        Graph::new(snapshot)?.verify_text_segment(target, 1, &expected)?;
1800        if trace {
1801            eprintln!(
1802                "consumer SEARCH order-preserving clone: accumulate/rewrite-scan={:.6}s sort=0.000000s pack+verify+graft={:.6}s logical-verify={:.6}s total={:.6}s",
1803                scan_elapsed.as_secs_f64(),
1804                graft_elapsed.as_secs_f64(),
1805                logical_verify_started.unwrap().elapsed().as_secs_f64(),
1806                total_started.unwrap().elapsed().as_secs_f64()
1807            );
1808        }
1809        Ok(true)
1810    }
1811
1812    fn index_text_inner(&mut self, field: u64, docid: u64, text: &str,
1813                        maintain_term_counts: bool, mut build_cache: Option<&mut TextBuildCache>,
1814                        mut accumulator: Option<&mut TextBuildAccumulator>) -> Result<()> {
1815        let tokens = tokenize(text);
1816        let total = tokens.len() as u64;
1817        if let Some(accumulator) = accumulator.as_deref_mut() { accumulator.document(total)?; }
1818        if maintain_term_counts {
1819            if let Some(v) = self.store_ref().get(&keys::text_norm_key(field, docid))? {
1820                let old = Norm::decode(&v)?;
1821                return self.replace_norm(field, docid, old, text);
1822            }
1823        }
1824        let mut meta = if maintain_term_counts { Some(self.ensure_field_meta(field)?) } else { None };
1825        let mut tf: HashMap<String, u64> = HashMap::new();
1826        for t in tokens { *tf.entry(t).or_insert(0) += 1; }
1827        let mut norm_terms = Vec::with_capacity(tf.len());
1828        for (term, count) in &tf {
1829            let mut v = Vec::new();
1830            write_varint(&mut v, *count);
1831            // the doc's token count rides in EVERY posting (2h A3): BM25
1832            // then needs zero length lookups -- the 1M term query paid 200K
1833            // cold point-gets for lengths and took 10.8s before this.
1834            write_varint(&mut v, total);
1835            self.store().put(&keys::text_head_key(field, term.as_bytes(), docid), &v)?;
1836            let id = if let Some(accumulator) = accumulator.as_deref_mut() {
1837                accumulator.term(term, docid)?
1838            } else if maintain_term_counts {
1839                self.adjust_term_df(field, term, 1)?
1840            } else if let Some(id) = build_cache.as_deref().and_then(|cache| cache.get(term)) {
1841                id
1842            } else {
1843                let id = self.intern_term_for_build(field, term)?;
1844                if let Some(cache) = build_cache.as_deref_mut() { cache.insert(term, id); }
1845                id
1846            };
1847            norm_terms.push((id, *count));
1848        }
1849        let old_head_len = if maintain_term_counts {
1850            self.store_ref().get(&keys::text_head_doc_key(field, docid))?
1851                .map(|v| { let mut p = 0; required_varint(&v, &mut p, "head document length is truncated") })
1852                .transpose()?
1853        } else {
1854            // A backfill follows clear_text, so no head-length row can exist.
1855            None
1856        };
1857        if total == 0 {
1858            let mut dv = Vec::new(); write_varint(&mut dv, total);
1859            self.store().put(&keys::text_head_doc_key(field, docid), &dv)?;
1860        } else if old_head_len.is_some() {
1861            self.store().delete(&keys::text_head_doc_key(field, docid))?;
1862        }
1863        norm_terms.sort_unstable_by_key(|&(id, _)| id);
1864        let norm = Norm { total, owner: 0, terms: norm_terms };
1865        self.store().put(&keys::text_norm_key(field, docid), &norm.encode())?;
1866        if maintain_term_counts {
1867            // Ordinary writes keep this one small row current. A known-empty
1868            // initial backfill publishes its accumulated totals once at the
1869            // end instead of rewriting the same page for every document.
1870            let mk = keys::text_meta_key(field, 0);
1871            let mut m = self.store().get(&mk)?.map(|v| SegMeta::decode(&v)).transpose()?
1872                .unwrap_or(SegMeta { doc_count: 0, total_tokens: 0, dead_tokens: 0, dead: Vec::new(),
1873                    level: 0, term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0 });
1874            let resurrecting_head = m.dead.contains(&docid);
1875            if resurrecting_head {
1876                let old_tokens = old_head_len.unwrap_or(0);
1877                m.total_tokens = m.total_tokens.saturating_sub(old_tokens) + total;
1878                m.dead_tokens = m.dead_tokens.saturating_sub(old_tokens);
1879                m.dead.retain(|&d| d != docid);
1880            } else {
1881                m.doc_count += 1;
1882                m.total_tokens += total;
1883                m.dead.retain(|&d| d != docid);
1884            }
1885            self.store().put(&mk, &m.encode())?;
1886            let meta = meta.as_mut().unwrap();
1887            meta.live_docs = meta.live_docs.checked_add(1).ok_or(Error::TooLarge)?;
1888            meta.total_tokens = meta.total_tokens.checked_add(total).ok_or(Error::TooLarge)?;
1889            self.put_field_meta(field, meta)?;
1890        }
1891        Ok(())
1892    }
1893
1894    fn put_head_build_meta(&mut self, field: u64, docs: u64, tokens: u64) -> Result<()> {
1895        self.store().put(&keys::text_meta_key(field, 0), &SegMeta {
1896            doc_count: docs, total_tokens: tokens, dead_tokens: 0, dead: Vec::new(),
1897            level: 0, term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0,
1898        }.encode())
1899    }
1900
1901    /// Finish an initial backfill by grouping compact `(term id, exact word)`
1902    /// counts outside the database. The accumulator owns 8 MiB regardless of
1903    /// corpus size and spills runs to scratch; only repeated terms need a
1904    /// second dictionary write because one-document terms already hold `1`.
1905    pub fn finish_text_build(&mut self, field: u64) -> Result<()> {
1906        use std::sync::atomic::{AtomicU64, Ordering};
1907        static SEQ: AtomicU64 = AtomicU64::new(0);
1908        let scratch = std::env::temp_dir().join(format!("text-stats-{}-{}",
1909            std::process::id(), SEQ.fetch_add(1, Ordering::Relaxed)));
1910        let mut sort = crate::bulk::ExternalSort::new(&scratch, 8 << 20)?;
1911        let prefix = keys::text_norm_prefix(field);
1912        let mut failure = None; let mut docs = 0u64; let mut tokens = 0u64;
1913        self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
1914            if !key.starts_with(&prefix) { return false; }
1915            match Norm::decode(val) {
1916                Ok(norm) => {
1917                    docs = match docs.checked_add(1) { Some(n) => n, None => { failure = Some(Error::TooLarge); return false; } };
1918                    tokens = match tokens.checked_add(norm.total) { Some(n) => n, None => { failure = Some(Error::TooLarge); return false; } };
1919                    for (id, _) in norm.terms {
1920                        if let Err(e) = sort.push(id.to_be_bytes().to_vec(), Vec::new()) {
1921                            failure = Some(e); return false;
1922                        }
1923                    }
1924                }
1925                Err(e) => { failure = Some(e); return false; }
1926            }
1927            true
1928        })?;
1929        if let Some(e) = failure { return Err(e); }
1930        let mut runs = sort.finish()?;
1931        let mut iter = runs.iter()?;
1932        let mut current = None; let mut count = 0u64;
1933        while let Some(item) = iter.next() {
1934            let (key, _, _) = item?;
1935            if key.len() != 8 { return Err(corrupt("text statistic sort emitted an invalid key")); }
1936            let id = u64::from_be_bytes(key.try_into().unwrap());
1937            if current.is_some_and(|old| old != id) {
1938                self.publish_built_term_count(field, current.unwrap(), count)?;
1939                count = 0;
1940            }
1941            current = Some(id); count = count.checked_add(1).ok_or(Error::TooLarge)?;
1942        }
1943        if let Some(id) = current { self.publish_built_term_count(field, id, count)?; }
1944        self.put_head_build_meta(field, docs, tokens)?;
1945        let mut meta = self.ensure_field_meta(field)?;
1946        meta.live_docs = docs;
1947        meta.total_tokens = tokens;
1948        meta.term_stats_ready = true;
1949        self.put_field_meta(field, &meta)
1950    }
1951
1952    pub fn finish_text_build_accum(&mut self, field: u64,
1953                                   mut accumulator: TextBuildAccumulator) -> Result<()> {
1954        if accumulator.sort.is_none() {
1955            let mut grouped: Vec<_> = accumulator.grouped.drain().collect();
1956            grouped.sort_unstable_by(|a, b| a.0.cmp(&b.0));
1957            let mut lex_keys = Vec::with_capacity(grouped.len());
1958            for (key, count) in grouped {
1959                if key.len() < 9 || count == 0 {
1960                    return Err(corrupt("text build accumulator contains an invalid term record"));
1961                }
1962                let raw = u64::from_be_bytes(key[..8].try_into().unwrap());
1963                let term = &key[8..];
1964                if term.is_empty() || std::str::from_utf8(term).is_err() {
1965                    return Err(corrupt("text build accumulator contains an invalid word"));
1966                }
1967                let actual = self.reserve_built_term(field, raw, term)?;
1968                self.publish_built_term_count(field, actual, count)?;
1969                if actual != raw { self.rewrite_built_norm_terms(field, term, raw, actual)?; }
1970                lex_keys.push(keys::text_term_lex_key(field, term));
1971            }
1972            lex_keys.sort_unstable();
1973            for batch in lex_keys.chunks(64) { self.store().put_empty_batch(batch)?; }
1974            self.put_head_build_meta(field, accumulator.docs, accumulator.tokens)?;
1975            let mut meta = self.ensure_field_meta(field)?;
1976            meta.live_docs = accumulator.docs;
1977            meta.total_tokens = accumulator.tokens;
1978            meta.term_stats_ready = true;
1979            return self.put_field_meta(field, &meta);
1980        }
1981        accumulator.flush()?;
1982        let TextBuildAccumulator { sort, docs, tokens, .. } = accumulator;
1983        let sort = sort.unwrap();
1984        let mut runs = sort.finish()?;
1985        let mut iter = runs.iter()?;
1986        let mut group: Option<(u64, Vec<u8>, u64, u64)> = None; // raw, term, actual, count
1987        let mut lex_batch = Vec::with_capacity(64);
1988        while let Some(item) = iter.next() {
1989            let (key, value, _) = item?;
1990            if key.len() < 9 || value.len() != 8 {
1991                return Err(corrupt("text build sort emitted an invalid term record"));
1992            }
1993            let raw = u64::from_be_bytes(key[..8].try_into().unwrap());
1994            let term = &key[8..];
1995            if term.is_empty() || std::str::from_utf8(term).is_err() {
1996                return Err(corrupt("text build sort emitted an invalid word"));
1997            }
1998            let occurrences = u64::from_be_bytes(value.try_into().unwrap());
1999            if occurrences == 0 { return Err(corrupt("text build sort emitted a zero count")); }
2000            let changed = group.as_ref().is_some_and(|(old_raw, old_term, _, _)|
2001                *old_raw != raw || old_term.as_slice() != term);
2002            if changed {
2003                let (old_raw, old_term, actual, count) = group.take().unwrap();
2004                self.publish_built_term_count(field, actual, count)?;
2005                if actual != old_raw {
2006                    self.rewrite_built_norm_terms(field, &old_term, old_raw, actual)?;
2007                }
2008                lex_batch.push(keys::text_term_lex_key(field, &old_term));
2009                if lex_batch.len() == 64 {
2010                    lex_batch.sort_unstable();
2011                    self.store().put_empty_batch(&lex_batch)?;
2012                    lex_batch.clear();
2013                }
2014            }
2015            if group.is_none() {
2016                let actual = self.reserve_built_term(field, raw, term)?;
2017                group = Some((raw, term.to_vec(), actual, 0));
2018            }
2019            let (_, _, _, count) = group.as_mut().unwrap();
2020            *count = count.checked_add(occurrences).ok_or(Error::TooLarge)?;
2021        }
2022        if let Some((raw, term, actual, count)) = group {
2023            self.publish_built_term_count(field, actual, count)?;
2024            if actual != raw { self.rewrite_built_norm_terms(field, &term, raw, actual)?; }
2025            lex_batch.push(keys::text_term_lex_key(field, &term));
2026        }
2027        if !lex_batch.is_empty() {
2028            lex_batch.sort_unstable();
2029            self.store().put_empty_batch(&lex_batch)?;
2030        }
2031        self.put_head_build_meta(field, docs, tokens)?;
2032        let mut meta = self.ensure_field_meta(field)?;
2033        meta.live_docs = docs; meta.total_tokens = tokens; meta.term_stats_ready = true;
2034        self.put_field_meta(field, &meta)
2035    }
2036
2037    fn reserve_built_term(&mut self, field: u64, raw: u64, term: &[u8]) -> Result<u64> {
2038        let mut id = raw;
2039        for _ in 0..1024 {
2040            let key = keys::text_term_id_key(field, id);
2041            match self.store_ref().get(&key)? {
2042                Some(v) => {
2043                    let (_, stored) = decode_term_row(&v)?;
2044                    if stored == term { return Ok(id); }
2045                    id = next_term_probe(id);
2046                }
2047                None => {
2048                    self.store().put(&key, &encode_term_row(term, 1))?;
2049                    return Ok(id);
2050                }
2051            }
2052        }
2053        Err(corrupt("built text term collision chain exceeds its bound"))
2054    }
2055
2056    fn rewrite_built_norm_term(&mut self, field: u64, docid: u64,
2057                               raw: u64, actual: u64) -> Result<()> {
2058        let key = keys::text_norm_key(field, docid);
2059        let v = self.store_ref().get(&key)?
2060            .ok_or_else(|| corrupt("colliding built term has no document norm"))?;
2061        let mut norm = Norm::decode(&v)?;
2062        let pos = norm.terms.binary_search_by_key(&raw, |&(id, _)| id)
2063            .map_err(|_| corrupt("colliding built term is absent from its document norm"))?;
2064        norm.terms[pos].0 = actual;
2065        norm.terms.sort_unstable_by_key(|&(id, _)| id);
2066        if norm.terms.windows(2).any(|pair| pair[0].0 >= pair[1].0) {
2067            return Err(corrupt("resolved built term numbers are not unique"));
2068        }
2069        self.store().put(&key, &norm.encode())
2070    }
2071
2072    /// A collision is discovered only after the external aggregation has
2073    /// grouped exact words.  The exact-word posting range is already present,
2074    /// so it supplies precisely the affected documents without retaining a
2075    /// corpus-sized side list during the build.
2076    fn rewrite_built_norm_terms(&mut self, field: u64, term: &[u8],
2077                                raw: u64, actual: u64) -> Result<()> {
2078        let mut prefix = keys::text_head_key(field, term, 0);
2079        prefix.truncate(prefix.len() - 8);
2080        let mut docs = Vec::new();
2081        self.store_ref().scan(&prefix)?.for_each_ref(|key, _| {
2082            if !key.starts_with(&prefix) { return false; }
2083            if key.len() != prefix.len() + 8 { return false; }
2084            docs.push(u64::from_be_bytes(key[prefix.len()..].try_into().unwrap()));
2085            true
2086        })?;
2087        for docid in docs { self.rewrite_built_norm_term(field, docid, raw, actual)?; }
2088        Ok(())
2089    }
2090
2091    fn publish_built_term_count(&mut self, field: u64, id: u64, count: u64) -> Result<()> {
2092        if count <= 1 { return Ok(()); }
2093        let key = keys::text_term_id_key(field, id);
2094        let v = self.store_ref().get(&key)?
2095            .ok_or_else(|| corrupt("built text term has no dictionary row"))?;
2096        let (_, term) = decode_term_row(&v)?;
2097        let term = term.to_vec();
2098        self.store().put(&key, &encode_term_row(&term, count))
2099    }
2100
2101    /// Replace one doc's text (the UPDATE path). The head is the mutable
2102    /// tier, so the old HEAD rows are removed physically -- dead-marking
2103    /// cannot express "these terms changed" -- while folded copies in
2104    /// segments are dead-marked as usual and dropped at the next fold.
2105    /// (Dead-mark-only replacement resurrected the doc's old head rows:
2106    /// "stale term still matches", caught by the ported e1 suite.)
2107    pub fn replace_text(&mut self, field: u64, docid: u64,
2108                        _old_text: Option<&str>, new_text: &str) -> Result<()> {
2109        let Some(v) = self.store_ref().get(&keys::text_norm_key(field, docid))? else {
2110            return self.index_text(field, docid, new_text);
2111        };
2112        self.replace_norm(field, docid, Norm::decode(&v)?, new_text)
2113    }
2114
2115    fn replace_norm(&mut self, field: u64, docid: u64, old: Norm, new_text: &str) -> Result<()> {
2116        let mut meta = self.ensure_field_meta(field)?;
2117        let old_owner = meta.resolve_owner(old.owner);
2118        let mut new_tf: std::collections::BTreeMap<String, u64> = std::collections::BTreeMap::new();
2119        for term in tokenize(new_text) { *new_tf.entry(term).or_insert(0) += 1; }
2120        let new_total: u64 = new_tf.values().copied().sum();
2121
2122        // Build the complete new head copy first.  Old segment/head rows remain
2123        // authoritative until the norm and owner metadata below move.
2124        for (term, tf) in &new_tf {
2125            let mut v = Vec::new(); write_varint(&mut v, *tf); write_varint(&mut v, new_total);
2126            self.store().put(&keys::text_head_key(field, term.as_bytes(), docid), &v)?;
2127        }
2128        if new_total == 0 {
2129            let mut dv = Vec::new(); write_varint(&mut dv, new_total);
2130            self.store().put(&keys::text_head_doc_key(field, docid), &dv)?;
2131        } else {
2132            self.store().delete(&keys::text_head_doc_key(field, docid))?;
2133        }
2134
2135        let old_ids: HashSet<u64> = old.terms.iter().map(|&(id, _)| id).collect();
2136        let mut new_norm_terms = Vec::with_capacity(new_tf.len());
2137        for (term, tf) in &new_tf {
2138            let id = match self.term_info(field, term)? {
2139                Some((_, id)) if old_ids.contains(&id) => id,
2140                Some(_) | None => self.adjust_term_df(field, term, 1)?,
2141            };
2142            new_norm_terms.push((id, *tf));
2143        }
2144        new_norm_terms.sort_unstable_by_key(|&(id, _)| id);
2145        let new_ids: HashSet<u64> = new_norm_terms.iter().map(|&(id, _)| id).collect();
2146        let mut removed_terms = Vec::new();
2147        for &(id, _) in old.terms.iter().filter(|(id, _)| !new_ids.contains(id)) {
2148            let term = self.term_by_id(field, id)?;
2149            self.adjust_term_df(field, &term, -1)?;
2150            removed_terms.push(term);
2151        }
2152
2153        let head_key = keys::text_meta_key(field, 0);
2154        let mut head = self.store_ref().get(&head_key)?.map(|v| SegMeta::decode(&v)).transpose()?
2155            .unwrap_or(SegMeta { doc_count: 0, total_tokens: 0, dead_tokens: 0, dead: Vec::new(),
2156                level: 0, term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0 });
2157        if old_owner == 0 {
2158            head.total_tokens = head.total_tokens.checked_sub(old.total)
2159                .and_then(|n| n.checked_add(new_total)).ok_or_else(|| corrupt("head token accounting overflows"))?;
2160            for term in &removed_terms {
2161                self.store().delete(&keys::text_head_key(field, term.as_bytes(), docid))?;
2162            }
2163        } else {
2164            let mut owner = self.segment_meta(field, old_owner)?;
2165            if !owner.dead.contains(&docid) {
2166                owner.dead.push(docid);
2167                owner.dead_tokens = owner.dead_tokens.checked_add(old.total).ok_or(Error::TooLarge)?;
2168                self.store().put(&keys::text_meta_key(field, old_owner), &owner.encode())?;
2169            }
2170            head.doc_count = head.doc_count.checked_add(1).ok_or(Error::TooLarge)?;
2171            head.total_tokens = head.total_tokens.checked_add(new_total).ok_or(Error::TooLarge)?;
2172        }
2173        head.dead.retain(|&id| id != docid);
2174        self.store().put(&head_key, &head.encode())?;
2175        meta.total_tokens = meta.total_tokens.checked_sub(old.total)
2176            .and_then(|n| n.checked_add(new_total)).ok_or_else(|| corrupt("field token accounting overflows"))?;
2177        let norm = Norm { total: new_total, owner: 0, terms: new_norm_terms };
2178        self.store().put(&keys::text_norm_key(field, docid), &norm.encode())?;
2179        self.put_field_meta(field, &meta)
2180    }
2181
2182    /// The deletion discipline (alive/dead bitmap, one value per segment):
2183    /// mark `docid` dead in the metadata row that owns its corpus counts.
2184    /// Postings are NOT touched -- folds drop dead docs physically (Law 3 shape).
2185    pub fn delete_text(&mut self, field: u64, docid: u64) -> Result<bool> {
2186        let Some(v) = self.store_ref().get(&keys::text_norm_key(field, docid))? else { return Ok(false); };
2187        let norm = Norm::decode(&v)?;
2188        let mut field_meta = self.ensure_field_meta(field)?;
2189        let owner = field_meta.resolve_owner(norm.owner);
2190        let mk = keys::text_meta_key(field, owner);
2191        let mut segment = self.segment_meta(field, owner)?;
2192        if segment.dead.contains(&docid) { return Ok(false); }
2193        segment.dead.push(docid);
2194        segment.dead_tokens = segment.dead_tokens.checked_add(norm.total).ok_or(Error::TooLarge)?;
2195        self.store().put(&mk, &segment.encode())?;
2196        if owner == 0 {
2197            // Ordinary non-empty head documents need no membership row: the
2198            // merge derives it from their postings.  A deletion keeps this
2199            // one small length row so a later resurrection can balance the
2200            // head's token totals after the norm is retired.
2201            let mut dv = Vec::new(); write_varint(&mut dv, norm.total);
2202            self.store().put(&keys::text_head_doc_key(field, docid), &dv)?;
2203        }
2204        for &(id, _) in &norm.terms {
2205            let term = self.term_by_id(field, id)?;
2206            self.adjust_term_df(field, &term, -1)?;
2207        }
2208        field_meta.live_docs = field_meta.live_docs.checked_sub(1)
2209            .ok_or_else(|| corrupt("text field document count underflows"))?;
2210        field_meta.total_tokens = field_meta.total_tokens.checked_sub(norm.total)
2211            .ok_or_else(|| corrupt("text field token count underflows"))?;
2212        self.put_field_meta(field, &field_meta)?;
2213        self.store().delete(&keys::text_norm_key(field, docid))?;
2214        Ok(true)
2215    }
2216
2217    fn text_posting_cursor(&self, field: u64, term: &str) -> Result<Option<TextPostingCursor<'_>>> {
2218        let Some(manifest) = self.field_meta(field)? else { return Ok(None) };
2219        let mut sources = Vec::new();
2220        for (seg, _) in manifest.active {
2221            let meta = self.segment_meta(field, seg)?;
2222            let dead = meta.dead.into_iter().collect();
2223            let exact_key = keys::text_seg_key(field, seg, term.as_bytes());
2224            let exact = self.store_ref().get(&exact_key)?;
2225            let mut pending = VecDeque::new();
2226            let mut blocks_read = 0;
2227            let mut postings_decoded = 0;
2228            let scan_more = if let Some(value) = exact {
2229                let posts = decode_postings(&value)?;
2230                blocks_read = 1;
2231                postings_decoded = posts.len() as u64;
2232                let more = posts.len() == POSTING_BLOCK;
2233                pending.extend(posts);
2234                more
2235            } else {
2236                true
2237            };
2238            let prefix = keys::text_seg_block_prefix(field, seg, term.as_bytes());
2239            let scan = scan_more.then(|| self.store_ref().scan(&prefix)).transpose()?;
2240            sources.push(PostingSource {
2241                scan, prefix, pending, dead, head: false, blocks_read,
2242                postings_decoded,
2243            });
2244        }
2245        if manifest.head_redirect.is_none() {
2246            if let Some(value) = self.store_ref().get(&keys::text_meta_key(field, 0))? {
2247                let meta = SegMeta::decode(&value)?;
2248                let mut prefix = keys::text_head_key(field, term.as_bytes(), 0);
2249                prefix.truncate(prefix.len() - 8);
2250                sources.push(PostingSource {
2251                    scan: Some(self.store_ref().scan(&prefix)?),
2252                    prefix,
2253                    pending: VecDeque::new(),
2254                    dead: meta.dead.into_iter().collect(),
2255                    head: true,
2256                    blocks_read: 0,
2257                    postings_decoded: 0,
2258                });
2259            }
2260        }
2261        let mut cursor = TextPostingCursor {
2262            heads: vec![None; sources.len()],
2263            sources,
2264        };
2265        for i in 0..cursor.sources.len() {
2266            cursor.heads[i] = cursor.sources[i].next_live()?;
2267        }
2268        Ok(Some(cursor))
2269    }
2270
2271    /// Count the union of exact query-term postings with one bounded block per
2272    /// active segment and term. This is the COUNT counterpart of ranked top-k:
2273    /// it never owns the matching document set.
2274    pub fn text_match_count(&self, field: u64, query: &str) -> Result<Option<u64>> {
2275        let diag = std::env::var_os("TEXT_DIAG").is_some();
2276        let started = std::time::Instant::now();
2277        let pages_before = self.store_ref().pool_stats();
2278        let mut terms = tokenize(query);
2279        terms.sort();
2280        terms.dedup();
2281        if terms.is_empty() { return Ok(Some(0)) }
2282        let mut cursors = Vec::with_capacity(terms.len());
2283        for term in &terms {
2284            let Some(cursor) = self.text_posting_cursor(field, term)? else { return Ok(None) };
2285            cursors.push(cursor);
2286        }
2287        let mut heads = Vec::with_capacity(cursors.len());
2288        for cursor in &mut cursors { heads.push(cursor.next()?); }
2289        let mut count = 0u64;
2290        while let Some(doc) = heads.iter().flatten().map(|posting| posting.0).min() {
2291            count = count.checked_add(1).ok_or(Error::TooLarge)?;
2292            for i in 0..cursors.len() {
2293                if heads[i].is_some_and(|posting| posting.0 == doc) {
2294                    heads[i] = cursors[i].next()?;
2295                }
2296            }
2297        }
2298        if diag {
2299            let (blocks, postings) = cursors.iter().fold((0u64, 0u64), |acc, cursor| {
2300                let current = cursor.counters();
2301                (acc.0 + current.0, acc.1 + current.1)
2302            });
2303            let after = self.store_ref().pool_stats();
2304            eprintln!(
2305                "TEXT_DIAG streaming_count={query:?} total_ms={:.3} blocks_read={} postings_scanned={} matched={} pool_logical={} pool_misses={}",
2306                started.elapsed().as_secs_f64() * 1000.0, blocks, postings, count,
2307                after.hits.saturating_add(after.misses)
2308                    .saturating_sub(pages_before.hits.saturating_add(pages_before.misses)),
2309                after.misses.saturating_sub(pages_before.misses),
2310            );
2311        }
2312        Ok(Some(count))
2313    }
2314
2315    /// First `limit` BM25 matches in document-id order. Collection membership
2316    /// uses that same order, so an unordered SQL filter with OFFSET/LIMIT can
2317    /// stop its posting merge without changing which rows the old collection
2318    /// scan returned. Scores are computed while the term heads are resident.
2319    pub fn text_match_doc_limit(&self, field: u64, query: &str,
2320                                min_score: f64, limit: usize)
2321        -> Result<Option<Vec<(u64, f64)>>>
2322    {
2323        let diag = std::env::var_os("TEXT_DIAG").is_some();
2324        let started = std::time::Instant::now();
2325        let pages_before = self.store_ref().pool_stats();
2326        if !self.field_meta(field)?.is_some_and(|meta| meta.term_stats_ready) {
2327            return Ok(None);
2328        }
2329        if limit == 0 { return Ok(Some(Vec::new())) }
2330        let mut terms = tokenize(query);
2331        terms.sort();
2332        terms.dedup();
2333        if terms.is_empty() { return Ok(Some(Vec::new())) }
2334        let (n_docs, avg_len) = self.text_stats(field)?;
2335        let mut cursors = Vec::new();
2336        let mut idfs = Vec::new();
2337        for term in &terms {
2338            let Some((df, _)) = self.term_info(field, term)? else { continue };
2339            if df == 0 { continue }
2340            let Some(cursor) = self.text_posting_cursor(field, term)? else { return Ok(None) };
2341            cursors.push(cursor);
2342            let df = df as f64;
2343            idfs.push(((n_docs - df + 0.5) / (df + 0.5) + 1.0).ln());
2344        }
2345        let mut heads = Vec::with_capacity(cursors.len());
2346        for cursor in &mut cursors { heads.push(cursor.next()?); }
2347        let mut out = Vec::with_capacity(limit.min(4096));
2348        let mut matches_scored = 0u64;
2349        while out.len() < limit {
2350            let Some(doc) = heads.iter().flatten().map(|posting| posting.0).min()
2351            else { break };
2352            let mut score = 0.0;
2353            for i in 0..cursors.len() {
2354                if let Some((posting_doc, tf, dl)) = heads[i] {
2355                    if posting_doc == doc {
2356                        let tf = tf as f64;
2357                        let dl = dl as f64;
2358                        let denom = tf + BM25_K1 as f64
2359                            * (1.0 - BM25_B as f64
2360                                + BM25_B as f64 * dl / avg_len.max(1.0));
2361                        score += idfs[i] * tf * (BM25_K1 as f64 + 1.0) / denom;
2362                        heads[i] = cursors[i].next()?;
2363                    }
2364                }
2365            }
2366            matches_scored += 1;
2367            if score > min_score { out.push((doc, score)); }
2368        }
2369        if diag {
2370            let (blocks, postings) = cursors.iter().fold((0u64, 0u64), |acc, cursor| {
2371                let current = cursor.counters();
2372                (acc.0 + current.0, acc.1 + current.1)
2373            });
2374            let after = self.store_ref().pool_stats();
2375            eprintln!(
2376                "TEXT_DIAG streaming_doc_limit={query:?} total_ms={:.3} blocks_read={} postings_scanned={} matches_scored={} returned={} pool_misses={}",
2377                started.elapsed().as_secs_f64() * 1000.0, blocks, postings,
2378                matches_scored, out.len(),
2379                after.misses.saturating_sub(pages_before.misses),
2380            );
2381        }
2382        Ok(Some(out))
2383    }
2384
2385    /// All (docid, tf) postings for one exact term in `field`, across the
2386    /// head rows and every folded segment, minus dead docs. O(matches).
2387    pub fn text_postings(&self, field: u64, term: &str) -> Result<Vec<(u64, u64, u64)>> {
2388        let diag = std::env::var_os("TEXT_DIAG").is_some();
2389        let started = std::time::Instant::now();
2390        let pages_before = self.store_ref().pool_stats();
2391        let mut blocks_read = 0u64;
2392        let mut postings_decoded = 0u64;
2393        let manifest = self.field_meta(field)?;
2394        let segs: Vec<u32> = match &manifest {
2395            Some(m) => m.active.iter().map(|&(seg, _)| seg).collect(),
2396            None => self.legacy_segment_metas(field)?.into_iter()
2397                .filter_map(|(seg, _)| (seg != 0).then_some(seg)).collect(),
2398        };
2399        let mut out: Vec<(u64, u64, u64)> = Vec::new();
2400        for seg in segs {
2401            let meta = self.segment_meta(field, seg)?;
2402            let dead: HashSet<u64> = meta.dead.into_iter().collect();
2403
2404            // The first bounded block keeps the historical exact-term key.
2405            // A rare term is therefore one point read; only a full first block
2406            // can have continuation rows and needs a prefix cursor.
2407            let exact = self.store_ref().get(&keys::text_seg_key(field, seg, term.as_bytes()))?;
2408            let scan_more = if let Some(v) = exact {
2409                blocks_read += 1;
2410                let posts = decode_postings(&v)?;
2411                postings_decoded += posts.len() as u64;
2412                let more = posts.len() == POSTING_BLOCK;
2413                for (doc, tf, dl) in posts {
2414                    if !dead.contains(&doc) { out.push((doc, tf, dl)); }
2415                }
2416                more
2417            } else { true }; // compatibility with the first block layout used during migration
2418            if scan_more {
2419                let prefix = keys::text_seg_block_prefix(field, seg, term.as_bytes());
2420                let mut failure = None;
2421                self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
2422                    if !key.starts_with(&prefix) { return false; }
2423                    if key.len() != prefix.len() + 8 { return false; }
2424                    match decode_postings(val) {
2425                        Ok(posts) => {
2426                            blocks_read += 1;
2427                            postings_decoded += posts.len() as u64;
2428                            for (doc, tf, dl) in posts {
2429                                if !dead.contains(&doc) { out.push((doc, tf, dl)); }
2430                            }
2431                        },
2432                        Err(e) => { failure = Some(e); return false; }
2433                    }
2434                    true
2435                })?;
2436                if let Some(e) = failure { return Err(e); }
2437            }
2438        }
2439
2440        // A published head fold sets head_redirect until the old head range is
2441        // retired, so an interrupted merge never exposes both copies.
2442        let mut head = Vec::new();
2443        let head_meta = if manifest.as_ref().is_none_or(|m| m.head_redirect.is_none()) {
2444            self.store_ref().get(&keys::text_meta_key(field, 0))?
2445                .map(|v| SegMeta::decode(&v)).transpose()?
2446        } else { None };
2447        if let Some(head_meta) = head_meta {
2448            let dead: HashSet<u64> = head_meta.dead.into_iter().collect();
2449            let mut prefix = keys::text_head_key(field, term.as_bytes(), 0);
2450            prefix.truncate(prefix.len() - 8);
2451            let mut failure = None;
2452            self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
2453                if !key.starts_with(&prefix) { return false; }
2454                if key.len() != prefix.len() + 8 { return false; }
2455                let doc = u64::from_be_bytes(key[prefix.len()..].try_into().unwrap());
2456                let mut pos = 0;
2457                let decoded = required_varint(val, &mut pos, "head posting frequency is truncated")
2458                    .and_then(|tf| required_varint(val, &mut pos, "head posting length is truncated")
2459                        .map(|dl| (tf, dl)));
2460                match decoded {
2461                    Ok((tf, dl)) if pos == val.len() && tf > 0 => {
2462                        blocks_read += 1;
2463                        postings_decoded += 1;
2464                        if !dead.contains(&doc) { head.push((doc, tf, dl)); }
2465                    }
2466                    Ok(_) => { failure = Some(corrupt("head posting has invalid trailing bytes or frequency")); return false; }
2467                    Err(e) => { failure = Some(e); return false; }
2468                }
2469                true
2470            })?;
2471            if let Some(e) = failure { return Err(e); }
2472        }
2473        if out.is_empty() {
2474            if diag {
2475                let after = self.store_ref().pool_stats();
2476                eprintln!(
2477                    "TEXT_DIAG term={term:?} posting_ms={:.3} blocks_read={} postings_decoded={} live_postings={} pool_misses={}",
2478                    started.elapsed().as_secs_f64() * 1000.0, blocks_read,
2479                    postings_decoded, head.len(), after.misses.saturating_sub(pages_before.misses),
2480                );
2481            }
2482            return Ok(head);
2483        }
2484        if !head.is_empty() {
2485            let head_ids: HashSet<u64> = head.iter().map(|&(doc, _, _)| doc).collect();
2486            out.retain(|(doc, _, _)| !head_ids.contains(doc));
2487            out.extend(head);
2488        }
2489        out.sort_unstable_by_key(|&(doc, _, _)| doc);
2490        out.dedup_by_key(|posting| posting.0);
2491        if diag {
2492            let after = self.store_ref().pool_stats();
2493            eprintln!(
2494                "TEXT_DIAG term={term:?} posting_ms={:.3} blocks_read={} postings_decoded={} live_postings={} pool_misses={}",
2495                started.elapsed().as_secs_f64() * 1000.0, blocks_read,
2496                postings_decoded, out.len(), after.misses.saturating_sub(pages_before.misses),
2497            );
2498        }
2499        Ok(out)
2500    }
2501
2502    /// BM25 statistics for `field`: (doc_count, avg_len) across segments,
2503    /// dead docs excluded from the count.
2504    fn text_stats(&self, field: u64) -> Result<(f64, f64)> {
2505        if let Some(meta) = self.field_meta(field)? {
2506            let docs = meta.live_docs.max(1);
2507            return Ok((docs as f64, meta.total_tokens as f64 / docs as f64));
2508        }
2509        let mut docs = 0u64; let mut tokens = 0u64;
2510        let mut dead = 0u64; let mut dead_tokens = 0u64;
2511        for seg in self.text_segments(field)? {
2512            if let Some(v) = self.store_ref().get(&keys::text_meta_key(field, seg))? {
2513                let m = SegMeta::decode(&v)?;
2514                docs += m.doc_count;
2515                tokens += m.total_tokens;
2516                dead += m.dead.len() as u64;
2517                dead_tokens += m.dead_tokens;
2518            }
2519        }
2520        let live = docs.saturating_sub(dead).max(1);
2521        let live_tokens = tokens.saturating_sub(dead_tokens);
2522        Ok((live as f64, live_tokens as f64 / live as f64))
2523    }
2524
2525    /// Incremental live corpus counters used by BM25.  This is a point read.
2526    pub fn text_live_stats(&self, field: u64) -> Result<(u64, u64)> {
2527        if let Some(meta) = self.field_meta(field)? { return Ok((meta.live_docs, meta.total_tokens)); }
2528        let (docs, avg) = self.text_stats(field)?;
2529        Ok((docs as u64, (docs * avg).round() as u64))
2530    }
2531
2532    /// Slow diagnostic oracle: recount current per-document norm rows from
2533    /// scratch.  Production scoring never calls it.
2534    pub fn text_recount_stats(&self, field: u64) -> Result<(u64, u64)> {
2535        let prefix = keys::text_norm_prefix(field);
2536        let mut docs = 0u64; let mut tokens = 0u64;
2537        let mut failure = None;
2538        self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
2539            if !key.starts_with(&prefix) { return false; }
2540            match Norm::decode(val) {
2541                Ok(norm) => {
2542                    docs = match docs.checked_add(1) { Some(n) => n, None => { failure = Some(Error::TooLarge); return false; } };
2543                    tokens = match tokens.checked_add(norm.total) { Some(n) => n, None => { failure = Some(Error::TooLarge); return false; } };
2544                }
2545                Err(e) => { failure = Some(e); return false; }
2546            }
2547            true
2548        })?;
2549        if let Some(e) = failure { return Err(e); }
2550        Ok((docs, tokens))
2551    }
2552
2553    /// Slow diagnostic oracle for one word's live-document count.  It scans
2554    /// per-document norms from scratch; production ranking uses the dictionary
2555    /// point row maintained by writes instead.
2556    pub fn text_recount_term_doc_freq(&self, field: u64, term: &str) -> Result<u64> {
2557        let Some((_, id)) = self.term_info(field, term)? else { return Ok(0); };
2558        let prefix = keys::text_norm_prefix(field);
2559        let mut docs = 0u64; let mut failure = None;
2560        self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
2561            if !key.starts_with(&prefix) { return false; }
2562            match Norm::decode(val) {
2563                Ok(norm) if norm.terms.binary_search_by_key(&id, |&(term_id, _)| term_id).is_ok() => {
2564                    docs = match docs.checked_add(1) {
2565                        Some(n) => n,
2566                        None => { failure = Some(Error::TooLarge); return false; }
2567                    };
2568                }
2569                Ok(_) => {}
2570                Err(e) => { failure = Some(e); return false; }
2571            }
2572            true
2573        })?;
2574        if let Some(e) = failure { return Err(e); }
2575        Ok(docs)
2576    }
2577
2578    pub fn text_term_doc_freq(&self, field: u64, term: &str) -> Result<Option<u64>> {
2579        if self.field_meta(field)?.is_some_and(|meta| !meta.term_stats_ready) { return Ok(None); }
2580        Ok(self.term_info(field, term)?.map(|(count, _)| count))
2581    }
2582
2583    /// BM25 for an already chosen candidate slice.  Current-format stores do
2584    /// one field-stat point read, one term-frequency point read per distinct
2585    /// query term, and one norm point read per candidate.  No posting extent
2586    /// is opened, so ten candidates cost ten document reads whether the term
2587    /// occurs in one hundred or ten million documents.
2588    pub fn text_score_candidates(&self, field: u64, query: &str, cands: &[u64]) -> Result<Vec<f32>> {
2589        let mut terms = tokenize(query);
2590        terms.sort(); terms.dedup();
2591        let mut out = vec![0.0f32; cands.len()];
2592        if terms.is_empty() || cands.is_empty() { return Ok(out); }
2593        let (n_docs, avg_len) = self.text_stats(field)?;
2594        let stats_ready = !self.field_meta(field)?.is_some_and(|meta| !meta.term_stats_ready);
2595        let mut scored_terms = Vec::with_capacity(terms.len());
2596        for term in terms {
2597            let info = if stats_ready { self.term_info(field, &term)? } else { None };
2598            let df = match info {
2599                Some((df, _)) => df,
2600                None => self.text_postings(field, &term)?.len() as u64, // legacy compatibility only
2601            } as f64;
2602            if df > 0.0 {
2603                scored_terms.push((term, info.map(|(_, id)| id),
2604                    ((n_docs - df + 0.5) / (df + 0.5) + 1.0).ln()));
2605            }
2606        }
2607        let mut norms: Vec<Option<Norm>> = (0..cands.len()).map(|_| None).collect();
2608        if cands.len().saturating_mul(16) >= n_docs as usize {
2609            // Dense candidate batches are still candidate-proportional when a
2610            // sequential norm walk visits at most sixteen rows per candidate.
2611            // This avoids N random descents without ever falling back to the
2612            // term's posting union.
2613            let index: HashMap<u64, usize> = cands.iter().enumerate().map(|(i, &id)| (id, i)).collect();
2614            let prefix = keys::text_norm_prefix(field);
2615            let mut failure = None;
2616            self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
2617                if !key.starts_with(&prefix) { return false; }
2618                if key.len() != prefix.len() + 8 { failure = Some(corrupt("text norm key has an invalid length")); return false; }
2619                let doc = u64::from_be_bytes(key[prefix.len()..].try_into().unwrap());
2620                if let Some(&i) = index.get(&doc) {
2621                    match Norm::decode(val) { Ok(norm) => norms[i] = Some(norm), Err(e) => { failure = Some(e); return false; } }
2622                }
2623                true
2624            })?;
2625            if let Some(e) = failure { return Err(e); }
2626        } else {
2627            for (i, &docid) in cands.iter().enumerate() {
2628                if let Some(v) = self.store_ref().get(&keys::text_norm_key(field, docid))? {
2629                    norms[i] = Some(Norm::decode(&v)?);
2630                }
2631            }
2632        }
2633        for (i, norm) in norms.into_iter().enumerate() {
2634            let Some(norm) = norm else { continue; };
2635            let docid = cands[i];
2636            if norm.terms.is_empty() || scored_terms.iter().any(|(_, id, _)| id.is_none()) {
2637                // Old norms did not carry frequencies; preserve readability
2638                // without putting this database-sized fallback on new data.
2639                for (term, _, term_idf) in &scored_terms {
2640                    if let Some((_, tf, dl)) = self.text_postings(field, term)?.into_iter().find(|p| p.0 == docid) {
2641                        let (tf, dl) = (tf as f64, dl as f64);
2642                        let denom = tf + BM25_K1 as f64 *
2643                            (1.0 - BM25_B as f64 + BM25_B as f64 * dl / avg_len.max(1.0));
2644                        out[i] += (*term_idf * tf
2645                            * (BM25_K1 as f64 + 1.0) / denom) as f32;
2646                    }
2647                }
2648                continue;
2649            }
2650            for (_, term_id, term_idf) in &scored_terms {
2651                let Some(term_id) = term_id else { continue; };
2652                let Ok(pos) = norm.terms.binary_search_by_key(term_id, |&(id, _)| id) else { continue; };
2653                let tf = norm.terms[pos].1;
2654                let (tf, dl) = (tf as f64, norm.total as f64);
2655                let denom = tf + BM25_K1 as f64 *
2656                    (1.0 - BM25_B as f64 + BM25_B as f64 * dl / avg_len.max(1.0));
2657                out[i] += (*term_idf * tf
2658                    * (BM25_K1 as f64 + 1.0) / denom) as f32;
2659            }
2660        }
2661        Ok(out)
2662    }
2663
2664    /// BM25 top-k for a multi-term query in one field. Current stores merge
2665    /// packed posting cursors in document order and retain only a k-entry heap.
2666    /// Each cursor owns at most one decoded 128-document block. Legacy stores
2667    /// without maintained term statistics use the materialising compatibility
2668    /// implementation below until they are rebuilt.
2669    pub fn text_search(&self, field: u64, query: &str, k: usize) -> Result<Vec<(u64, f64)>> {
2670        if self.field_meta(field)?.is_some_and(|meta| meta.term_stats_ready) {
2671            return self.text_search_streaming(field, query, k);
2672        }
2673        let diag = std::env::var_os("TEXT_DIAG").is_some();
2674        let total_started = std::time::Instant::now();
2675        let pages_before = self.store_ref().pool_stats();
2676        let terms = tokenize(query);
2677        if terms.is_empty() { return Ok(Vec::new()); }
2678        let stats_started = std::time::Instant::now();
2679        let (n_docs, avg_len) = self.text_stats(field)?;
2680        let stats_elapsed = stats_started.elapsed();
2681        let mut acc: std::collections::HashMap<u64, f64> = std::collections::HashMap::new();
2682        let mut seen: std::collections::HashSet<&String> = std::collections::HashSet::new();
2683        let mut posting_elapsed = std::time::Duration::ZERO;
2684        let mut scoring_elapsed = std::time::Duration::ZERO;
2685        let mut postings_scanned = 0usize;
2686        for term in &terms {
2687            if !seen.insert(term) { continue; } // repeated query words count once
2688            let posting_started = std::time::Instant::now();
2689            let posts = self.text_postings(field, term)?;
2690            posting_elapsed += posting_started.elapsed();
2691            if posts.is_empty() { continue; }
2692            postings_scanned += posts.len();
2693            let df = posts.len() as f64;
2694            let idf = ((n_docs - df + 0.5) / (df + 0.5) + 1.0).ln();
2695            let scoring_started = std::time::Instant::now();
2696            for (docid, tf, dl) in posts {
2697                let dl = dl as f64;
2698                let tf = tf as f64;
2699                let denom = tf + (BM25_K1 as f64) * (1.0 - BM25_B as f64
2700                    + (BM25_B as f64) * dl / avg_len.max(1.0));
2701                *acc.entry(docid).or_insert(0.0) += idf * tf * (BM25_K1 as f64 + 1.0) / denom;
2702            }
2703            scoring_elapsed += scoring_started.elapsed();
2704        }
2705        let matched = acc.len();
2706        let rank_started = std::time::Instant::now();
2707        let mut ranked: Vec<(u64, f64)> = acc.into_iter().collect();
2708        ranked.sort_by(|a, b| b.1.total_cmp(&a.1).then(a.0.cmp(&b.0)));
2709        ranked.truncate(k);
2710        let rank_elapsed = rank_started.elapsed();
2711        if diag {
2712            let after = self.store_ref().pool_stats();
2713            eprintln!(
2714                "TEXT_DIAG search={query:?} total_ms={:.3} stats_ms={:.3} postings_ms={:.3} scoring_ms={:.3} ranking_ms={:.3} postings_scanned={} matched={} returned={} pool_misses={}",
2715                total_started.elapsed().as_secs_f64() * 1000.0,
2716                stats_elapsed.as_secs_f64() * 1000.0,
2717                posting_elapsed.as_secs_f64() * 1000.0,
2718                scoring_elapsed.as_secs_f64() * 1000.0,
2719                rank_elapsed.as_secs_f64() * 1000.0,
2720                postings_scanned, matched, ranked.len(),
2721                after.misses.saturating_sub(pages_before.misses),
2722            );
2723        }
2724        Ok(ranked)
2725    }
2726
2727    fn text_search_streaming(&self, field: u64, query: &str, k: usize)
2728        -> Result<Vec<(u64, f64)>>
2729    {
2730        let diag = std::env::var_os("TEXT_DIAG").is_some();
2731        let started = std::time::Instant::now();
2732        let pages_before = self.store_ref().pool_stats();
2733        let mut terms = tokenize(query);
2734        terms.sort();
2735        terms.dedup();
2736        if terms.is_empty() || k == 0 { return Ok(Vec::new()) }
2737        let (n_docs, avg_len) = self.text_stats(field)?;
2738        let mut cursors = Vec::new();
2739        let mut idfs = Vec::new();
2740        for term in &terms {
2741            let Some((df, _)) = self.term_info(field, term)? else { continue };
2742            if df == 0 { continue }
2743            let Some(cursor) = self.text_posting_cursor(field, term)? else {
2744                return Err(corrupt("current text manifest disappeared during search"));
2745            };
2746            cursors.push(cursor);
2747            let df = df as f64;
2748            idfs.push(((n_docs - df + 0.5) / (df + 0.5) + 1.0).ln());
2749        }
2750        let mut heads = Vec::with_capacity(cursors.len());
2751        for cursor in &mut cursors { heads.push(cursor.next()?); }
2752        let mut heap = BinaryHeap::new();
2753        let mut matches = 0u64;
2754        while let Some(doc) = heads.iter().flatten().map(|posting| posting.0).min() {
2755            let mut score = 0.0;
2756            for i in 0..cursors.len() {
2757                if let Some((posting_doc, tf, dl)) = heads[i] {
2758                    if posting_doc == doc {
2759                        let tf = tf as f64;
2760                        let dl = dl as f64;
2761                        let denom = tf + BM25_K1 as f64 *
2762                            (1.0 - BM25_B as f64 + BM25_B as f64 * dl / avg_len.max(1.0));
2763                        score += idfs[i] * tf * (BM25_K1 as f64 + 1.0) / denom;
2764                        heads[i] = cursors[i].next()?;
2765                    }
2766                }
2767            }
2768            matches += 1;
2769            let hit = RankedText { id: doc, score };
2770            if heap.len() < k {
2771                heap.push(hit);
2772            } else {
2773                let worst = heap.peek().unwrap();
2774                if score > worst.score || (score == worst.score && doc < worst.id) {
2775                    *heap.peek_mut().unwrap() = hit;
2776                }
2777            }
2778        }
2779        let mut ranked: Vec<(u64, f64)> = heap.into_iter()
2780            .map(|hit| (hit.id, hit.score)).collect();
2781        ranked.sort_by(|a, b| b.1.total_cmp(&a.1).then(a.0.cmp(&b.0)));
2782        if diag {
2783            let (blocks, postings) = cursors.iter().fold((0u64, 0u64), |acc, cursor| {
2784                let current = cursor.counters();
2785                (acc.0 + current.0, acc.1 + current.1)
2786            });
2787            let after = self.store_ref().pool_stats();
2788            eprintln!(
2789                "TEXT_DIAG streaming_search={query:?} total_ms={:.3} blocks_read={} postings_scanned={} matched={} returned={} pool_misses={}",
2790                started.elapsed().as_secs_f64() * 1000.0, blocks, postings,
2791                matches, ranked.len(), after.misses.saturating_sub(pages_before.misses),
2792            );
2793        }
2794        Ok(ranked)
2795    }
2796
2797    /// Fold the mutable head, then run size-tiered levels.  At most seven
2798    /// segments remain at any level: the eighth is k-way merged into the next
2799    /// level.  Every candidate is checkpointed and independently reopened and
2800    /// counted before the one-row active manifest can publish it.
2801    pub fn fold_text(&mut self, field: u64) -> Result<()> {
2802        self.fold_text_with_before_publish(field, |_, _| Ok(()))
2803    }
2804
2805    fn fold_text_with_before_publish<F>(&mut self, field: u64, before_publish: F) -> Result<()>
2806    where F: FnOnce(&mut Graph, u32) -> Result<()> {
2807        if self.store_ref().get(&keys::text_meta_key(field, 0))?.is_none() { return Ok(()); }
2808        let mut manifest = self.ensure_field_meta(field)?;
2809        let new_seg = manifest.next_seg;
2810        manifest.next_seg = manifest.next_seg.checked_add(1).filter(|n| *n < keys::TEXT_TERM_STATS_SEG)
2811            .ok_or(Error::TooLarge)?;
2812        // Reserve the namespace before any candidate row is written.  A crash
2813        // before publication may leak that inactive range, but can never make
2814        // a later merge reuse it and mix stale rows into a replacement.
2815        self.put_field_meta(field, &manifest)?;
2816        self.commit()?;
2817        let built = self.build_text_segment(field, &[0], new_seg, 0)?;
2818        before_publish(self, new_seg)?;
2819
2820        // PUBLISH: old head remains on disk but is ignored while redirect is
2821        // present.  A crash on either side of this checkpoint therefore sees
2822        // exactly one complete copy.
2823        manifest.active.push((new_seg, 0));
2824        manifest.head_redirect = Some(new_seg);
2825        manifest.generation = manifest.generation.checked_add(1).ok_or(Error::TooLarge)?;
2826        self.put_field_meta(field, &manifest)?;
2827        self.commit()?; self.checkpoint()?;
2828
2829        self.rewrite_norm_owners(field, &[0], new_seg)?;
2830        let head_prefix = keys::text_seg_key(field, 0, b"");
2831        self.store().delete_prefix(&head_prefix)?;
2832        self.store().delete(&keys::text_meta_key(field, 0))?;
2833        self.commit()?; self.checkpoint()?;
2834        manifest.head_redirect = None;
2835        self.put_field_meta(field, &manifest)?;
2836        self.commit()?; self.checkpoint()?;
2837        debug_assert_eq!(built.doc_count, self.segment_meta(field, new_seg)?.doc_count);
2838
2839        loop {
2840            let mut choice = None;
2841            let max_level = manifest.active.iter().map(|&(_, level)| level).max().unwrap_or(0);
2842            for level in 0..=max_level {
2843                let same: Vec<u32> = manifest.active.iter().filter_map(|&(seg, l)| (l == level).then_some(seg)).collect();
2844                if same.len() >= MERGE_FANOUT { choice = Some((level, same[..MERGE_FANOUT].to_vec())); break; }
2845            }
2846            let Some((level, sources)) = choice else { break; };
2847            manifest = self.merge_text_segments(field, manifest, &sources, level + 1)?;
2848        }
2849        Ok(())
2850    }
2851
2852    fn build_text_segment(&mut self, field: u64, sources: &[u32], new_seg: u32, level: u32) -> Result<SegMeta> {
2853        use std::sync::atomic::{AtomicU64, Ordering};
2854        static SEQ: AtomicU64 = AtomicU64::new(0);
2855        let scratch = std::env::temp_dir().join(format!("text-merge-{}-{}", std::process::id(), SEQ.fetch_add(1, Ordering::Relaxed)));
2856        let mut sort = crate::bulk::ExternalSort::new(&scratch, 8 << 20)?;
2857        let mut expected = SegMeta { doc_count: 0, total_tokens: 0, dead_tokens: 0, dead: Vec::new(),
2858            level, term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0 };
2859
2860        for &source in sources {
2861            let source_meta = self.segment_meta(field, source)?;
2862            let dead: HashSet<u64> = source_meta.dead.into_iter().collect();
2863            let push_doc = |sort: &mut crate::bulk::ExternalSort, doc: u64, dl: u64| -> Result<()> {
2864                let mut k = vec![0]; k.extend_from_slice(&doc.to_be_bytes());
2865                let mut v = Vec::new(); write_varint(&mut v, dl); sort.push(k, v)
2866            };
2867            if source != 0 {
2868                let doc_prefix = keys::text_seg_doc_prefix(field, source);
2869                let mut doc_failure = None;
2870                self.store_ref().scan(&doc_prefix)?.for_each_ref(|key, val| {
2871                    if !key.starts_with(&doc_prefix) { return false; }
2872                    let result = if key.len() != doc_prefix.len() + 8 {
2873                        Err(corrupt("segment document key has an invalid length"))
2874                    } else {
2875                        let doc = u64::from_be_bytes(key[doc_prefix.len()..].try_into().unwrap());
2876                        let mut p = 0;
2877                        required_varint(val, &mut p, "segment document length is truncated").and_then(|dl|
2878                            if p == val.len() { push_doc(&mut sort, doc, dl) }
2879                            else { Err(corrupt("segment document length has trailing bytes")) })
2880                    };
2881                    if let Err(e) = result { doc_failure = Some(e); return false; }
2882                    true
2883                })?;
2884                if let Some(e) = doc_failure { return Err(e); }
2885            }
2886            let prefix = keys::text_seg_key(field, source, b"");
2887            let mut failure = None;
2888            self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
2889                if !key.starts_with(&prefix) { return false; }
2890                let body = &key[prefix.len()..];
2891                let push_post = |sort: &mut crate::bulk::ExternalSort, expected: &mut SegMeta,
2892                                 term: &[u8], doc: u64, tf: u64, dl: u64| -> Result<()> {
2893                    let mut k = Vec::with_capacity(term.len() + 10);
2894                    k.push(1); k.extend_from_slice(term); k.push(0); k.extend_from_slice(&doc.to_be_bytes());
2895                    let mut v = Vec::new(); write_varint(&mut v, tf); write_varint(&mut v, dl);
2896                    sort.push(k, v)?;
2897                    expected.posting_count = expected.posting_count.checked_add(1).ok_or(Error::TooLarge)?;
2898                    let h = posting_hash(term, doc, tf, dl);
2899                    expected.logical_xor ^= h; expected.logical_sum = expected.logical_sum.wrapping_add(h);
2900                    Ok(())
2901                };
2902                let result = if body.len() == 9 && body[0] == 0 {
2903                    let doc = u64::from_be_bytes(body[1..].try_into().unwrap());
2904                    if dead.contains(&doc) { Ok(()) } else {
2905                        let mut p = 0; required_varint(val, &mut p, "segment document length is truncated")
2906                            .and_then(|dl| if p == val.len() { push_doc(&mut sort, doc, dl) }
2907                                else { Err(corrupt("segment document length has trailing bytes")) })
2908                    }
2909                } else if source == 0 {
2910                    if body.len() < 10 || body[body.len() - 9] != 0 { Err(corrupt("head posting key is malformed")) } else {
2911                        let term = &body[..body.len() - 9];
2912                        let doc = u64::from_be_bytes(body[body.len() - 8..].try_into().unwrap());
2913                        if dead.contains(&doc) { Ok(()) } else {
2914                            let mut p = 0;
2915                            required_varint(val, &mut p, "head posting frequency is truncated").and_then(|tf|
2916                                required_varint(val, &mut p, "head posting length is truncated").and_then(|dl|
2917                                    if p == val.len() {
2918                                        push_doc(&mut sort, doc, dl)?;
2919                                        push_post(&mut sort, &mut expected, term, doc, tf, dl)
2920                                    }
2921                                    else { Err(corrupt("head posting has trailing bytes")) }))
2922                        }
2923                    }
2924                } else {
2925                    let (term, is_block) = if body.len() >= 10 && body[body.len() - 9] == 0 {
2926                        (&body[..body.len() - 9], true)
2927                    } else { (body, false) };
2928                    decode_postings(val).and_then(|posts| {
2929                        if is_block && posts.len() > POSTING_BLOCK { return Err(corrupt("text posting block exceeds its bound")); }
2930                        for (doc, tf, dl) in posts {
2931                            if !dead.contains(&doc) { push_post(&mut sort, &mut expected, term, doc, tf, dl)?; }
2932                        }
2933                        Ok(())
2934                    })
2935                };
2936                if let Err(e) = result { failure = Some(e); return false; }
2937                true
2938            })?;
2939            if let Some(e) = failure { return Err(e); }
2940        }
2941
2942        let mut runs = sort.finish()?;
2943        let mut iter = runs.iter()?;
2944        let mut last_doc = None;
2945        let mut block_term: Vec<u8> = Vec::new();
2946        let mut block: Vec<(u64, u64, u64)> = Vec::with_capacity(POSTING_BLOCK);
2947        let mut blocks_for_term = 0usize;
2948        while let Some(item) = iter.next() {
2949            let (key, val, _) = item?;
2950            if key.first() == Some(&0) {
2951                if key.len() != 9 { return Err(corrupt("merged document key is malformed")); }
2952                let doc = u64::from_be_bytes(key[1..].try_into().unwrap());
2953                let mut p = 0; let dl = required_varint(&val, &mut p, "merged document length is truncated")?;
2954                if p != val.len() { return Err(corrupt("merged document length has trailing bytes")); }
2955                if last_doc == Some(doc) { continue; }
2956                if last_doc.is_some_and(|previous| previous > doc) { return Err(corrupt("merged segment document order regressed")); }
2957                self.store().put(&keys::text_seg_doc_key(field, new_seg, doc), &val)?;
2958                expected.doc_count = expected.doc_count.checked_add(1).ok_or(Error::TooLarge)?;
2959                expected.total_tokens = expected.total_tokens.checked_add(dl).ok_or(Error::TooLarge)?;
2960                last_doc = Some(doc);
2961                continue;
2962            }
2963            if key.first() != Some(&1) || key.len() < 11 || key[key.len() - 9] != 0 {
2964                return Err(corrupt("merged posting key is malformed"));
2965            }
2966            let term = &key[1..key.len() - 9];
2967            let doc = u64::from_be_bytes(key[key.len() - 8..].try_into().unwrap());
2968            let mut p = 0; let tf = required_varint(&val, &mut p, "merged posting frequency is truncated")?;
2969            let dl = required_varint(&val, &mut p, "merged posting length is truncated")?;
2970            if p != val.len() { return Err(corrupt("merged posting has trailing bytes")); }
2971            if !block.is_empty() && (block_term.as_slice() != term || block.len() == POSTING_BLOCK) {
2972                let changed_term = block_term.as_slice() != term;
2973                expected.term_rows = expected.term_rows.checked_add(1).ok_or(Error::TooLarge)?;
2974                self.write_posting_block(field, new_seg, &block_term, &block, blocks_for_term == 0)?;
2975                blocks_for_term += 1;
2976                block.clear();
2977                if changed_term { blocks_for_term = 0; }
2978            }
2979            if block.is_empty() { block_term = term.to_vec(); }
2980            if block.last().is_some_and(|(previous, _, _)| *previous >= doc) {
2981                return Err(corrupt("merged term contains duplicate documents"));
2982            }
2983            block.push((doc, tf, dl));
2984        }
2985        if !block.is_empty() {
2986            expected.term_rows = expected.term_rows.checked_add(1).ok_or(Error::TooLarge)?;
2987            self.write_posting_block(field, new_seg, &block_term, &block, blocks_for_term == 0)?;
2988        }
2989        self.store().put(&keys::text_meta_key(field, new_seg), &expected.encode())?;
2990        self.commit()?; self.checkpoint()?;
2991
2992        let cfg = crate::store::Config { budget_bytes: 16 * crate::page::PAGE_SIZE,
2993            io: self.store_ref().io_mode(), sync: crate::store::SyncMode::Off };
2994        let snapshot = crate::store::Store::open_snapshot(self.store_ref().dir(), cfg)?;
2995        let verifier = Graph::new(snapshot)?;
2996        verifier.verify_text_segment(field, new_seg, &expected)?;
2997        Ok(expected)
2998    }
2999
3000    fn write_posting_block(&mut self, field: u64, seg: u32, term: &[u8],
3001                           block: &[(u64, u64, u64)], first: bool) -> Result<()> {
3002        if block.is_empty() || block.len() > POSTING_BLOCK { return Err(Error::TooLarge); }
3003        let mut value = Vec::new(); let mut last = 0u64;
3004        for &(doc, tf, dl) in block {
3005            write_varint(&mut value, doc.checked_sub(last).ok_or_else(|| corrupt("posting block order regressed"))?);
3006            write_varint(&mut value, tf); write_varint(&mut value, dl); last = doc;
3007        }
3008        let key = if first { keys::text_seg_key(field, seg, term) }
3009            else { keys::text_seg_block_key(field, seg, term, block[0].0) };
3010        self.store().put(&key, &value)
3011    }
3012
3013    fn verify_text_segment(&self, field: u64, seg: u32, expected: &SegMeta) -> Result<()> {
3014        let actual = self.segment_meta(field, seg)?;
3015        if &actual != expected { return Err(corrupt("reopened text segment metadata differs from its builder manifest")); }
3016        let mut got = SegMeta { doc_count: 0, total_tokens: 0, dead_tokens: 0, dead: Vec::new(),
3017            level: expected.level, term_rows: 0, posting_count: 0, logical_xor: 0, logical_sum: 0 };
3018        let mut failure = None;
3019        let doc_prefix = keys::text_seg_doc_prefix(field, seg);
3020        self.store_ref().scan(&doc_prefix)?.for_each_ref(|key, val| {
3021            if !key.starts_with(&doc_prefix) { return false; }
3022            let result = if key.len() != doc_prefix.len() + 8 {
3023                Err(corrupt("verified segment document key has an invalid length"))
3024            } else {
3025                let mut p = 0;
3026                required_varint(val, &mut p, "verified segment document length is truncated").and_then(|dl| {
3027                    if p != val.len() { return Err(corrupt("verified segment document length has trailing bytes")); }
3028                    got.doc_count = got.doc_count.checked_add(1).ok_or(Error::TooLarge)?;
3029                    got.total_tokens = got.total_tokens.checked_add(dl).ok_or(Error::TooLarge)?; Ok(())
3030                })
3031            };
3032            if let Err(e) = result { failure = Some(e); return false; }
3033            true
3034        })?;
3035        if let Some(e) = failure.take() { return Err(e); }
3036
3037        let prefix = keys::text_seg_key(field, seg, b"");
3038        self.store_ref().scan(&prefix)?.for_each_ref(|key, val| {
3039            if !key.starts_with(&prefix) { return false; }
3040            let body = &key[prefix.len()..];
3041            let result = if body.len() >= 10 && body[body.len() - 9] == 0 {
3042                let term = &body[..body.len() - 9];
3043                let first = u64::from_be_bytes(body[body.len() - 8..].try_into().unwrap());
3044                decode_postings(val).and_then(|posts| {
3045                    if posts.is_empty() || posts.len() > POSTING_BLOCK || posts[0].0 != first {
3046                        return Err(corrupt("verified text posting block manifest is invalid"));
3047                    }
3048                    got.term_rows = got.term_rows.checked_add(1).ok_or(Error::TooLarge)?;
3049                    for (doc, tf, dl) in posts {
3050                        got.posting_count = got.posting_count.checked_add(1).ok_or(Error::TooLarge)?;
3051                        let h = posting_hash(term, doc, tf, dl); got.logical_xor ^= h; got.logical_sum = got.logical_sum.wrapping_add(h);
3052                    }
3053                    Ok(())
3054                })
3055            } else if !body.is_empty() {
3056                let term = body;
3057                decode_postings(val).and_then(|posts| {
3058                    if posts.is_empty() || posts.len() > POSTING_BLOCK {
3059                        return Err(corrupt("verified exact text posting block exceeds its bound"));
3060                    }
3061                    got.term_rows = got.term_rows.checked_add(1).ok_or(Error::TooLarge)?;
3062                    for (doc, tf, dl) in posts {
3063                        got.posting_count = got.posting_count.checked_add(1).ok_or(Error::TooLarge)?;
3064                        let h = posting_hash(term, doc, tf, dl); got.logical_xor ^= h; got.logical_sum = got.logical_sum.wrapping_add(h);
3065                    }
3066                    Ok(())
3067                })
3068            } else { Err(corrupt("verified text segment contains an unknown row shape")) };
3069            if let Err(e) = result { failure = Some(e); return false; }
3070            true
3071        })?;
3072        if let Some(e) = failure { return Err(e); }
3073        if got != *expected { return Err(corrupt("reopened text segment counts or logical checksum differ")); }
3074        Ok(())
3075    }
3076
3077    fn rewrite_norm_owners(&mut self, field: u64, old: &[u32], new_seg: u32) -> Result<()> {
3078        let prefix = keys::text_seg_doc_prefix(field, new_seg);
3079        let mut cursor = prefix.clone();
3080        loop {
3081            let mut docs = Vec::with_capacity(8192);
3082            self.store_ref().scan(&cursor)?.for_each_ref(|key, _| {
3083                if !key.starts_with(&prefix) || key.len() != prefix.len() + 8 { return false; }
3084                docs.push(u64::from_be_bytes(key[key.len() - 8..].try_into().unwrap()));
3085                docs.len() < 8192
3086            })?;
3087            if docs.is_empty() { break; }
3088            for &doc in &docs {
3089                let nk = keys::text_norm_key(field, doc);
3090                if let Some(v) = self.store_ref().get(&nk)? {
3091                    let mut norm = Norm::decode(&v)?;
3092                    if old.contains(&norm.owner) { norm.owner = new_seg; self.store().put(&nk, &norm.encode())?; }
3093                }
3094            }
3095            let Some(next) = docs.last().copied().and_then(|doc| doc.checked_add(1)) else { break; };
3096            cursor = keys::text_seg_doc_key(field, new_seg, next);
3097            if docs.len() < 8192 { break; }
3098        }
3099        Ok(())
3100    }
3101
3102    fn merge_text_segments(&mut self, field: u64, mut manifest: FieldMeta,
3103                           sources: &[u32], level: u32) -> Result<FieldMeta> {
3104        let new_seg = manifest.next_seg;
3105        manifest.next_seg = manifest.next_seg.checked_add(1).filter(|n| *n < keys::TEXT_TERM_STATS_SEG)
3106            .ok_or(Error::TooLarge)?;
3107        self.put_field_meta(field, &manifest)?;
3108        self.commit()?;
3109        self.build_text_segment(field, sources, new_seg, level)?;
3110        manifest.active.retain(|(seg, _)| !sources.contains(seg));
3111        manifest.active.push((new_seg, level));
3112        manifest.redirects = sources.iter().map(|&old| (old, new_seg)).collect();
3113        manifest.generation = manifest.generation.checked_add(1).ok_or(Error::TooLarge)?;
3114        self.put_field_meta(field, &manifest)?;
3115        self.commit()?; self.checkpoint()?; // PUBLISH before retirement
3116        self.rewrite_norm_owners(field, sources, new_seg)?;
3117        for &source in sources {
3118            self.store().delete_prefix(&keys::text_seg_key(field, source, b""))?;
3119            self.store().delete_prefix(&keys::text_seg_doc_prefix(field, source))?;
3120            self.store().delete(&keys::text_meta_key(field, source))?;
3121        }
3122        self.commit()?; self.checkpoint()?;
3123        manifest.redirects.clear();
3124        self.put_field_meta(field, &manifest)?;
3125        self.commit()?; self.checkpoint()?;
3126        Ok(manifest)
3127    }
3128
3129    /// Enumerate every distinct term of `field` starting at `from` (bytes),
3130    /// across the head (deduped) and every folded segment, in NO global
3131    /// order (per-source order only) -- callers union/dedupe. The callback
3132    /// returns false to stop that source's walk.
3133    fn for_each_term(&self, field: u64, from: &[u8],
3134                     mut f: impl FnMut(&[u8]) -> bool) -> Result<()> {
3135        if self.field_meta(field)?.is_some_and(|meta| meta.term_stats_ready) {
3136            let prefix = keys::text_term_lex_prefix(field);
3137            let mut start = prefix.clone();
3138            start.extend_from_slice(from);
3139            self.store_ref().scan(&start)?.for_each_ref(|key, _| {
3140                if !key.starts_with(&prefix) { return false; }
3141                f(&key[prefix.len()..])
3142            })?;
3143            return Ok(());
3144        }
3145        let segs = self.text_segments(field)?;
3146        for &seg in &segs {
3147            if seg == 0 {
3148                let mut prefix = keys::text_head_key(field, b"", 0);
3149                prefix.truncate(prefix.len() - 9);
3150                let mut start = prefix.clone();
3151                start.extend_from_slice(from);
3152                let it = self.store_ref().scan(&start)?;
3153                let mut last: Vec<u8> = Vec::new();
3154                let mut go = true;
3155                it.for_each_ref(|key, _| {
3156                    if !key.starts_with(&prefix) { return false; }
3157                    let body = &key[prefix.len()..];
3158                    if body.len() < 10 || body[body.len() - 9] != 0x00 { return true; }
3159                    let term = &body[..body.len() - 9];
3160                    if term != last.as_slice() {
3161                        last = term.to_vec();
3162                        go = f(term);
3163                    }
3164                    go
3165                })?;
3166            } else {
3167                let prefix = keys::text_seg_key(field, seg, b"");
3168                let mut start = prefix.clone();
3169                start.extend_from_slice(from);
3170                let it = self.store_ref().scan(&start)?;
3171                let mut last = Vec::new();
3172                it.for_each_ref(|key, _| {
3173                    if !key.starts_with(&prefix) { return false; }
3174                    let body = &key[prefix.len()..];
3175                    if body.len() == 9 && body[0] == 0 { return true; }
3176                    let term = if body.len() >= 10 && body[body.len() - 9] == 0 {
3177                        &body[..body.len() - 9]
3178                    } else { body };
3179                    if term == last.as_slice() { return true; }
3180                    last = term.to_vec();
3181                    f(term)
3182                })?;
3183            }
3184        }
3185        Ok(())
3186    }
3187
3188    /// All terms of `field` starting with `prefix`, clipped to the
3189    /// `limit` most frequent (Manticore's expansion rule: the rare tail of
3190    /// an expansion is mostly misspellings; keep the popular head).
3191    pub fn text_prefix_terms(&self, field: u64, prefix: &str, limit: usize)
3192        -> Result<Vec<String>>
3193    {
3194        let p = prefix.as_bytes();
3195        let mut terms: Vec<String> = Vec::new();
3196        self.for_each_term(field, p, |t| {
3197            if !t.starts_with(p) { return false; }
3198            if let Ok(s) = std::str::from_utf8(t) { terms.push(s.to_string()); }
3199            true
3200        })?;
3201        terms.sort_unstable(); terms.dedup();
3202        if terms.len() > limit {
3203            let mut by_df: Vec<(usize, String)> = terms.into_iter()
3204                .map(|t| {
3205                    let df = self.text_term_doc_freq(field, &t).ok().flatten()
3206                        .map(|n| n as usize)
3207                        .unwrap_or_else(|| self.text_postings(field, &t).map(|p| p.len()).unwrap_or(0));
3208                    (df, t)
3209                })
3210                .collect();
3211            by_df.sort_by(|a, b| b.0.cmp(&a.0));
3212            by_df.truncate(limit);
3213            terms = by_df.into_iter().map(|(_, t)| t).collect();
3214        }
3215        Ok(terms)
3216    }
3217
3218    /// Terms of `field` within `max_edits` (Levenshtein) of `word` -- the
3219    /// typo walk. Terms are visited in sorted order per source; a full DP
3220    /// row per term with shared-prefix reuse is O(|term| * |word|) worst
3221    /// case but the row's min bound prunes: when every cell of the row
3222    /// exceeds max_edits the whole SUBTREE of terms sharing that prefix is
3223    /// dead, and the walk seeks past it (restart scan at prefix-successor)
3224    /// instead of visiting each term (typesense's trie-DP on our sorted
3225    /// keys; caps: typos by length 0/<5, 1/<9, 2 else).
3226    pub fn text_fuzzy_terms(&self, field: u64, word: &str, max_edits: u32, limit: usize)
3227        -> Result<Vec<(String, u32)>>
3228    {
3229        let w: Vec<char> = word.chars().collect();
3230        let n = w.len();
3231        let mut found: Vec<(String, u32)> = Vec::new();
3232        // DP over rows: row[j] = edits between term-prefix and word[..j]
3233        let dp_next = |row: &Vec<u32>, ch: char| -> Vec<u32> {
3234            let mut nr = vec![row[0] + 1];
3235            for j in 1..=n {
3236                let cost = if w[j - 1] == ch { 0 } else { 1 };
3237                nr.push((row[j] + 1).min(nr[j - 1] + 1).min(row[j - 1] + cost));
3238            }
3239            nr
3240        };
3241        let mut visit = |term: &[u8]| -> Vec<u8> /* seek-to, empty = continue */ {
3242            let Ok(t) = std::str::from_utf8(term) else { return Vec::new() };
3243            let mut row: Vec<u32> = (0..=n as u32).collect();
3244            let mut alive_prefix = 0usize; // chars consumed while any cell <= max
3245            let mut bytes_at_alive = 0usize;
3246            for (ci, ch) in t.chars().enumerate() {
3247                row = dp_next(&row, ch);
3248                if row.iter().min().copied().unwrap_or(u32::MAX) > max_edits {
3249                    // dead prefix: everything sharing t[..=ci] is dead; seek
3250                    // to the successor of that byte prefix.
3251                    let dead_bytes = t.char_indices().nth(ci + 1)
3252                        .map(|(i, _)| i).unwrap_or(t.len());
3253                    let mut succ = term[..dead_bytes].to_vec();
3254                    while let Some(last) = succ.pop() {
3255                        if last < 0xFE { succ.push(last + 1); break; }
3256                    }
3257                    return succ;
3258                }
3259                alive_prefix = ci + 1;
3260                bytes_at_alive = t.char_indices().nth(ci + 1).map(|(i, _)| i).unwrap_or(t.len());
3261            }
3262            let _ = (alive_prefix, bytes_at_alive);
3263            if row[n] <= max_edits {
3264                found.push((t.to_string(), row[n]));
3265            }
3266            Vec::new()
3267        };
3268        // walk each source with seek-restarts
3269        let mut cursor: Vec<u8> = Vec::new();
3270        loop {
3271            let mut seek: Option<Vec<u8>> = None;
3272            self.for_each_term(field, &cursor, |t| {
3273                if t.as_ref() < cursor.as_slice() { return true; } // other source lag
3274                let s = visit(t);
3275                if s.is_empty() { true } else { seek = Some(s); false }
3276            })?;
3277            match seek {
3278                Some(sk) if sk > cursor => cursor = sk,
3279                _ => break,
3280            }
3281        }
3282        found.sort();
3283        found.dedup();
3284        found.sort_by(|a, b| a.1.cmp(&b.1));
3285        found.truncate(limit);
3286        Ok(found)
3287    }
3288
3289    /// Instant search (2h, the as-you-type shape): every token exact-or-typo
3290    /// expanded, the LAST token also prefix-expanded, documents must match
3291    /// ALL tokens (AND), ranked by (fewest edits used, then BM25 over the
3292    /// matched expansions). Typo budget by token length: 0 under 5 chars,
3293    /// 1 under 9, else 2.
3294    pub fn text_search_instant(&self, field: u64, query: &str, k: usize)
3295        -> Result<Vec<(u64, f64)>>
3296    {
3297        self.text_search_instant_typo(field, query, k, None)
3298    }
3299
3300    /// The instant walk with the edit budget FORCED (SQL's `typo => n`): the
3301    /// length ladder is a default, not a floor -- a 4-char token gets 0 edits
3302    /// by default, so "warz" can only reach "wars" when the caller raises it.
3303    pub fn text_search_instant_typo(&self, field: u64, query: &str, k: usize,
3304                                    forced_edits: Option<u32>)
3305        -> Result<Vec<(u64, f64)>>
3306    {
3307        let tokens = tokenize(query);
3308        if tokens.is_empty() { return Ok(Vec::new()); }
3309        let budget = |t: &str| -> u32 {
3310            if let Some(f) = forced_edits { return f; }
3311            let l = t.chars().count();
3312            if l < 5 { 0 } else if l < 9 { 1 } else { 2 }
3313        };
3314        let mut per_token: Vec<Vec<(String, u32)>> = Vec::new();
3315        for (i, tok) in tokens.iter().enumerate() {
3316            let last = i + 1 == tokens.len();
3317            let mut cands: Vec<(String, u32)> = Vec::new();
3318            if last {
3319                cands.extend(self.text_prefix_terms(field, tok, 50)?
3320                    .into_iter().map(|t| (t, 0u32)));
3321            }
3322            if cands.is_empty() || !last {
3323                let b = budget(tok);
3324                cands.extend(self.text_fuzzy_terms(field, tok, b, 50)?);
3325            }
3326            if cands.is_empty() { return Ok(Vec::new()); } // AND semantics
3327            cands.sort(); cands.dedup();
3328            per_token.push(cands);
3329        }
3330        // gather per-token doc sets with best edit cost
3331        let (n_docs, _avg) = self.text_stats(field)?;
3332        let mut doc_sets: Vec<std::collections::HashMap<u64, (u32, f64)>> = Vec::new();
3333        for cands in &per_token {
3334            let mut m: std::collections::HashMap<u64, (u32, f64)> = std::collections::HashMap::new();
3335            for (term, edits) in cands {
3336                let posts = self.text_postings(field, term)?;
3337                if posts.is_empty() { continue; }
3338                let df = posts.len() as f64;
3339                let idf = ((n_docs - df + 0.5) / (df + 0.5) + 1.0).ln();
3340                for (docid, _tf, _dl) in posts {
3341                    let e = m.entry(docid).or_insert((*edits, idf));
3342                    if *edits < e.0 { *e = (*edits, idf); }
3343                }
3344            }
3345            doc_sets.push(m);
3346        }
3347        // AND-intersect, score = sum over tokens of (2 - edits)*1000 + idf
3348        let (first, rest) = doc_sets.split_first().unwrap();
3349        let mut out: Vec<(u64, f64)> = Vec::new();
3350        'doc: for (&docid, &(e0, idf0)) in first {
3351            let mut score = (2.0 - e0 as f64) * 1000.0 + idf0;
3352            for m in rest {
3353                let Some(&(e, idf)) = m.get(&docid) else { continue 'doc };
3354                score += (2.0 - e as f64) * 1000.0 + idf;
3355            }
3356            out.push((docid, score));
3357        }
3358        out.sort_by(|a, b| b.1.total_cmp(&a.1).then(a.0.cmp(&b.0)));
3359        out.truncate(k);
3360        Ok(out)
3361    }
3362
3363    /// Every segment id present for `field` (0 = head, if it has meta).
3364    pub fn text_segments(&self, field: u64) -> Result<Vec<u32>> {
3365        if let Some(meta) = self.field_meta(field)? {
3366            let mut segs = Vec::new();
3367            if meta.head_redirect.is_none()
3368                && self.store_ref().get(&keys::text_meta_key(field, 0))?.is_some()
3369            { segs.push(0); }
3370            segs.extend(meta.active.into_iter().map(|(seg, _)| seg));
3371            segs.sort_unstable();
3372            return Ok(segs);
3373        }
3374        let mut segs = Vec::new();
3375        let from = keys::text_meta_key(field, 0);
3376        let it = self.store_ref().scan(&from)?;
3377        it.for_each_ref(|key, _| {
3378            if key.first() != Some(&keys::TAG_TEXTMETA) || key.len() != 13 { return false; }
3379            let f = u64::from_be_bytes(key[1..9].try_into().unwrap());
3380            if f != field { return false; }
3381            let seg = u32::from_be_bytes(key[9..13].try_into().unwrap());
3382            if seg < keys::TEXT_TERM_STATS_SEG { segs.push(seg); }
3383            true
3384        })?;
3385        Ok(segs)
3386    }
3387}
3388
3389#[cfg(test)]
3390mod merge_interruption_tests {
3391    use super::*;
3392    use crate::io::IoMode;
3393    use crate::store::{Config, Store, SyncMode};
3394
3395    fn cfg() -> Config {
3396        Config { budget_bytes: 8 << 20, io: IoMode::Buffered, sync: SyncMode::Off }
3397    }
3398
3399    #[test]
3400    fn an_interrupted_k_way_merge_leaves_every_old_batch_findable() {
3401        let d = tempfile::TempDir::new().unwrap();
3402        let before;
3403        {
3404            let mut g = Graph::new(Store::create(d.path(), cfg()).unwrap()).unwrap();
3405            for round in 0..3u64 {
3406                for i in 0..20u64 {
3407                    g.index_text(1, round * 100 + i + 1,
3408                        &format!("common batch{round} token{i}")).unwrap();
3409                }
3410                g.fold_text(1).unwrap();
3411            }
3412            before = g.text_search(1, "common token7", 100).unwrap();
3413            let manifest = g.field_meta(1).unwrap().unwrap();
3414            let sources: Vec<u32> = manifest.active.iter().map(|&(seg, _)| seg).collect();
3415
3416            // This is the exact pre-publication boundary: the candidate has
3417            // been built, checkpointed, reopened and verified, but the active
3418            // manifest still names only the old batches. Simulate process loss
3419            // by dropping the writer here.
3420            g.build_text_segment(1, &sources, manifest.next_seg, 1).unwrap();
3421        }
3422        let g = Graph::new(Store::open(d.path(), cfg()).unwrap()).unwrap();
3423        assert_eq!(g.text_search(1, "common token7", 100).unwrap(), before);
3424        assert_eq!(g.text_segments(1).unwrap().len(), 3,
3425            "an unpublished candidate must not displace any old batch");
3426    }
3427
3428    #[test]
3429    fn an_interruption_after_publish_reads_the_verified_replacement() {
3430        let d = tempfile::TempDir::new().unwrap();
3431        let before;
3432        {
3433            let mut g = Graph::new(Store::create(d.path(), cfg()).unwrap()).unwrap();
3434            for round in 0..3u64 {
3435                for i in 0..20u64 {
3436                    g.index_text(1, round * 100 + i + 1,
3437                        &format!("common batch{round} token{i}")).unwrap();
3438                }
3439                g.fold_text(1).unwrap();
3440            }
3441            before = g.text_search(1, "common token7", 100).unwrap();
3442            let mut manifest = g.field_meta(1).unwrap().unwrap();
3443            let sources: Vec<u32> = manifest.active.iter().map(|&(seg, _)| seg).collect();
3444            let replacement = manifest.next_seg;
3445            manifest.next_seg += 1;
3446            g.put_field_meta(1, &manifest).unwrap();
3447            g.commit().unwrap();
3448            g.build_text_segment(1, &sources, replacement, 1).unwrap();
3449
3450            // Publish the independently verified replacement, then simulate a
3451            // crash before owner rewrites and old-batch retirement.  Redirects
3452            // keep updates valid; reads use the complete replacement.
3453            manifest.active.retain(|(seg, _)| !sources.contains(seg));
3454            manifest.active.push((replacement, 1));
3455            manifest.redirects = sources.iter().map(|&old| (old, replacement)).collect();
3456            g.put_field_meta(1, &manifest).unwrap();
3457            g.commit().unwrap(); g.checkpoint().unwrap();
3458        }
3459        let g = Graph::new(Store::open(d.path(), cfg()).unwrap()).unwrap();
3460        assert_eq!(g.text_search(1, "common token7", 100).unwrap(), before);
3461        assert_eq!(g.text_segments(1).unwrap(), vec![4],
3462            "a published verified replacement must be the sole visible batch");
3463    }
3464}
3465
3466#[cfg(test)]
3467mod build_accumulator_tests {
3468    use super::*;
3469
3470    fn build(open_budget: usize) -> Vec<(Vec<u8>, Vec<u8>, bool)> {
3471        let dir = tempfile::TempDir::new().unwrap();
3472        let mut accumulator = TextPostingAccumulator::new(
3473            &dir.path().join("blocks"), open_budget, 1024).unwrap();
3474        for doc in 1..=600u64 {
3475            accumulator.push(b"common".to_vec(), doc, 1, 4).unwrap();
3476            accumulator.push(format!("bucket{}", doc % 11).into_bytes(), doc, 2, 4).unwrap();
3477            accumulator.push(format!("unique{doc}").into_bytes(), doc, 1, 4).unwrap();
3478        }
3479        let (mut runs, postings, _, _, _) = accumulator.finish().unwrap();
3480        let rows: Vec<_> = TextBlockIter::new(runs.iter().unwrap(), 77)
3481            .collect::<Result<Vec<_>>>().unwrap();
3482        let decoded = rows.iter().map(|(_, value, _)| decode_postings(value).unwrap().len() as u64)
3483            .sum::<u64>();
3484        assert_eq!(decoded, postings);
3485        rows
3486    }
3487
3488    #[test]
3489    fn forced_spill_is_byte_identical_to_unspilled_posting_blocks() {
3490        let one = TextPostingAccumulator::entry_bytes(b"unique600");
3491        let spilled = build(one * 3);
3492        let unspilled = build(8 << 20);
3493        assert_eq!(spilled, unspilled);
3494    }
3495}