1use std::sync::atomic::{AtomicU64, Ordering};
10use std::time::{Duration, Instant};
11
12use rusqlite::{params, Connection, ErrorCode};
13use serde::Serialize;
14
15const DEFAULT_MERGE_INTERVAL: Duration = Duration::from_secs(300);
16const DEFAULT_MERGE_PAGES: u32 = 500;
17const MAX_MERGE_PAGES: u32 = i32::MAX as u32;
25const DEFAULT_MINIMUM_SEGMENTS: u64 = 2;
26
27#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
28pub struct FtsLevelStructure {
29 pub level: u64,
30 pub merge_input_segments: u64,
31 pub segment_count: u64,
32}
33
34#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
35pub struct FtsIndexStructure {
36 pub cookie: u32,
37 pub level_count: u64,
38 pub segment_count: u64,
39 pub level_zero_segments_written: u64,
40 pub levels: Vec<FtsLevelStructure>,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
44pub struct FtsSegmentDiagnostics {
45 pub entities: FtsIndexStructure,
46 pub notes: FtsIndexStructure,
47 pub total_segments: u64,
48}
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
51pub struct FtsMaintenanceCounters {
52 pub checks: u64,
53 pub attempts: u64,
54 pub work_steps: u64,
55 pub noops: u64,
56 pub below_threshold_skips: u64,
57 pub busy_skips: u64,
58 pub errors: u64,
59 pub requested_pages: u64,
60}
61
62static FTS_CHECKS: AtomicU64 = AtomicU64::new(0);
63static FTS_ATTEMPTS: AtomicU64 = AtomicU64::new(0);
64static FTS_WORK_STEPS: AtomicU64 = AtomicU64::new(0);
65static FTS_NOOPS: AtomicU64 = AtomicU64::new(0);
66static FTS_BELOW_THRESHOLD_SKIPS: AtomicU64 = AtomicU64::new(0);
67static FTS_BUSY_SKIPS: AtomicU64 = AtomicU64::new(0);
68static FTS_ERRORS: AtomicU64 = AtomicU64::new(0);
69static FTS_REQUESTED_PAGES: AtomicU64 = AtomicU64::new(0);
70
71pub fn fts_maintenance_counters() -> FtsMaintenanceCounters {
72 FtsMaintenanceCounters {
73 checks: FTS_CHECKS.load(Ordering::Relaxed),
74 attempts: FTS_ATTEMPTS.load(Ordering::Relaxed),
75 work_steps: FTS_WORK_STEPS.load(Ordering::Relaxed),
76 noops: FTS_NOOPS.load(Ordering::Relaxed),
77 below_threshold_skips: FTS_BELOW_THRESHOLD_SKIPS.load(Ordering::Relaxed),
78 busy_skips: FTS_BUSY_SKIPS.load(Ordering::Relaxed),
79 errors: FTS_ERRORS.load(Ordering::Relaxed),
80 requested_pages: FTS_REQUESTED_PAGES.load(Ordering::Relaxed),
81 }
82}
83
84#[derive(Debug, Clone)]
85pub(crate) struct FtsMaintenanceConfig {
86 pub(crate) enabled: bool,
87 pub(crate) interval: Duration,
88 pub(crate) merge_pages: u32,
89 pub(crate) minimum_segments: u64,
90}
91
92impl Default for FtsMaintenanceConfig {
93 fn default() -> Self {
94 Self {
95 enabled: true,
96 interval: DEFAULT_MERGE_INTERVAL,
97 merge_pages: DEFAULT_MERGE_PAGES,
98 minimum_segments: DEFAULT_MINIMUM_SEGMENTS,
99 }
100 }
101}
102
103impl FtsMaintenanceConfig {
104 pub(crate) fn from_env() -> Self {
105 let mut config = Self::default();
106 if let Ok(value) = std::env::var("KHIVE_FTS_MERGE_ENABLED") {
107 match value.trim().to_ascii_lowercase().as_str() {
108 "1" | "true" | "yes" | "on" => config.enabled = true,
109 "0" | "false" | "no" | "off" => config.enabled = false,
110 _ => tracing::warn!(
111 value,
112 fallback = config.enabled,
113 "invalid KHIVE_FTS_MERGE_ENABLED; using compiled default"
114 ),
115 }
116 }
117 if let Ok(value) = std::env::var("KHIVE_FTS_MERGE_INTERVAL_SECS") {
118 match value.parse::<u64>() {
119 Ok(seconds) if seconds > 0 => config.interval = Duration::from_secs(seconds),
120 _ => tracing::warn!(
121 value,
122 fallback_secs = config.interval.as_secs(),
123 "invalid KHIVE_FTS_MERGE_INTERVAL_SECS; using compiled default"
124 ),
125 }
126 }
127 if let Ok(value) = std::env::var("KHIVE_FTS_MERGE_PAGES") {
128 match parse_merge_pages(&value) {
129 Some(pages) => config.merge_pages = pages,
130 None => tracing::warn!(
131 value,
132 fallback_pages = config.merge_pages,
133 max_pages = MAX_MERGE_PAGES,
134 "invalid KHIVE_FTS_MERGE_PAGES (expected 1..=max_pages); using compiled \
135 default"
136 ),
137 }
138 }
139 if let Ok(value) = std::env::var("KHIVE_FTS_MERGE_MIN_SEGMENTS") {
140 match value.parse::<u64>() {
141 Ok(segments) if segments >= 2 => config.minimum_segments = segments,
142 _ => tracing::warn!(
143 value,
144 fallback_segments = config.minimum_segments,
145 "invalid KHIVE_FTS_MERGE_MIN_SEGMENTS; using compiled default"
146 ),
147 }
148 }
149 config
150 }
151}
152
153#[derive(Debug, Clone, Copy)]
154struct FtsTable {
155 name: &'static str,
156 structure_sql: &'static str,
157 merge_sql: &'static str,
158}
159
160const FTS_TABLES: [FtsTable; 2] = [
161 FtsTable {
162 name: "fts_entities",
163 structure_sql: "SELECT block FROM fts_entities_data WHERE id = 10",
164 merge_sql: "INSERT INTO fts_entities(fts_entities, rank) VALUES('merge', ?1)",
165 },
166 FtsTable {
167 name: "fts_notes",
168 structure_sql: "SELECT block FROM fts_notes_data WHERE id = 10",
169 merge_sql: "INSERT INTO fts_notes(fts_notes, rank) VALUES('merge', ?1)",
170 },
171];
172
173#[derive(Debug)]
174pub(crate) struct FtsMaintenanceState {
175 last_check: Instant,
176 next_table: usize,
177 in_progress: [bool; FTS_TABLES.len()],
178}
179
180impl FtsMaintenanceState {
181 pub(crate) fn new(now: Instant) -> Self {
182 Self {
183 last_check: now,
184 next_table: 0,
185 in_progress: [false; FTS_TABLES.len()],
186 }
187 }
188
189 pub(crate) fn is_due(&self, config: &FtsMaintenanceConfig, now: Instant) -> bool {
193 config.enabled && now.saturating_duration_since(self.last_check) >= config.interval
194 }
195
196 #[cfg(test)]
197 fn table_in_progress(&self, table: &str) -> bool {
198 FTS_TABLES
199 .iter()
200 .position(|candidate| candidate.name == table)
201 .is_some_and(|index| self.in_progress[index])
202 }
203}
204
205#[derive(Debug, Clone, Copy, PartialEq, Eq)]
206pub(crate) enum FtsMaintenanceOutcome {
207 Worked,
208 Noop,
209 BelowThreshold,
210 Busy,
211}
212
213#[derive(Debug, Clone, PartialEq, Eq)]
214pub(crate) struct FtsMaintenanceStep {
215 pub(crate) table: &'static str,
216 pub(crate) requested_pages: i64,
217 pub(crate) segments_before: u64,
218 pub(crate) segments_after: u64,
219 pub(crate) outcome: FtsMaintenanceOutcome,
220}
221
222fn read_varint(bytes: &[u8], cursor: &mut usize) -> Result<u64, String> {
223 let mut value = 0_u64;
224 for index in 0..9 {
225 let byte = *bytes
226 .get(*cursor)
227 .ok_or_else(|| "truncated FTS5 structure varint".to_string())?;
228 *cursor += 1;
229 if index == 8 {
230 return value
231 .checked_shl(8)
232 .map(|prefix| prefix | u64::from(byte))
233 .ok_or_else(|| "FTS5 structure varint overflow".to_string());
234 }
235 value = value
236 .checked_shl(7)
237 .map(|prefix| prefix | u64::from(byte & 0x7f))
238 .ok_or_else(|| "FTS5 structure varint overflow".to_string())?;
239 if byte & 0x80 == 0 {
240 return Ok(value);
241 }
242 }
243 unreachable!("nine-byte SQLite varint returns inside the loop")
244}
245
246const FTS5_STRUCTURE_V2_MARKER: [u8; 4] = [0xFF, 0x00, 0x00, 0x01];
259
260pub(crate) fn parse_structure_record(bytes: &[u8]) -> Result<FtsIndexStructure, String> {
261 let cookie_bytes: [u8; 4] = bytes
262 .get(..4)
263 .ok_or_else(|| "FTS5 structure record is shorter than its 4-byte cookie".to_string())?
264 .try_into()
265 .expect("slice length checked above");
266 let cookie = u32::from_be_bytes(cookie_bytes);
267 let mut cursor = 4;
268 let is_v2 = bytes.get(cursor..cursor + 4) == Some(&FTS5_STRUCTURE_V2_MARKER[..]);
269 if is_v2 {
270 cursor += 4;
271 }
272 let level_count = read_varint(bytes, &mut cursor)?;
273 let segment_count = read_varint(bytes, &mut cursor)?;
274 let level_zero_segments_written = read_varint(bytes, &mut cursor)?;
275 let level_capacity = usize::try_from(level_count)
276 .map_err(|_| "FTS5 level count does not fit this platform".to_string())?;
277 if level_capacity > bytes.len().saturating_sub(cursor) / 2 {
278 return Err(format!(
279 "FTS5 structure declares {level_count} levels but only {} bytes remain",
280 bytes.len().saturating_sub(cursor)
281 ));
282 }
283
284 let bytes_per_segment: u64 = if is_v2 { 8 } else { 3 };
287
288 let mut levels = Vec::with_capacity(level_capacity);
289 let mut parsed_segments = 0_u64;
290 for level in 0..level_count {
291 let merge_input_segments = read_varint(bytes, &mut cursor)?;
292 let level_segments = read_varint(bytes, &mut cursor)?;
293 if merge_input_segments > level_segments {
298 return Err(format!(
299 "FTS5 level {level} declares {merge_input_segments} merge input segments but only \
300 {level_segments} segments"
301 ));
302 }
303 let remaining = bytes.len().saturating_sub(cursor);
304 if level_segments > remaining as u64 / bytes_per_segment {
305 return Err(format!(
306 "FTS5 level {level} declares {level_segments} segments but its record is truncated"
307 ));
308 }
309 for _ in 0..level_segments {
310 let _segment_id = read_varint(bytes, &mut cursor)?;
311 let first_leaf = read_varint(bytes, &mut cursor)?;
312 let last_leaf = read_varint(bytes, &mut cursor)?;
313 if last_leaf < first_leaf {
323 return Err(format!(
324 "FTS5 level {level} has invalid leaf range {first_leaf}..={last_leaf}"
325 ));
326 }
327 if is_v2 {
328 for _ in 0..5 {
333 read_varint(bytes, &mut cursor)?;
334 }
335 }
336 }
337 parsed_segments = parsed_segments
338 .checked_add(level_segments)
339 .ok_or_else(|| "FTS5 per-level segment count overflow".to_string())?;
340 levels.push(FtsLevelStructure {
341 level,
342 merge_input_segments,
343 segment_count: level_segments,
344 });
345 }
346 if parsed_segments != segment_count {
347 return Err(format!(
348 "FTS5 structure segment count mismatch: header declares {segment_count}, levels contain {parsed_segments}"
349 ));
350 }
351
352 Ok(FtsIndexStructure {
353 cookie,
354 level_count,
355 segment_count,
356 level_zero_segments_written,
357 levels,
358 })
359}
360
361fn inspect_table(conn: &Connection, table: FtsTable) -> Result<FtsIndexStructure, String> {
362 let structure = conn
363 .query_row(table.structure_sql, [], |row| row.get::<_, Vec<u8>>(0))
364 .map_err(|error| format!("{} structure read failed: {error}", table.name))?;
365 parse_structure_record(&structure)
366 .map_err(|error| format!("{} structure record is invalid: {error}", table.name))
367}
368
369fn has_active_merge(structure: &FtsIndexStructure) -> bool {
370 structure
371 .levels
372 .iter()
373 .any(|level| level.merge_input_segments > 0)
374}
375
376fn fts_table(name: &str) -> FtsTable {
382 FTS_TABLES
383 .iter()
384 .find(|table| table.name == name)
385 .copied()
386 .unwrap_or_else(|| unreachable!("FTS_TABLES must define a {name} entry"))
387}
388
389pub(crate) fn inspect_fts_segments(conn: &Connection) -> Result<FtsSegmentDiagnostics, String> {
390 let entities = inspect_table(conn, fts_table("fts_entities"))?;
391 let notes = inspect_table(conn, fts_table("fts_notes"))?;
392 let total_segments = entities.segment_count.saturating_add(notes.segment_count);
393 Ok(FtsSegmentDiagnostics {
394 entities,
395 notes,
396 total_segments,
397 })
398}
399
400fn is_busy(error: &rusqlite::Error) -> bool {
401 matches!(
402 error.sqlite_error_code(),
403 Some(ErrorCode::DatabaseBusy | ErrorCode::DatabaseLocked)
404 )
405}
406
407fn parse_merge_pages(value: &str) -> Option<u32> {
411 value
412 .trim()
413 .parse::<u32>()
414 .ok()
415 .filter(|pages| (1..=MAX_MERGE_PAGES).contains(pages))
416}
417
418fn merge_page_budget(config: &FtsMaintenanceConfig, merge_in_progress: bool) -> i64 {
423 let page_budget = i64::from(config.merge_pages.clamp(1, MAX_MERGE_PAGES));
424 if merge_in_progress {
425 page_budget
426 } else {
427 -page_budget
428 }
429}
430
431fn run_merge_statement(
432 conn: &Connection,
433 table: FtsTable,
434 requested_pages: i64,
435) -> Result<FtsMaintenanceOutcome, String> {
436 let prior_busy_ms = conn
437 .query_row("PRAGMA busy_timeout", [], |row| row.get::<_, i64>(0))
438 .map_err(|error| format!("read busy_timeout before FTS maintenance: {error}"))?;
439 conn.busy_timeout(Duration::ZERO)
440 .map_err(|error| format!("disable busy wait for FTS maintenance: {error}"))?;
441 FTS_ATTEMPTS.fetch_add(1, Ordering::Relaxed);
442 FTS_REQUESTED_PAGES.fetch_add(requested_pages.unsigned_abs(), Ordering::Relaxed);
443 let changes_before = conn.total_changes();
444 let execution = conn.execute(table.merge_sql, params![requested_pages]);
445 let restore = conn.busy_timeout(Duration::from_millis(prior_busy_ms.max(0) as u64));
446 if let Err(error) = restore {
447 FTS_ERRORS.fetch_add(1, Ordering::Relaxed);
448 return Err(format!(
449 "restore busy_timeout after {} maintenance: {error}",
450 table.name
451 ));
452 }
453 match execution {
454 Ok(_) => {
455 let changes = conn.total_changes().saturating_sub(changes_before);
456 if changes >= 2 {
457 FTS_WORK_STEPS.fetch_add(1, Ordering::Relaxed);
458 Ok(FtsMaintenanceOutcome::Worked)
459 } else {
460 FTS_NOOPS.fetch_add(1, Ordering::Relaxed);
461 Ok(FtsMaintenanceOutcome::Noop)
462 }
463 }
464 Err(error) if is_busy(&error) => {
465 FTS_BUSY_SKIPS.fetch_add(1, Ordering::Relaxed);
466 Ok(FtsMaintenanceOutcome::Busy)
467 }
468 Err(error) => {
469 FTS_ERRORS.fetch_add(1, Ordering::Relaxed);
470 Err(format!("{} bounded merge failed: {error}", table.name))
471 }
472 }
473}
474
475pub(crate) fn run_if_due(
476 conn: &Connection,
477 config: &FtsMaintenanceConfig,
478 state: &mut FtsMaintenanceState,
479 now: Instant,
480) -> Result<Option<FtsMaintenanceStep>, String> {
481 if !state.is_due(config, now) {
482 return Ok(None);
483 }
484 state.last_check = now;
485 let table_index = state.next_table;
486 state.next_table = (state.next_table + 1) % FTS_TABLES.len();
487 let table = FTS_TABLES[table_index];
488 FTS_CHECKS.fetch_add(1, Ordering::Relaxed);
489
490 let before = match inspect_table(conn, table) {
491 Ok(structure) => structure,
492 Err(error) => {
493 FTS_ERRORS.fetch_add(1, Ordering::Relaxed);
494 return Err(error);
495 }
496 };
497 let merge_in_progress = state.in_progress[table_index] || has_active_merge(&before);
501 if !merge_in_progress && before.segment_count < config.minimum_segments {
502 FTS_BELOW_THRESHOLD_SKIPS.fetch_add(1, Ordering::Relaxed);
503 return Ok(Some(FtsMaintenanceStep {
504 table: table.name,
505 requested_pages: 0,
506 segments_before: before.segment_count,
507 segments_after: before.segment_count,
508 outcome: FtsMaintenanceOutcome::BelowThreshold,
509 }));
510 }
511
512 let requested_pages = merge_page_budget(config, merge_in_progress);
513 let outcome = run_merge_statement(conn, table, requested_pages)?;
514 if outcome == FtsMaintenanceOutcome::Busy {
515 return Ok(Some(FtsMaintenanceStep {
516 table: table.name,
517 requested_pages,
518 segments_before: before.segment_count,
519 segments_after: before.segment_count,
520 outcome,
521 }));
522 }
523
524 let after = match inspect_table(conn, table) {
525 Ok(structure) => structure,
526 Err(error) => {
527 state.in_progress[table_index] = outcome == FtsMaintenanceOutcome::Worked;
528 FTS_ERRORS.fetch_add(1, Ordering::Relaxed);
529 return Err(error);
530 }
531 };
532 state.in_progress[table_index] = outcome == FtsMaintenanceOutcome::Worked
533 && (after.segment_count > 1 || has_active_merge(&after));
534 Ok(Some(FtsMaintenanceStep {
535 table: table.name,
536 requested_pages,
537 segments_before: before.segment_count,
538 segments_after: after.segment_count,
539 outcome,
540 }))
541}
542
543#[cfg(test)]
544mod tests {
545 use std::time::{Duration, Instant};
546
547 use rusqlite::{params, Connection};
548 use tempfile::TempDir;
549
550 use super::*;
551
552 fn create_fts_tables(conn: &Connection) {
553 conn.execute_batch(
554 "CREATE VIRTUAL TABLE fts_entities USING fts5(namespace UNINDEXED, subject_id UNINDEXED, title, body, tokenize='trigram');
555 CREATE VIRTUAL TABLE fts_notes USING fts5(namespace UNINDEXED, subject_id UNINDEXED, title, body, tokenize='trigram');
556 INSERT INTO fts_entities(fts_entities, rank) VALUES('automerge', 0);
557 INSERT INTO fts_notes(fts_notes, rank) VALUES('automerge', 0);",
558 )
559 .expect("create FTS fixtures");
560 }
561
562 fn seed_segments(conn: &Connection, table: &str, count: usize) {
563 let sql = format!(
564 "INSERT INTO {table}(namespace, subject_id, title, body) VALUES(?1, ?2, ?3, ?4)"
565 );
566 for index in 0..count {
567 let body = format!(
568 "segment fixture {index} keeps enough repeated production recall text to span pages {}",
569 "memory query corpus ".repeat(40)
570 );
571 conn.execute(
572 &sql,
573 params![
574 "local",
575 format!("id-{index}"),
576 format!("title {index}"),
577 body
578 ],
579 )
580 .expect("one autocommit FTS write");
581 }
582 }
583
584 fn file_fixture() -> (TempDir, Connection, Connection) {
585 let dir = tempfile::tempdir().expect("tempdir");
586 let path = dir.path().join("fts-maintenance.db");
587 let maintenance = Connection::open(&path).expect("maintenance connection");
588 maintenance
589 .pragma_update(None, "journal_mode", "WAL")
590 .expect("enable WAL");
591 create_fts_tables(&maintenance);
592 let writer = Connection::open(&path).expect("writer connection");
593 (dir, maintenance, writer)
594 }
595
596 fn push_varint(mut value: u64, out: &mut Vec<u8>) {
597 let mut bytes = [0_u8; 10];
598 let mut cursor = bytes.len();
599 cursor -= 1;
600 bytes[cursor] = (value & 0x7f) as u8;
601 value >>= 7;
602 while value != 0 {
603 cursor -= 1;
604 bytes[cursor] = ((value & 0x7f) as u8) | 0x80;
605 value >>= 7;
606 }
607 out.extend_from_slice(&bytes[cursor..]);
608 }
609
610 #[test]
611 fn parses_documented_structure_record_with_multibyte_segment_count() {
612 let mut bytes = vec![0, 0, 0, 7];
613 push_varint(2, &mut bytes);
614 push_varint(130, &mut bytes);
615 push_varint(321, &mut bytes);
616 push_varint(3, &mut bytes);
617 push_varint(70, &mut bytes);
618 for segment in 0..70 {
619 push_varint(segment + 1, &mut bytes);
620 push_varint(1, &mut bytes);
621 push_varint(2, &mut bytes);
622 }
623 push_varint(0, &mut bytes);
624 push_varint(60, &mut bytes);
625 for segment in 0..60 {
626 push_varint(segment + 71, &mut bytes);
627 push_varint(1, &mut bytes);
628 push_varint(1, &mut bytes);
629 }
630
631 let parsed = parse_structure_record(&bytes).expect("valid documented structure record");
632
633 assert_eq!(parsed.cookie, 7);
634 assert_eq!(parsed.level_count, 2);
635 assert_eq!(parsed.segment_count, 130);
636 assert_eq!(parsed.level_zero_segments_written, 321);
637 assert_eq!(parsed.levels[0].merge_input_segments, 3);
638 assert_eq!(parsed.levels[0].segment_count, 70);
639 assert_eq!(parsed.levels[1].segment_count, 60);
640 }
641
642 #[test]
643 fn parses_v2_contentless_delete_structure_record() {
644 let bytes = [
664 0x00, 0x00, 0x00, 0x00, 0xff, 0x00, 0x00, 0x01, 0x01, 0x02, 0x02, 0x00, 0x02, 0x01,
665 0x01, 0x01, 0x01, 0x01, 0x00, 0x00, 0x01, 0x02, 0x01, 0x01, 0x02, 0x02, 0x00, 0x00,
666 0x01,
667 ];
668
669 let parsed = parse_structure_record(&bytes).expect("valid V2 structure record");
670
671 assert_eq!(parsed.cookie, 0);
672 assert_eq!(parsed.level_count, 1);
673 assert_eq!(parsed.segment_count, 2);
674 assert_eq!(parsed.level_zero_segments_written, 2);
675 assert_eq!(parsed.levels.len(), 1);
676 assert_eq!(parsed.levels[0].merge_input_segments, 0);
677 assert_eq!(parsed.levels[0].segment_count, 2);
678 }
679
680 #[test]
681 fn v1_structure_record_is_unaffected_by_v2_marker_detection() {
682 let mut bytes = vec![0, 0, 0, 9];
687 push_varint(1, &mut bytes); push_varint(1, &mut bytes); push_varint(5, &mut bytes); push_varint(0, &mut bytes); push_varint(1, &mut bytes); push_varint(1, &mut bytes); push_varint(1, &mut bytes); push_varint(1, &mut bytes); let parsed = parse_structure_record(&bytes).expect("valid V1 structure record");
697 assert_eq!(parsed.cookie, 9);
698 assert_eq!(parsed.segment_count, 1);
699 }
700
701 #[test]
702 fn rejects_truncated_or_self_inconsistent_structure_records() {
703 assert!(parse_structure_record(&[0, 0, 0]).is_err());
704
705 let inconsistent = [0, 0, 0, 0, 1, 2, 0, 0, 1, 1, 1, 1];
706 let error = parse_structure_record(&inconsistent)
707 .expect_err("declared total must equal per-level total");
708 assert!(error.contains("segment count"), "unexpected error: {error}");
709 }
710
711 #[test]
712 fn rejects_merge_input_exceeding_the_level_segment_count() {
713 let mut bytes = vec![0, 0, 0, 0];
717 push_varint(1, &mut bytes); push_varint(2, &mut bytes); push_varint(2, &mut bytes); push_varint(3, &mut bytes); push_varint(2, &mut bytes); for segment in 1..=2 {
723 push_varint(segment, &mut bytes);
724 push_varint(1, &mut bytes);
725 push_varint(1, &mut bytes);
726 }
727
728 let error = parse_structure_record(&bytes)
729 .expect_err("merge input segments must not exceed the level's segment count");
730 assert!(error.contains("merge input"), "unexpected error: {error}");
731
732 bytes[7] = 2;
735 let parsed = parse_structure_record(&bytes).expect("consistent record parses");
736 assert_eq!(parsed.levels[0].merge_input_segments, 2);
737 assert!(has_active_merge(&parsed));
738 }
739
740 #[test]
741 fn accepts_a_trimmed_merge_input_segment() {
742 let mut bytes = vec![0, 0, 0, 0];
751 push_varint(1, &mut bytes); push_varint(2, &mut bytes); push_varint(5, &mut bytes); push_varint(2, &mut bytes); push_varint(2, &mut bytes); push_varint(1, &mut bytes); push_varint(0, &mut bytes); push_varint(0, &mut bytes); push_varint(2, &mut bytes); push_varint(5, &mut bytes); push_varint(9, &mut bytes); let parsed =
764 parse_structure_record(&bytes).expect("a trimmed input segment must not be corrupt");
765
766 assert_eq!(parsed.levels[0].segment_count, 2);
767 assert_eq!(parsed.levels[0].merge_input_segments, 2);
768 assert!(has_active_merge(&parsed));
769 }
770
771 #[test]
772 fn rejects_last_leaf_before_first_leaf() {
773 let mut bytes = vec![0, 0, 0, 0];
777 push_varint(1, &mut bytes); push_varint(1, &mut bytes); push_varint(2, &mut bytes); push_varint(0, &mut bytes); push_varint(1, &mut bytes); push_varint(1, &mut bytes); push_varint(5, &mut bytes); push_varint(3, &mut bytes); let error = parse_structure_record(&bytes)
787 .expect_err("last_leaf below first_leaf must still be rejected");
788 assert!(
789 error.contains("invalid leaf range"),
790 "unexpected error: {error}"
791 );
792 }
793
794 #[test]
795 fn segment_diagnostics_use_one_structure_row_and_never_scan_the_idx_table() {
796 let conn = Connection::open_in_memory().expect("in-memory sqlite");
797 create_fts_tables(&conn);
798 seed_segments(&conn, "fts_notes", 8);
799
800 let mut statement = conn
801 .prepare(&format!(
802 "EXPLAIN QUERY PLAN {}",
803 FTS_TABLES[1].structure_sql
804 ))
805 .expect("explain structure lookup");
806 let details: Vec<String> = statement
807 .query_map([], |row| row.get(3))
808 .expect("query plan rows")
809 .collect::<Result<_, _>>()
810 .expect("query plan details");
811
812 assert!(
813 details
814 .iter()
815 .all(|detail| !detail.contains("fts_notes_idx")),
816 "segment diagnostics must not scan the large idx shadow table: {details:?}"
817 );
818 assert!(
819 !FTS_TABLES[1]
820 .structure_sql
821 .to_ascii_uppercase()
822 .contains("COUNT"),
823 "the segment count must come from the structure header, not COUNT(*)"
824 );
825 assert_eq!(
826 inspect_fts_segments(&conn)
827 .expect("decode structure record")
828 .notes
829 .segment_count,
830 8
831 );
832 }
833
834 #[test]
835 fn bounded_steps_converge_to_one_segment_without_changing_match_results() {
836 let conn = Connection::open_in_memory().expect("in-memory sqlite");
837 create_fts_tables(&conn);
838 seed_segments(&conn, "fts_notes", 80);
839 let before = inspect_fts_segments(&conn).expect("inspect fragmented fixture");
840 assert!(before.notes.segment_count > 1, "fixture must be fragmented");
841 let expected_rows: Vec<String> = conn
842 .prepare("SELECT subject_id FROM fts_notes WHERE fts_notes MATCH 'production recall' ORDER BY subject_id")
843 .unwrap()
844 .query_map([], |row| row.get(0))
845 .unwrap()
846 .collect::<Result<_, _>>()
847 .unwrap();
848
849 let config = FtsMaintenanceConfig {
850 enabled: true,
851 interval: Duration::ZERO,
852 merge_pages: 8,
853 minimum_segments: 2,
854 };
855 let mut state = FtsMaintenanceState::new(Instant::now());
856 for _ in 0..512 {
857 let now = Instant::now();
858 run_if_due(&conn, &config, &mut state, now).expect("bounded maintenance step");
859 if inspect_fts_segments(&conn).unwrap().notes.segment_count == 1 {
860 break;
861 }
862 }
863
864 let after = inspect_fts_segments(&conn).expect("inspect optimized fixture");
865 assert_eq!(after.notes.segment_count, 1);
866 assert!(
867 !has_active_merge(&after.notes),
868 "convergence must leave no active merge behind"
869 );
870 assert!(
871 !state.table_in_progress("fts_notes"),
872 "convergence must clear the in-progress bookkeeping"
873 );
874 let actual_rows: Vec<String> = conn
875 .prepare("SELECT subject_id FROM fts_notes WHERE fts_notes MATCH 'production recall' ORDER BY subject_id")
876 .unwrap()
877 .query_map([], |row| row.get(0))
878 .unwrap()
879 .collect::<Result<_, _>>()
880 .unwrap();
881 assert_eq!(actual_rows, expected_rows);
882 assert_eq!(
883 conn.query_row("SELECT COUNT(*) FROM fts_notes", [], |row| row
884 .get::<_, i64>(0))
885 .unwrap(),
886 80
887 );
888 }
889
890 #[test]
891 fn maintenance_round_robins_and_continues_with_positive_page_budget() {
892 let conn = Connection::open_in_memory().expect("in-memory sqlite");
893 create_fts_tables(&conn);
894 seed_segments(&conn, "fts_entities", 80);
895 seed_segments(&conn, "fts_notes", 80);
896 let config = FtsMaintenanceConfig {
897 enabled: true,
898 interval: Duration::ZERO,
899 merge_pages: 1,
900 minimum_segments: 2,
901 };
902 let mut state = FtsMaintenanceState::new(Instant::now());
903
904 let first = run_if_due(&conn, &config, &mut state, Instant::now())
905 .unwrap()
906 .expect("first due step");
907 let second = run_if_due(&conn, &config, &mut state, Instant::now())
908 .unwrap()
909 .expect("second due step");
910 let third = run_if_due(&conn, &config, &mut state, Instant::now())
911 .unwrap()
912 .expect("third due step");
913
914 assert_eq!(first.table, "fts_entities");
915 assert_eq!(first.requested_pages, -1);
916 assert_eq!(second.table, "fts_notes");
917 assert_eq!(second.requested_pages, -1);
918 assert_eq!(third.table, "fts_entities");
919 assert_eq!(third.requested_pages, 1);
920 }
921
922 #[test]
923 fn persisted_incremental_merge_resumes_after_maintenance_state_restart() {
924 let conn = Connection::open_in_memory().expect("in-memory sqlite");
925 create_fts_tables(&conn);
926 seed_segments(&conn, "fts_entities", 80);
927 let config = FtsMaintenanceConfig {
928 enabled: true,
929 interval: Duration::ZERO,
930 merge_pages: 1,
931 minimum_segments: 2,
932 };
933 let mut first_process = FtsMaintenanceState::new(Instant::now());
934
935 let starter = run_if_due(&conn, &config, &mut first_process, Instant::now())
936 .expect("starter succeeds")
937 .expect("starter is due");
938 assert_eq!(starter.requested_pages, -1);
939 assert!(
940 has_active_merge(&inspect_table(&conn, FTS_TABLES[0]).expect("structure")),
941 "one-page starter must persist an incremental merge in FTS5's structure record"
942 );
943
944 let mut restarted_process = FtsMaintenanceState::new(Instant::now());
945 let resumed = run_if_due(&conn, &config, &mut restarted_process, Instant::now())
946 .expect("resume succeeds")
947 .expect("resume is due");
948 assert_eq!(resumed.table, "fts_entities");
949 assert_eq!(
950 resumed.requested_pages, 1,
951 "persisted nMerge state must continue instead of restarting optimize"
952 );
953 }
954
955 #[test]
956 fn held_writer_is_counted_as_busy_without_losing_the_optimize_cycle() {
957 let (_dir, maintenance, writer) = file_fixture();
958 seed_segments(&maintenance, "fts_entities", 20);
959 writer
960 .execute_batch("BEGIN IMMEDIATE")
961 .expect("hold writer");
962 let counters_before = fts_maintenance_counters();
963 let config = FtsMaintenanceConfig {
964 enabled: true,
965 interval: Duration::ZERO,
966 merge_pages: 8,
967 minimum_segments: 2,
968 };
969 let mut state = FtsMaintenanceState::new(Instant::now());
970
971 let step = run_if_due(&maintenance, &config, &mut state, Instant::now())
972 .expect("busy is an outcome, not a task failure")
973 .expect("due step");
974
975 assert_eq!(step.outcome, FtsMaintenanceOutcome::Busy);
976 assert_eq!(
977 fts_maintenance_counters().busy_skips,
978 counters_before.busy_skips + 1
979 );
980 assert!(
981 !state.table_in_progress("fts_entities"),
982 "a busy initial step must retry the negative starter next cadence"
983 );
984 writer.execute_batch("ROLLBACK").expect("release writer");
985 }
986
987 #[test]
988 fn merge_pages_override_must_fit_the_fts5_int32_argument() {
989 assert_eq!(parse_merge_pages("500"), Some(500));
990 assert_eq!(parse_merge_pages(" 1 "), Some(1));
991 assert_eq!(parse_merge_pages("2147483647"), Some(i32::MAX as u32));
992 assert_eq!(parse_merge_pages("2147483648"), None);
995 assert_eq!(parse_merge_pages("4294967295"), None);
996 assert_eq!(parse_merge_pages("0"), None);
997 assert_eq!(parse_merge_pages("-5"), None);
998 assert_eq!(parse_merge_pages("many"), None);
999 }
1000
1001 #[test]
1002 fn merge_page_budget_fits_the_fts5_int32_argument_in_either_sign() {
1003 let oversized = FtsMaintenanceConfig {
1004 merge_pages: u32::MAX,
1005 ..FtsMaintenanceConfig::default()
1006 };
1007 let start = merge_page_budget(&oversized, false);
1008 let resume = merge_page_budget(&oversized, true);
1009 assert_eq!(start, -i64::from(i32::MAX));
1010 assert_eq!(resume, i64::from(i32::MAX));
1011 assert!(i32::try_from(start).is_ok() && i32::try_from(resume).is_ok());
1012
1013 let zero = FtsMaintenanceConfig {
1014 merge_pages: 0,
1015 ..FtsMaintenanceConfig::default()
1016 };
1017 assert_eq!(merge_page_budget(&zero, false), -1);
1018
1019 let ordinary = FtsMaintenanceConfig {
1020 merge_pages: 8,
1021 ..FtsMaintenanceConfig::default()
1022 };
1023 assert_eq!(merge_page_budget(&ordinary, false), -8);
1024 assert_eq!(merge_page_budget(&ordinary, true), 8);
1025 }
1026}