1use 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
27pub 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
44pub 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 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 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
166pub 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
276pub struct TextBuildCache {
280 slots: Vec<Option<(String, u64)>>,
281}
282
283pub 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#[derive(Clone, Debug, PartialEq, Eq)]
367pub struct SegMeta {
368 pub doc_count: u64,
369 pub total_tokens: u64,
370 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 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)>, 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 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
659struct 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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; 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 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 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 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 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 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 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 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 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 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 }; 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 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 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 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 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 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 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, } 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 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 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 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; } 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 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 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 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()?; 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 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 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 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 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> {
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; 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 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 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; } 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 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 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()); } cands.sort(); cands.dedup();
3328 per_token.push(cands);
3329 }
3330 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 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 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 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 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}