1use std::sync::Arc;
4
5use async_trait::async_trait;
6use chrono::{DateTime, TimeZone, Utc};
7use uuid::Uuid;
8
9use khive_score::DeterministicScore;
10use khive_storage::error::StorageError;
11use khive_storage::types::{
12 BatchWriteSummary, IndexRebuildScope, SqlStatement, SqlValue, TextDocument, TextFilter,
13 TextGatherMode, TextIndexStats, TextQueryMode, TextSearchHit, TextSearchOptions,
14 TextSearchRequest, TextTermStats, TextTermStatsRequest,
15};
16use khive_storage::usage::{UsageContext, UsageUnit};
17use khive_storage::StorageCapability;
18use khive_storage::TextSearch;
19use khive_types::SubstrateKind;
20
21use crate::error::SqliteError;
22use crate::pool::ConnectionPool;
23use crate::sql_bridge::bind_params;
24use crate::writer_task::WriterTaskHandle;
25
26pub fn rowid_map_table(table: &str) -> String {
37 format!("{table}_rowids")
38}
39
40pub fn rowid_map_state_table(table: &str) -> String {
45 format!("{}_state", rowid_map_table(table))
46}
47
48pub const ROWID_MAP_BACKFILL_COMPLETE: &str = "complete";
52
53pub fn rowid_map_ddl(table: &str) -> String {
58 let map = rowid_map_table(table);
59 let state = rowid_map_state_table(table);
60 format!(
61 "CREATE TABLE IF NOT EXISTS {map} (\
62 namespace TEXT NOT NULL, \
63 subject_id TEXT NOT NULL, \
64 rowid INTEGER NOT NULL, \
65 PRIMARY KEY (namespace, subject_id)\
66 ) WITHOUT ROWID; \
67 CREATE TABLE IF NOT EXISTS {state} (\
68 key TEXT PRIMARY KEY, \
69 value TEXT NOT NULL\
70 ) WITHOUT ROWID"
71 )
72}
73
74pub fn delete_document_statement(table: &str, namespace: &str, subject_id: Uuid) -> SqlStatement {
97 let map = rowid_map_table(table);
98 SqlStatement {
99 sql: format!(
100 "DELETE FROM {table} WHERE rowid IN \
101 (SELECT rowid FROM {map} WHERE namespace = ?1 AND subject_id = ?2) \
102 AND namespace = ?1 AND subject_id = ?2"
103 ),
104 params: vec![
105 SqlValue::Text(namespace.to_string()),
106 SqlValue::Text(subject_id.to_string()),
107 ],
108 label: Some(format!("fts-delete-{table}")),
109 }
110}
111
112fn delete_document_statement_scan_fallback(
119 table: &str,
120 namespace: &str,
121 subject_id: Uuid,
122) -> SqlStatement {
123 SqlStatement {
124 sql: format!("DELETE FROM {table} WHERE namespace = ?1 AND subject_id = ?2"),
125 params: vec![
126 SqlValue::Text(namespace.to_string()),
127 SqlValue::Text(subject_id.to_string()),
128 ],
129 label: Some(format!("fts-delete-scan-fallback-{table}")),
130 }
131}
132
133pub fn delete_document_map_statement(
140 table: &str,
141 namespace: &str,
142 subject_id: Uuid,
143) -> SqlStatement {
144 let map = rowid_map_table(table);
145 SqlStatement {
146 sql: format!("DELETE FROM {map} WHERE namespace = ?1 AND subject_id = ?2"),
147 params: vec![
148 SqlValue::Text(namespace.to_string()),
149 SqlValue::Text(subject_id.to_string()),
150 ],
151 label: Some(format!("fts-delete-map-{table}")),
152 }
153}
154
155pub fn delete_document_statements(
163 table: &str,
164 namespace: &str,
165 subject_id: Uuid,
166) -> [SqlStatement; 2] {
167 [
168 delete_document_statement(table, namespace, subject_id),
169 delete_document_map_statement(table, namespace, subject_id),
170 ]
171}
172
173pub fn insert_document_statement(table: &str, document: &TextDocument) -> SqlStatement {
182 let tags_json = tags_to_json(&document.tags);
183 let metadata_json = document.metadata.as_ref().map(|v| v.to_string());
184 SqlStatement {
185 sql: format!(
186 "INSERT INTO {table} \
187 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, record_kind) \
188 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)"
189 ),
190 params: vec![
191 SqlValue::Text(document.subject_id.to_string()),
192 SqlValue::Text(document.kind.to_string()),
193 SqlValue::Text(document.title.clone().unwrap_or_default()),
194 SqlValue::Text(document.body.clone()),
195 SqlValue::Text(tags_json),
196 SqlValue::Text(document.namespace.clone()),
197 match metadata_json {
198 Some(m) => SqlValue::Text(m),
199 None => SqlValue::Null,
200 },
201 SqlValue::Integer(dt_to_micros(&document.updated_at)),
202 match &document.record_kind {
203 Some(kind) => SqlValue::Text(kind.clone()),
204 None => SqlValue::Null,
205 },
206 ],
207 label: Some(format!("fts-insert-{table}")),
208 }
209}
210
211pub fn insert_document_map_statement(
221 table: &str,
222 namespace: &str,
223 subject_id: Uuid,
224) -> SqlStatement {
225 let map = rowid_map_table(table);
226 SqlStatement {
227 sql: format!(
228 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
229 VALUES (?1, ?2, last_insert_rowid())"
230 ),
231 params: vec![
232 SqlValue::Text(namespace.to_string()),
233 SqlValue::Text(subject_id.to_string()),
234 ],
235 label: Some(format!("fts-insert-map-{table}")),
236 }
237}
238
239pub fn insert_document_statements(table: &str, document: &TextDocument) -> [SqlStatement; 2] {
246 [
247 insert_document_statement(table, document),
248 insert_document_map_statement(table, &document.namespace, document.subject_id),
249 ]
250}
251
252#[cfg(test)]
257pub(crate) fn ensure_fts5_schema(
258 conn: &rusqlite::Connection,
259 table_key: &str,
260) -> Result<(), rusqlite::Error> {
261 let table_name = format!("fts_{}", table_key);
262 let ddl = format!(
263 "CREATE VIRTUAL TABLE IF NOT EXISTS {} USING fts5(\
264 subject_id UNINDEXED, \
265 kind UNINDEXED, \
266 title, \
267 body, \
268 tags UNINDEXED, \
269 namespace UNINDEXED, \
270 metadata UNINDEXED, \
271 updated_at UNINDEXED, \
272 record_kind\
273 )",
274 table_name
275 );
276 conn.execute_batch(&ddl)?;
277 conn.execute_batch(&rowid_map_ddl(&table_name))
278}
279
280fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
281 StorageError::driver(StorageCapability::Text, op, e)
282}
283
284fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
285 StorageError::driver(StorageCapability::Text, op, e)
286}
287
288fn count_fts_pass(context: Option<&UsageContext>) {
289 if let Some(context) = context {
290 context.add(UsageUnit::FtsPasses, 1);
291 }
292}
293
294pub struct Fts5TextSearch {
301 pool: Arc<ConnectionPool>,
302 is_file_backed: bool,
303 table_name: String,
304 writer_task: Option<WriterTaskHandle>,
305 scan_fallback: bool,
312}
313
314impl Fts5TextSearch {
315 pub(crate) fn new(pool: Arc<ConnectionPool>, is_file_backed: bool, table_key: String) -> Self {
320 Self::new_with_mode(pool, is_file_backed, table_key, false)
321 }
322
323 pub(crate) fn new_scan_fallback(
332 pool: Arc<ConnectionPool>,
333 is_file_backed: bool,
334 table_key: String,
335 ) -> Self {
336 Self::new_with_mode(pool, is_file_backed, table_key, true)
337 }
338
339 fn new_with_mode(
340 pool: Arc<ConnectionPool>,
341 is_file_backed: bool,
342 table_key: String,
343 scan_fallback: bool,
344 ) -> Self {
345 let table_name = format!("fts_{}", table_key);
346 let writer_task = pool.writer_task_handle().ok().flatten();
352 Self {
353 pool,
354 is_file_backed,
355 table_name,
356 writer_task,
357 scan_fallback,
358 }
359 }
360
361 fn open_standalone_writer(&self) -> Result<rusqlite::Connection, StorageError> {
362 self.pool
363 .open_standalone_writer()
364 .map_err(|e| map_sqlite_err(e, "open_fts_writer"))
365 }
366
367 fn current_writer_task(
377 &self,
378 operation: &'static str,
379 ) -> Result<Option<WriterTaskHandle>, StorageError> {
380 self.pool
381 .writer_task_for_write(self.writer_task.as_ref(), operation)
382 }
383
384 async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
392 where
393 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
394 R: Send + 'static,
395 {
396 if let Some(writer_task) = self.current_writer_task(op)? {
397 return writer_task
398 .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
399 .await;
400 }
401
402 self.pool
403 .record_direct_route(crate::timeout_sink::Site::DirectRouteFtsGeneralWrite);
404 self.with_writer_unmanaged(op, f).await
405 }
406
407 async fn with_writer_unmanaged<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
416 where
417 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
418 R: Send + 'static,
419 {
420 let result = if self.is_file_backed {
421 let conn = self.open_standalone_writer()?;
422 tokio::task::spawn_blocking(move || f(&conn).map_err(|e| map_err(e, op)))
423 .await
424 .map_err(|e| StorageError::driver(StorageCapability::Text, op, e))?
425 } else {
426 let pool = Arc::clone(&self.pool);
427 tokio::task::spawn_blocking(move || {
428 let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
429 f(guard.conn()).map_err(|e| map_err(e, op))
430 })
431 .await
432 .map_err(|e| StorageError::driver(StorageCapability::Text, op, e))?
433 };
434 if let Err(err) = &result {
435 let msg = err.to_string();
436 if msg.contains("locked") || msg.contains("busy") {
437 if self.is_file_backed {
438 crate::timeout_sink::emit_timeout(
443 &crate::timeout_sink::db_label(&self.pool),
444 crate::timeout_sink::Site::StandaloneText,
445 &msg,
446 None,
447 );
448 }
449 let open: Vec<String> = khive_storage::tx_registry::snapshot()
450 .into_iter()
451 .map(|(age, label)| {
452 format!(
453 "{}@{}ms",
454 label.as_deref().unwrap_or("unlabeled"),
455 age.as_millis()
456 )
457 })
458 .collect();
459 tracing::warn!(
460 op,
461 open_tx_count = open.len(),
462 open_txs = %open.join(","),
463 "text write starved on the SQLite write lock; open registered \
464 transactions listed (empty means the holder issued no \
465 registered BEGIN IMMEDIATE in this process)"
466 );
467 }
468 }
469 result
470 }
471
472 async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
473 where
474 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
475 R: Send + 'static,
476 {
477 super::run_pooled_store_read(
478 Arc::clone(&self.pool),
479 StorageCapability::Text,
480 op,
481 move |conn| f(conn).map_err(|error| map_err(error, op)),
482 )
483 .await
484 }
485}
486
487fn tags_to_json(tags: &[String]) -> String {
490 serde_json::to_string(tags).unwrap_or_else(|_| "[]".to_string())
491}
492
493fn tags_from_json(s: &str) -> Vec<String> {
494 serde_json::from_str(s).unwrap_or_default()
495}
496
497fn dt_to_micros(dt: &DateTime<Utc>) -> i64 {
498 dt.timestamp_micros()
499}
500
501fn micros_to_dt(micros: i64) -> DateTime<Utc> {
502 Utc.timestamp_micros(micros)
503 .single()
504 .unwrap_or_else(Utc::now)
505}
506
507fn sanitize_fts5_query(query: &str) -> String {
514 let spaced: String = query
518 .chars()
519 .map(|c| {
520 if matches!(c, '(' | ')' | ',' | ':' | '-' | '.' | '/') {
521 ' '
522 } else {
523 c
524 }
525 })
526 .collect();
527
528 let sanitized: String = spaced
535 .chars()
536 .filter(|c| {
537 !matches!(c, '*' | '"' | '\'' | '+' | '^' | '~' | '!' | '$' | '\0') && !c.is_control()
538 })
539 .collect();
540
541 sanitized
543 .split_whitespace()
544 .filter(|t| {
545 !matches!(
546 t.to_ascii_uppercase().as_str(),
547 "AND" | "OR" | "NOT" | "NEAR"
548 )
549 })
550 .collect::<Vec<_>>()
551 .join(" ")
552}
553
554fn sanitize_fts5_query_legacy_merged(query: &str) -> String {
561 let spaced: String = query
562 .chars()
563 .map(|c| {
564 if matches!(c, '(' | ')' | ',' | ':' | '/') {
565 ' '
566 } else {
567 c
568 }
569 })
570 .collect();
571
572 let sanitized: String = spaced
573 .chars()
574 .filter(|c| {
575 !matches!(
576 c,
577 '*' | '"' | '\'' | '+' | '-' | '^' | '.' | '~' | '!' | '$' | '\0'
578 ) && !c.is_control()
579 })
580 .collect();
581
582 sanitized
583 .split_whitespace()
584 .filter(|t| {
585 !matches!(
586 t.to_ascii_uppercase().as_str(),
587 "AND" | "OR" | "NOT" | "NEAR"
588 )
589 })
590 .collect::<Vec<_>>()
591 .join(" ")
592}
593
594fn sanitize_fts5_phrase_literal(token: &str) -> Option<String> {
603 let mut literal = String::with_capacity(token.len());
604 for c in token
605 .chars()
606 .filter(|c| !matches!(c, '*' | '\0') && !c.is_control())
607 {
608 if c == '"' {
609 literal.push_str("\"\"");
610 } else {
611 literal.push(c);
612 }
613 }
614 if literal.is_empty() {
615 None
616 } else {
617 Some(literal)
618 }
619}
620
621const FTS5_TRIGRAM_MIN_SAFE_LEN: usize = 3;
629
630fn is_fts5_bareword_safe(s: &str) -> bool {
634 !s.is_empty() && s.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
635}
636
637fn sanitize_fts5_token_group(token: &str) -> Option<String> {
645 let split = sanitize_fts5_query(token);
646 let split_terms: Vec<&str> = split.split_whitespace().collect();
647 if split_terms.is_empty() {
648 return None;
649 }
650
651 let all_bareword_safe = split_terms.iter().all(|t| is_fts5_bareword_safe(t));
652 if split_terms.len() == 1 && is_fts5_bareword_safe(token) {
653 return Some(split_terms[0].to_string());
654 }
655
656 let has_trigram_unsafe_segment = split_terms
657 .iter()
658 .any(|t| t.chars().count() < FTS5_TRIGRAM_MIN_SAFE_LEN);
659
660 let mut alternatives = Vec::new();
661 if all_bareword_safe && !has_trigram_unsafe_segment {
662 alternatives.push(format!("({})", split_terms.join(" ")));
663 }
664
665 let merged = sanitize_fts5_query_legacy_merged(token);
676 let merged_terms: Vec<&str> = merged.split_whitespace().collect();
677 let merged_all_bareword_safe =
678 !merged_terms.is_empty() && merged_terms.iter().all(|t| is_fts5_bareword_safe(t));
679 let merged_has_unsafe_segment = merged_terms.len() > 1
680 && merged_terms
681 .iter()
682 .any(|t| t.chars().count() < FTS5_TRIGRAM_MIN_SAFE_LEN);
683 let phrase = sanitize_fts5_phrase_literal(token);
692 let duplicates_split = merged == split_terms.join(" ");
693 let duplicates_phrase = phrase.as_deref() == Some(merged.as_str());
694 if merged_all_bareword_safe
695 && !merged.is_empty()
696 && !duplicates_split
697 && !duplicates_phrase
698 && !merged_has_unsafe_segment
699 {
700 alternatives.push(merged);
701 }
702
703 if let Some(phrase) = phrase {
704 alternatives.push(format!("\"{}\"", phrase));
705 }
706
707 match alternatives.len() {
708 0 => None,
709 1 => alternatives.into_iter().next(),
710 _ => Some(format!("({})", alternatives.join(" OR "))),
711 }
712}
713
714fn join_plain_groups(groups: &[String]) -> String {
729 let mut expr = String::new();
730 for (i, group) in groups.iter().enumerate() {
731 if i > 0 {
732 let prev_compound = groups[i - 1].starts_with('(');
733 let this_compound = group.starts_with('(');
734 expr.push_str(if prev_compound || this_compound {
735 " AND "
736 } else {
737 " "
738 });
739 }
740 expr.push_str(group);
741 }
742 expr
743}
744
745fn build_match_expr(query: &str, mode: TextQueryMode) -> Option<String> {
762 match mode {
763 TextQueryMode::AnyTerm => {
764 let groups: Vec<String> = query
765 .split_whitespace()
766 .filter_map(sanitize_fts5_token_group)
767 .collect();
768 if groups.is_empty() {
769 None
770 } else {
771 Some(groups.join(" OR "))
772 }
773 }
774 TextQueryMode::Plain => {
775 let groups: Vec<String> = query
776 .split_whitespace()
777 .filter_map(sanitize_fts5_token_group)
778 .collect();
779 if groups.is_empty() {
780 None
781 } else {
782 Some(join_plain_groups(&groups))
783 }
784 }
785 TextQueryMode::Phrase => {
786 sanitize_fts5_phrase_literal(query).map(|literal| format!("\"{}\"", literal))
787 }
788 }
789}
790
791fn quote_fts5_phrase(value: &str) -> String {
792 format!("\"{}\"", value.replace('"', "\"\""))
793}
794
795fn record_kind_match_expr(filter: Option<&TextFilter>) -> Option<String> {
802 let record_kinds = &filter?.record_kinds;
803 if record_kinds.is_empty() || record_kinds.iter().any(|kind| kind.chars().count() < 3) {
804 return None;
805 }
806
807 let clauses: Vec<String> = record_kinds
808 .iter()
809 .map(|kind| format!("record_kind : {}", quote_fts5_phrase(kind)))
810 .collect();
811 if clauses.len() == 1 {
812 clauses.into_iter().next()
813 } else {
814 Some(format!("({})", clauses.join(" OR ")))
815 }
816}
817
818fn build_filtered_match_expr(
819 query: &str,
820 mode: TextQueryMode,
821 filter: Option<&TextFilter>,
822) -> Option<String> {
823 let query = format!("{{title body}} : ({})", build_match_expr(query, mode)?);
827 match record_kind_match_expr(filter) {
828 Some(classifier) => Some(format!("{classifier} AND ({query})")),
829 None => Some(query),
830 }
831}
832
833const LEXICAL_BM25_RANK: &str = "bm25(0.0, 0.0, 1.0, 1.0, 0.0, 0.0, 0.0, 0.0, 0.0)";
840
841fn build_filter_clause(
846 filter: &TextFilter,
847 table: &str,
848 start_idx: usize,
849) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
850 let mut conditions: Vec<String> = Vec::new();
851 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
852 let mut idx = start_idx;
853
854 if !filter.ids.is_empty() {
855 let placeholders: Vec<String> = filter
856 .ids
857 .iter()
858 .map(|_| {
859 let p = format!("?{}", idx);
860 idx += 1;
861 p
862 })
863 .collect();
864 conditions.push(format!(
865 "{}.subject_id IN ({})",
866 table,
867 placeholders.join(", ")
868 ));
869 for id in &filter.ids {
870 params.push(Box::new(id.to_string()));
871 }
872 }
873
874 if !filter.kinds.is_empty() {
875 let placeholders: Vec<String> = filter
876 .kinds
877 .iter()
878 .map(|_| {
879 let p = format!("?{}", idx);
880 idx += 1;
881 p
882 })
883 .collect();
884 conditions.push(format!("{}.kind IN ({})", table, placeholders.join(", ")));
885 for kind in &filter.kinds {
886 params.push(Box::new(kind.to_string()));
887 }
888 }
889
890 if !filter.record_kinds.is_empty() {
891 let placeholders: Vec<String> = filter
892 .record_kinds
893 .iter()
894 .map(|_| {
895 let p = format!("?{}", idx);
896 idx += 1;
897 p
898 })
899 .collect();
900 conditions.push(format!(
901 "{}.record_kind IN ({})",
902 table,
903 placeholders.join(", ")
904 ));
905 for kind in &filter.record_kinds {
906 params.push(Box::new(kind.clone()));
907 }
908 }
909
910 if !filter.namespaces.is_empty() {
911 let placeholders: Vec<String> = filter
912 .namespaces
913 .iter()
914 .map(|_| {
915 let p = format!("?{}", idx);
916 idx += 1;
917 p
918 })
919 .collect();
920 conditions.push(format!(
921 "{}.namespace IN ({})",
922 table,
923 placeholders.join(", ")
924 ));
925 for ns in &filter.namespaces {
926 params.push(Box::new(ns.clone()));
927 }
928 }
929
930 if conditions.is_empty() {
931 (String::new(), params)
932 } else {
933 (format!(" AND {}", conditions.join(" AND ")), params)
934 }
935}
936
937#[cfg(test)]
961pub(crate) mod delete_dml_test_seam {
962 use std::collections::HashSet;
963 use std::sync::Mutex;
964 use uuid::Uuid;
965
966 pub(crate) static TARGETS: Mutex<Option<HashSet<(String, Uuid)>>> = Mutex::new(None);
967
968 pub(crate) fn arm(namespace: &str, subject_id: Uuid) {
969 TARGETS
970 .lock()
971 .unwrap()
972 .get_or_insert_with(HashSet::new)
973 .insert((namespace.to_string(), subject_id));
974 }
975
976 pub(crate) fn disarm(namespace: &str, subject_id: Uuid) {
977 if let Some(targets) = TARGETS.lock().unwrap().as_mut() {
978 targets.remove(&(namespace.to_string(), subject_id));
979 }
980 }
981
982 pub(crate) fn matches(namespace: &str, subject_id: Uuid) -> bool {
983 TARGETS
984 .lock()
985 .unwrap()
986 .as_ref()
987 .is_some_and(|targets| targets.contains(&(namespace.to_string(), subject_id)))
988 }
989}
990
991fn delete_document_dml(
1000 conn: &rusqlite::Connection,
1001 table: &str,
1002 namespace: &str,
1003 subject_id: Uuid,
1004) -> Result<bool, SqliteError> {
1005 let [fts_statement, map_statement] = delete_document_statements(table, namespace, subject_id);
1006
1007 let mut stmt = conn.prepare(&fts_statement.sql)?;
1008 bind_params(&mut stmt, &fts_statement.params)?;
1009 let deleted = stmt.raw_execute()? > 0;
1010 drop(stmt);
1011
1012 #[cfg(test)]
1013 if delete_dml_test_seam::matches(namespace, subject_id) {
1014 return Err(SqliteError::InvalidData(
1015 "delete_document_dml test seam: forced failure between the FTS-row \
1016 delete and the map-row delete"
1017 .to_string(),
1018 ));
1019 }
1020
1021 let mut map_stmt = conn.prepare(&map_statement.sql)?;
1022 bind_params(&mut map_stmt, &map_statement.params)?;
1023 map_stmt.raw_execute()?;
1024
1025 Ok(deleted)
1026}
1027
1028fn upsert_document_dml(
1034 conn: &rusqlite::Connection,
1035 table: &str,
1036 document: &TextDocument,
1037) -> Result<(), rusqlite::Error> {
1038 let statement = delete_document_statement(table, &document.namespace, document.subject_id);
1042 let mut stmt = conn.prepare(&statement.sql)?;
1043 bind_params(&mut stmt, &statement.params)?;
1044 stmt.raw_execute()?;
1045
1046 for statement in insert_document_statements(table, document) {
1047 let mut stmt = conn.prepare(&statement.sql)?;
1048 bind_params(&mut stmt, &statement.params)?;
1049 stmt.raw_execute()?;
1050 }
1051 Ok(())
1052}
1053
1054fn batch_upsert_documents_dml(
1063 conn: &rusqlite::Connection,
1064 table: &str,
1065 documents: &[TextDocument],
1066 attempted: u64,
1067) -> Result<BatchWriteSummary, rusqlite::Error> {
1068 let map = rowid_map_table(table);
1069 let del_sql = format!(
1076 "DELETE FROM {table} WHERE rowid IN \
1077 (SELECT rowid FROM {map} WHERE namespace = ?1 AND subject_id = ?2) \
1078 AND namespace = ?1 AND subject_id = ?2"
1079 );
1080 let ins_sql = format!(
1081 "INSERT INTO {} \
1082 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, record_kind) \
1083 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
1084 table
1085 );
1086 let map_ins_sql = format!(
1089 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
1090 VALUES (?1, ?2, last_insert_rowid())"
1091 );
1092
1093 let mut summary = BatchWriteSummary {
1094 attempted,
1095 ..BatchWriteSummary::default()
1096 };
1097
1098 for (index, doc) in documents.iter().enumerate() {
1099 conn.execute_batch("SAVEPOINT fts_upsert_doc")?;
1100 let id_str = doc.subject_id.to_string();
1101 let namespace = &doc.namespace;
1102 let result = (|| {
1103 conn.execute(&del_sql, rusqlite::params![namespace, &id_str])?;
1104
1105 let tags_json = tags_to_json(&doc.tags);
1106 let metadata_json: Option<String> = doc.metadata.as_ref().map(|v| v.to_string());
1107
1108 conn.execute(
1109 &ins_sql,
1110 rusqlite::params![
1111 &id_str,
1112 &doc.kind.to_string(),
1113 doc.title.as_deref().unwrap_or(""),
1114 &doc.body,
1115 &tags_json,
1116 namespace,
1117 &metadata_json,
1118 dt_to_micros(&doc.updated_at),
1119 &doc.record_kind,
1120 ],
1121 )?;
1122 conn.execute(&map_ins_sql, rusqlite::params![namespace, &id_str])?;
1123 Ok::<(), rusqlite::Error>(())
1124 })();
1125
1126 match result {
1127 Ok(()) => {
1128 conn.execute_batch("RELEASE SAVEPOINT fts_upsert_doc")?;
1129 summary.affected = summary.affected.saturating_add(1);
1130 }
1131 Err(e) => {
1132 let _ = conn.execute_batch("ROLLBACK TO SAVEPOINT fts_upsert_doc");
1133 let _ = conn.execute_batch("RELEASE SAVEPOINT fts_upsert_doc");
1134 let (class, retryability) = super::classify_batch_sqlite_error(&e);
1135 summary.record_failure(index, Some(id_str), class, retryability, e.to_string());
1136 }
1137 }
1138 }
1139
1140 Ok(summary)
1141}
1142
1143#[async_trait]
1144impl TextSearch for Fts5TextSearch {
1145 async fn upsert_document(&self, document: TextDocument) -> Result<(), StorageError> {
1146 let table = self.table_name.clone();
1147
1148 if let Some(writer_task) = self.current_writer_task("fts_upsert")? {
1159 let table2 = table.clone();
1160 return writer_task
1161 .send_bounded(move |conn| {
1162 upsert_document_dml(conn, &table2, &document)
1163 .map_err(|e| map_err(e, "fts_upsert"))
1164 })
1165 .await;
1166 }
1167
1168 let origin = self.pool.origin();
1171 self.with_writer("fts_upsert", move |conn| {
1172 conn.execute_batch("BEGIN IMMEDIATE")?;
1173 let _tx_handle = khive_storage::tx_registry::register_scoped(
1174 Some("text_upsert_document".to_string()),
1175 origin,
1176 );
1177
1178 if let Err(e) = upsert_document_dml(conn, &table, &document) {
1179 let _ = conn.execute_batch("ROLLBACK");
1180 return Err(e);
1181 }
1182
1183 conn.execute_batch("COMMIT")?;
1184 Ok(())
1185 })
1186 .await
1187 }
1188
1189 async fn upsert_documents(
1190 &self,
1191 documents: Vec<TextDocument>,
1192 ) -> Result<BatchWriteSummary, StorageError> {
1193 let table = self.table_name.clone();
1194 let attempted = documents.len() as u64;
1195
1196 if let Some(writer_task) = self.current_writer_task("fts_upsert_batch")? {
1204 let table2 = table.clone();
1205 return writer_task
1206 .send_bounded(move |conn| {
1207 batch_upsert_documents_dml(conn, &table2, &documents, attempted)
1208 .map_err(|e| map_err(e, "fts_upsert_batch"))
1209 })
1210 .await;
1211 }
1212
1213 let origin = self.pool.origin();
1216 self.with_writer("fts_upsert_batch", move |conn| {
1217 conn.execute_batch("BEGIN IMMEDIATE")?;
1218 let _tx_handle = khive_storage::tx_registry::register_scoped(
1219 Some("text_upsert_batch".to_string()),
1220 origin,
1221 );
1222
1223 let summary = batch_upsert_documents_dml(conn, &table, &documents, attempted)?;
1224
1225 conn.execute_batch("COMMIT")?;
1226
1227 Ok(summary)
1228 })
1229 .await
1230 }
1231
1232 async fn delete_document(
1233 &self,
1234 namespace: &str,
1235 subject_id: Uuid,
1236 ) -> Result<bool, StorageError> {
1237 let table = self.table_name.clone();
1238 let namespace = namespace.to_string();
1239
1240 if self.scan_fallback {
1241 let statement = delete_document_statement_scan_fallback(&table, &namespace, subject_id);
1242 return self
1243 .with_writer("fts_delete_scan_fallback", move |conn| {
1244 let mut stmt = conn.prepare(&statement.sql)?;
1245 bind_params(&mut stmt, &statement.params)?;
1246 Ok(stmt.raw_execute()? > 0)
1247 })
1248 .await;
1249 }
1250
1251 if let Some(writer_task) = self.current_writer_task("fts_delete")? {
1259 let table2 = table.clone();
1260 let namespace2 = namespace.clone();
1261 return writer_task
1262 .send_bounded(move |conn| {
1263 delete_document_dml(conn, &table2, &namespace2, subject_id)
1264 .map_err(|e| map_sqlite_err(e, "fts_delete"))
1265 })
1266 .await;
1267 }
1268
1269 let origin = self.pool.origin();
1270 self.with_writer("fts_delete", move |conn| {
1271 conn.execute_batch("BEGIN IMMEDIATE")?;
1272 let _tx_handle = khive_storage::tx_registry::register_scoped(
1273 Some("text_delete_document".to_string()),
1274 origin,
1275 );
1276
1277 match delete_document_dml(conn, &table, &namespace, subject_id) {
1278 Ok(deleted) => {
1279 conn.execute_batch("COMMIT")?;
1280 Ok(deleted)
1281 }
1282 Err(e) => {
1283 let _ = conn.execute_batch("ROLLBACK");
1284 Err(match e {
1285 SqliteError::Rusqlite(inner) => inner,
1286 other => rusqlite::Error::InvalidParameterName(other.to_string()),
1287 })
1288 }
1289 }
1290 })
1291 .await
1292 }
1293
1294 async fn get_document(
1295 &self,
1296 namespace: &str,
1297 subject_id: Uuid,
1298 ) -> Result<Option<TextDocument>, StorageError> {
1299 let namespace = namespace.to_string();
1300 let table = self.table_name.clone();
1301 let scan_fallback = self.scan_fallback;
1302
1303 self.with_reader("fts_get", move |conn| {
1304 let sql = if scan_fallback {
1305 format!(
1306 "SELECT subject_id, kind, title, body, tags, namespace, \
1307 metadata, updated_at, record_kind \
1308 FROM {table} WHERE namespace = ?1 AND subject_id = ?2"
1309 )
1310 } else {
1311 let map = rowid_map_table(&table);
1312 format!(
1313 "SELECT t.subject_id, t.kind, t.title, t.body, t.tags, t.namespace, \
1314 t.metadata, t.updated_at, t.record_kind \
1315 FROM {table} AS t JOIN {map} AS m ON m.rowid = t.rowid \
1316 WHERE m.namespace = ?1 AND m.subject_id = ?2 \
1317 AND t.namespace = ?1 AND t.subject_id = ?2"
1318 )
1319 };
1320 let mut stmt = conn.prepare(&sql)?;
1321 let mut rows = stmt.query(rusqlite::params![namespace, subject_id.to_string()])?;
1322
1323 match rows.next()? {
1324 Some(row) => {
1325 let id_str: String = row.get(0)?;
1326 let kind_str: String = row.get(1)?;
1327 let title: String = row.get(2)?;
1328 let body: String = row.get(3)?;
1329 let tags_json: String = row.get(4)?;
1330 let ns: String = row.get(5)?;
1331 let metadata_json: Option<String> = row.get(6)?;
1332 let updated_at_micros: i64 = row.get(7)?;
1333 let record_kind: Option<String> = row.get(8)?;
1334
1335 let sid = Uuid::parse_str(&id_str).map_err(|e| {
1336 rusqlite::Error::FromSqlConversionFailure(
1337 0,
1338 rusqlite::types::Type::Text,
1339 Box::new(e),
1340 )
1341 })?;
1342
1343 let kind = kind_str.parse::<SubstrateKind>().map_err(|e| {
1344 rusqlite::Error::FromSqlConversionFailure(
1345 1,
1346 rusqlite::types::Type::Text,
1347 Box::new(e),
1348 )
1349 })?;
1350
1351 Ok(Some(TextDocument {
1352 subject_id: sid,
1353 kind,
1354 record_kind,
1355 title: if title.is_empty() { None } else { Some(title) },
1356 body,
1357 tags: tags_from_json(&tags_json),
1358 namespace: ns,
1359 metadata: metadata_json.and_then(|s| serde_json::from_str(&s).ok()),
1360 updated_at: micros_to_dt(updated_at_micros),
1361 }))
1362 }
1363 None => Ok(None),
1364 }
1365 })
1366 .await
1367 }
1368
1369 async fn search(&self, request: TextSearchRequest) -> Result<Vec<TextSearchHit>, StorageError> {
1370 let table = self.table_name.clone();
1371 let usage = khive_storage::usage::current();
1372
1373 let result = self
1374 .with_reader("fts_search", move |conn| {
1375 let match_expr = match build_filtered_match_expr(
1376 &request.query,
1377 request.mode,
1378 request.filter.as_ref(),
1379 ) {
1380 Some(expr) => expr,
1381 None => return Ok(Vec::new()),
1382 };
1383
1384 let snippet_expr = if request.snippet_chars == 0 {
1391 "NULL AS snippet".to_string()
1392 } else {
1393 let chars = i32::try_from(request.snippet_chars).unwrap_or(i32::MAX);
1394 format!("snippet({table}, 3, '', '', '...', {chars})")
1395 };
1396
1397 let (filter_clause, filter_params) = if let Some(ref filter) = request.filter {
1398 build_filter_clause(filter, &table, 3)
1399 } else {
1400 (String::new(), Vec::new())
1401 };
1402
1403 let sql = format!(
1404 "SELECT subject_id, rank, title, {snippet_expr} \
1405 FROM {table} WHERE {table} MATCH ?1 \
1406 AND rank MATCH '{LEXICAL_BM25_RANK}'{filter_clause} \
1407 ORDER BY rank LIMIT ?2",
1408 );
1409
1410 let mut stmt = conn.prepare(&sql)?;
1411 count_fts_pass(usage.as_ref());
1412 stmt.raw_bind_parameter(1, &match_expr)?;
1413 stmt.raw_bind_parameter(2, request.top_k as i64)?;
1414
1415 for (i, param) in filter_params.iter().enumerate() {
1416 param
1417 .to_sql()
1418 .map(|val| stmt.raw_bind_parameter(3 + i, val))
1419 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1420 }
1421
1422 let mut hits = Vec::new();
1423 let mut rows = stmt.raw_query();
1424 let mut rank_idx = 0u32;
1425
1426 while let Some(row) = rows.next()? {
1427 let id_str: String = row.get(0)?;
1428 let fts_rank: f64 = row.get(1)?;
1429 let title: String = row.get(2)?;
1430 let snippet: Option<String> = row.get(3)?;
1431
1432 let subject_id = Uuid::parse_str(&id_str).map_err(|e| {
1433 rusqlite::Error::FromSqlConversionFailure(
1434 0,
1435 rusqlite::types::Type::Text,
1436 Box::new(e),
1437 )
1438 })?;
1439
1440 rank_idx += 1;
1441 hits.push((subject_id, fts_rank, rank_idx, title, snippet));
1442 }
1443
1444 let min_rank = hits.iter().map(|h| h.1).fold(f64::INFINITY, f64::min);
1447 let max_rank = hits.iter().map(|h| h.1).fold(f64::NEG_INFINITY, f64::max);
1448 let range = max_rank - min_rank;
1449
1450 let results = hits
1451 .into_iter()
1452 .map(|(subject_id, raw_rank, rank, title, snippet)| {
1453 let score = if range.abs() < 1e-12 {
1454 1.0
1455 } else {
1456 let t = (max_rank - raw_rank) / range;
1457 0.05 + 0.95 * t
1458 };
1459 TextSearchHit {
1460 subject_id,
1461 score: DeterministicScore::from_f64(score),
1462 rank,
1463 title: if title.is_empty() { None } else { Some(title) },
1464 snippet: snippet.filter(|s| !s.is_empty()),
1465 }
1466 })
1467 .collect();
1468
1469 Ok(results)
1470 })
1471 .await;
1472
1473 result
1474 }
1475
1476 async fn count(&self, filter: TextFilter) -> Result<u64, StorageError> {
1477 let table = self.table_name.clone();
1478
1479 self.with_reader("fts_count", move |conn| {
1480 let indexed_classifier = record_kind_match_expr(Some(&filter));
1481 let filter_start = if indexed_classifier.is_some() { 2 } else { 1 };
1482 let (filter_clause, filter_params) = build_filter_clause(&filter, &table, filter_start);
1483
1484 let sql = if indexed_classifier.is_some() {
1485 format!("SELECT COUNT(*) FROM {table} WHERE {table} MATCH ?1{filter_clause}")
1486 } else if filter_clause.is_empty() {
1487 format!("SELECT COUNT(*) FROM {table}")
1488 } else {
1489 let where_part = filter_clause.trim_start_matches(" AND ");
1490 format!("SELECT COUNT(*) FROM {table} WHERE {where_part}")
1491 };
1492
1493 let mut stmt = conn.prepare(&sql)?;
1494
1495 if let Some(classifier) = indexed_classifier {
1496 stmt.raw_bind_parameter(1, classifier)?;
1497 }
1498
1499 for (i, param) in filter_params.iter().enumerate() {
1500 param
1501 .to_sql()
1502 .map(|val| stmt.raw_bind_parameter(filter_start + i, val))
1503 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1504 }
1505
1506 let mut rows = stmt.raw_query();
1507 match rows.next()? {
1508 Some(row) => {
1509 let count: i64 = row.get(0)?;
1510 Ok(count as u64)
1511 }
1512 None => Ok(0),
1513 }
1514 })
1515 .await
1516 }
1517
1518 async fn stats(&self) -> Result<TextIndexStats, StorageError> {
1519 let table = self.table_name.clone();
1520
1521 self.with_reader("fts_stats", move |conn| {
1522 let sql = format!("SELECT COUNT(*) FROM {}", table);
1523 let count: i64 = conn.query_row(&sql, [], |row| row.get(0))?;
1524
1525 Ok(TextIndexStats {
1526 document_count: count as u64,
1527 needs_rebuild: false,
1528 last_rebuild_at: None,
1529 })
1530 })
1531 .await
1532 }
1533
1534 async fn search_with_options(
1535 &self,
1536 request: TextSearchRequest,
1537 options: TextSearchOptions,
1538 ) -> Result<Vec<TextSearchHit>, StorageError> {
1539 match options.gather_mode {
1540 TextGatherMode::Ranked => self.search(request).await,
1541 TextGatherMode::Unranked => self.search_unranked(request).await,
1542 TextGatherMode::RankWithinCap => {
1543 let gather_limit = options
1544 .gather_limit
1545 .unwrap_or(request.top_k)
1546 .max(request.top_k);
1547 self.search_rank_within_cap(request, gather_limit).await
1548 }
1549 }
1550 }
1551
1552 async fn term_stats(
1553 &self,
1554 request: TextTermStatsRequest,
1555 ) -> Result<Vec<TextTermStats>, StorageError> {
1556 let table = self.table_name.clone();
1557
1558 self.with_reader("fts_term_stats", move |conn| {
1559 let filter = request.filter.as_ref();
1560 let indexed_classifier = record_kind_match_expr(filter);
1561
1562 let count_filter_start = if indexed_classifier.is_some() { 2 } else { 1 };
1563 let (count_filter_clause, count_filter_params) = if let Some(f) = filter {
1564 build_filter_clause(f, &table, count_filter_start)
1565 } else {
1566 (String::new(), Vec::new())
1567 };
1568
1569 let document_count: u64 = {
1570 let count_sql = if indexed_classifier.is_some() {
1571 format!(
1572 "SELECT COUNT(*) FROM {table} \
1573 WHERE {table} MATCH ?1{count_filter_clause}"
1574 )
1575 } else if count_filter_clause.is_empty() {
1576 format!("SELECT COUNT(*) FROM {table}")
1577 } else {
1578 let where_part = count_filter_clause.trim_start_matches(" AND ");
1579 format!("SELECT COUNT(*) FROM {table} WHERE {where_part}")
1580 };
1581 let mut stmt = conn.prepare(&count_sql)?;
1582 if let Some(ref classifier) = indexed_classifier {
1583 stmt.raw_bind_parameter(1, classifier)?;
1584 }
1585 for (i, param) in count_filter_params.iter().enumerate() {
1586 param
1587 .to_sql()
1588 .map(|val| stmt.raw_bind_parameter(count_filter_start + i, val))
1589 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1590 }
1591 let mut rows = stmt.raw_query();
1592 match rows.next()? {
1593 Some(row) => {
1594 let c: i64 = row.get(0)?;
1595 c as u64
1596 }
1597 None => 0,
1598 }
1599 };
1600
1601 let mut results = Vec::with_capacity(request.terms.len());
1602 for term in &request.terms {
1603 let sanitized = sanitize_fts5_token_group(term).unwrap_or_default();
1605 if sanitized.is_empty() {
1606 results.push(TextTermStats {
1607 term: term.clone(),
1608 sanitized_term: sanitized,
1609 document_frequency: 0,
1610 document_count,
1611 inverse_document_frequency: 0.0,
1612 });
1613 continue;
1614 }
1615
1616 let (term_filter_clause, term_filter_params) = if let Some(f) = filter {
1618 build_filter_clause(f, &table, 2)
1619 } else {
1620 (String::new(), Vec::new())
1621 };
1622
1623 let text_term = format!("{{title body}} : ({sanitized})");
1624 let term_match = match &indexed_classifier {
1625 Some(classifier) => format!("{classifier} AND ({text_term})"),
1626 None => text_term,
1627 };
1628 let count_sql = format!(
1629 "SELECT COUNT(*) FROM {table} WHERE {table} MATCH ?1{term_filter_clause}"
1630 );
1631 let mut stmt = conn.prepare(&count_sql)?;
1632 stmt.raw_bind_parameter(1, &term_match)?;
1633 for (i, param) in term_filter_params.iter().enumerate() {
1634 param
1635 .to_sql()
1636 .map(|val| stmt.raw_bind_parameter(2 + i, val))
1637 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1638 }
1639
1640 let df: u64 = {
1641 let mut rows = stmt.raw_query();
1642 match rows.next()? {
1643 Some(row) => {
1644 let c: i64 = row.get(0)?;
1645 c as u64
1646 }
1647 None => 0,
1648 }
1649 };
1650
1651 let idf = Fts5TextSearch::bm25_idf(df, document_count);
1652 results.push(TextTermStats {
1653 term: term.clone(),
1654 sanitized_term: sanitized,
1655 document_frequency: df,
1656 document_count,
1657 inverse_document_frequency: idf,
1658 });
1659 }
1660
1661 Ok(results)
1662 })
1663 .await
1664 }
1665
1666 async fn rebuild(&self, _scope: IndexRebuildScope) -> Result<TextIndexStats, StorageError> {
1667 let table = self.table_name.clone();
1668
1669 self.with_writer("fts_rebuild", move |conn| {
1670 let sql = format!("INSERT INTO {}({}) VALUES('rebuild')", table, table);
1672 conn.execute(&sql, [])?;
1673
1674 let count_sql = format!("SELECT COUNT(*) FROM {}", table);
1675 let count: i64 = conn.query_row(&count_sql, [], |row| row.get(0))?;
1676
1677 Ok(TextIndexStats {
1678 document_count: count as u64,
1679 needs_rebuild: false,
1680 last_rebuild_at: Some(Utc::now()),
1681 })
1682 })
1683 .await
1684 }
1685}
1686
1687impl Fts5TextSearch {
1688 fn bm25_idf(df: u64, document_count: u64) -> f64 {
1690 let n = document_count as f64;
1691 let f = df as f64;
1692 ((n - f + 0.5) / (f + 0.5) + 1.0).ln()
1693 }
1694
1695 async fn search_unranked(
1697 &self,
1698 request: TextSearchRequest,
1699 ) -> Result<Vec<TextSearchHit>, StorageError> {
1700 let table = self.table_name.clone();
1701 let usage = khive_storage::usage::current();
1702
1703 self.with_reader("fts_search_unranked", move |conn| {
1704 let match_expr = match build_filtered_match_expr(
1705 &request.query,
1706 request.mode,
1707 request.filter.as_ref(),
1708 ) {
1709 Some(expr) => expr,
1710 None => return Ok(Vec::new()),
1711 };
1712
1713 let (filter_clause, filter_params) = if let Some(ref filter) = request.filter {
1714 build_filter_clause(filter, &table, 3)
1715 } else {
1716 (String::new(), Vec::new())
1717 };
1718
1719 let sql = format!(
1721 "SELECT subject_id, title \
1722 FROM {table} WHERE {table} MATCH ?1{filter_clause} \
1723 LIMIT ?2",
1724 );
1725
1726 let mut stmt = conn.prepare(&sql)?;
1727 count_fts_pass(usage.as_ref());
1728 stmt.raw_bind_parameter(1, &match_expr)?;
1729 stmt.raw_bind_parameter(2, request.top_k as i64)?;
1730
1731 for (i, param) in filter_params.iter().enumerate() {
1732 param
1733 .to_sql()
1734 .map(|val| stmt.raw_bind_parameter(3 + i, val))
1735 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1736 }
1737
1738 let mut results = Vec::new();
1739 let mut rows = stmt.raw_query();
1740 let mut rank_idx = 0u32;
1741
1742 while let Some(row) = rows.next()? {
1743 let id_str: String = row.get(0)?;
1744 let title: String = row.get(1)?;
1745
1746 let subject_id = Uuid::parse_str(&id_str).map_err(|e| {
1747 rusqlite::Error::FromSqlConversionFailure(
1748 0,
1749 rusqlite::types::Type::Text,
1750 Box::new(e),
1751 )
1752 })?;
1753
1754 rank_idx += 1;
1755 results.push(TextSearchHit {
1756 subject_id,
1757 score: DeterministicScore::from_f64(1.0),
1758 rank: rank_idx,
1759 title: if title.is_empty() { None } else { Some(title) },
1760 snippet: None,
1761 });
1762 }
1763
1764 Ok(results)
1765 })
1766 .await
1767 }
1768
1769 async fn search_rank_within_cap(
1771 &self,
1772 request: TextSearchRequest,
1773 gather_limit: u32,
1774 ) -> Result<Vec<TextSearchHit>, StorageError> {
1775 let table = self.table_name.clone();
1776 let usage = khive_storage::usage::current();
1777
1778 self.with_reader("fts_search_rank_within_cap", move |conn| {
1779 let match_expr = match build_filtered_match_expr(
1780 &request.query,
1781 request.mode,
1782 request.filter.as_ref(),
1783 ) {
1784 Some(expr) => expr,
1785 None => return Ok(Vec::new()),
1786 };
1787
1788 let (filter_clause, filter_params) = if let Some(ref filter) = request.filter {
1789 build_filter_clause(filter, &table, 3)
1790 } else {
1791 (String::new(), Vec::new())
1792 };
1793
1794 let gather_sql = format!(
1796 "SELECT subject_id FROM {table} WHERE {table} MATCH ?1{filter_clause} LIMIT ?2"
1797 );
1798
1799 let mut stmt = conn.prepare(&gather_sql)?;
1800 count_fts_pass(usage.as_ref());
1801 stmt.raw_bind_parameter(1, &match_expr)?;
1802 stmt.raw_bind_parameter(2, gather_limit as i64)?;
1803 for (i, param) in filter_params.iter().enumerate() {
1804 param
1805 .to_sql()
1806 .map(|val| stmt.raw_bind_parameter(3 + i, val))
1807 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1808 }
1809
1810 let mut gathered_ids: Vec<String> = Vec::new();
1811 let mut rows = stmt.raw_query();
1812 while let Some(row) = rows.next()? {
1813 gathered_ids.push(row.get::<_, String>(0)?);
1814 }
1815
1816 if gathered_ids.is_empty() {
1817 return Ok(Vec::new());
1818 }
1819
1820 let snippet_expr = if request.snippet_chars == 0 {
1822 "NULL AS snippet".to_string()
1823 } else {
1824 let chars = i32::try_from(request.snippet_chars).unwrap_or(i32::MAX);
1825 format!("snippet({table}, 3, '', '', '...', {chars})")
1826 };
1827
1828 let id_placeholders: Vec<String> = gathered_ids
1830 .iter()
1831 .enumerate()
1832 .map(|(i, _)| format!("?{}", 3 + i))
1833 .collect();
1834 let in_clause = id_placeholders.join(", ");
1835
1836 let rank_sql = format!(
1837 "SELECT subject_id, rank, title, {snippet_expr} \
1838 FROM {table} WHERE {table} MATCH ?1 \
1839 AND rank MATCH '{LEXICAL_BM25_RANK}' \
1840 AND subject_id IN ({in_clause}) \
1841 ORDER BY rank LIMIT ?2",
1842 );
1843
1844 let mut stmt2 = conn.prepare(&rank_sql)?;
1845 count_fts_pass(usage.as_ref());
1846 stmt2.raw_bind_parameter(1, &match_expr)?;
1847 stmt2.raw_bind_parameter(2, request.top_k as i64)?;
1848 for (i, id_str) in gathered_ids.iter().enumerate() {
1849 stmt2.raw_bind_parameter(3 + i, id_str.as_str())?;
1850 }
1851
1852 let mut hits = Vec::new();
1853 let mut rows2 = stmt2.raw_query();
1854 let mut rank_idx = 0u32;
1855
1856 while let Some(row) = rows2.next()? {
1857 let id_str: String = row.get(0)?;
1858 let fts_rank: f64 = row.get(1)?;
1859 let title: String = row.get(2)?;
1860 let snippet: Option<String> = row.get(3)?;
1861
1862 let subject_id = Uuid::parse_str(&id_str).map_err(|e| {
1863 rusqlite::Error::FromSqlConversionFailure(
1864 0,
1865 rusqlite::types::Type::Text,
1866 Box::new(e),
1867 )
1868 })?;
1869
1870 rank_idx += 1;
1871 hits.push((subject_id, fts_rank, rank_idx, title, snippet));
1872 }
1873
1874 let min_rank = hits.iter().map(|h| h.1).fold(f64::INFINITY, f64::min);
1876 let max_rank = hits.iter().map(|h| h.1).fold(f64::NEG_INFINITY, f64::max);
1877 let range = max_rank - min_rank;
1878
1879 let results = hits
1880 .into_iter()
1881 .map(|(subject_id, raw_rank, rank, title, snippet)| {
1882 let score = if range.abs() < 1e-12 {
1883 1.0
1884 } else {
1885 let t = (max_rank - raw_rank) / range;
1886 0.05 + 0.95 * t
1887 };
1888 TextSearchHit {
1889 subject_id,
1890 score: DeterministicScore::from_f64(score),
1891 rank,
1892 title: if title.is_empty() { None } else { Some(title) },
1893 snippet: snippet.filter(|s| !s.is_empty()),
1894 }
1895 })
1896 .collect();
1897
1898 Ok(results)
1899 })
1900 .await
1901 }
1902
1903 #[allow(dead_code)]
1914 pub(crate) async fn rename_namespace(
1915 &self,
1916 old_namespace: &str,
1917 new_namespace: &str,
1918 ) -> Result<u64, StorageError> {
1919 if old_namespace == new_namespace {
1920 return Ok(0);
1921 }
1922 let table = self.table_name.clone();
1923 let old_ns = old_namespace.to_string();
1924 let new_ns = new_namespace.to_string();
1925
1926 if let Some(writer_task) = self.current_writer_task("fts_rename_namespace")? {
1934 let table2 = table.clone();
1935 return writer_task
1936 .send_bounded(move |conn| {
1937 rename_namespace_dml(conn, &table2, &old_ns, &new_ns)
1938 .map_err(|e| map_err(e, "fts_rename_namespace"))
1939 })
1940 .await;
1941 }
1942
1943 self.pool
1944 .record_direct_route(crate::timeout_sink::Site::DirectRouteFtsRenameNamespace);
1945
1946 let origin = self.pool.origin();
1947 self.with_writer_unmanaged("fts_rename_namespace", move |conn| {
1948 conn.execute_batch("BEGIN IMMEDIATE")?;
1949 let _tx_handle = khive_storage::tx_registry::register_scoped(
1950 Some("text_rename_namespace".to_string()),
1951 origin,
1952 );
1953 match rename_namespace_dml(conn, &table, &old_ns, &new_ns) {
1957 Ok(moved) => {
1958 conn.execute_batch("COMMIT")?;
1959 Ok(moved)
1960 }
1961 Err(e) => {
1962 let _ = conn.execute_batch("ROLLBACK");
1963 Err(e)
1964 }
1965 }
1966 })
1967 .await
1968 }
1969}
1970
1971struct FtsRenameRow {
1972 subject_id: String,
1973 kind: String,
1974 title: String,
1975 body: String,
1976 tags: String,
1977 metadata: Option<String>,
1978 updated_at: i64,
1979 record_kind: Option<String>,
1980}
1981
1982fn rename_namespace_dml(
1990 conn: &rusqlite::Connection,
1991 table: &str,
1992 old_ns: &str,
1993 new_ns: &str,
1994) -> Result<u64, rusqlite::Error> {
1995 let sel_sql = format!(
1996 "SELECT subject_id, kind, title, body, tags, metadata, updated_at, record_kind \
1997 FROM {} WHERE namespace = ?1",
1998 table
1999 );
2000 let rows: Vec<FtsRenameRow> = {
2001 let mut stmt = conn.prepare(&sel_sql)?;
2002 let iter = stmt.query_map(rusqlite::params![old_ns], |row| {
2003 Ok(FtsRenameRow {
2004 subject_id: row.get(0)?,
2005 kind: row.get(1)?,
2006 title: row.get(2)?,
2007 body: row.get(3)?,
2008 tags: row.get(4)?,
2009 metadata: row.get(5)?,
2010 updated_at: row.get(6)?,
2011 record_kind: row.get(7)?,
2012 })
2013 })?;
2014 iter.collect::<Result<Vec<_>, _>>()?
2015 };
2016 let moved = rows.len() as u64;
2017 if moved == 0 {
2018 return Ok(0);
2019 }
2020
2021 let del_sql = format!("DELETE FROM {} WHERE namespace = ?1", table);
2022 conn.execute(&del_sql, rusqlite::params![old_ns])?;
2023
2024 let map = rowid_map_table(table);
2028 let map_del_sql = format!("DELETE FROM {map} WHERE namespace = ?1");
2029 conn.execute(&map_del_sql, rusqlite::params![old_ns])?;
2030
2031 let ins_sql = format!(
2032 "INSERT INTO {} \
2033 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, record_kind) \
2034 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
2035 table
2036 );
2037 let map_ins_sql = format!(
2038 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
2039 VALUES (?1, ?2, last_insert_rowid())"
2040 );
2041 for row in &rows {
2042 conn.execute(
2043 &ins_sql,
2044 rusqlite::params![
2045 row.subject_id,
2046 row.kind,
2047 row.title,
2048 row.body,
2049 row.tags,
2050 new_ns,
2051 row.metadata,
2052 row.updated_at,
2053 row.record_kind,
2054 ],
2055 )?;
2056 conn.execute(&map_ins_sql, rusqlite::params![new_ns, row.subject_id])?;
2057 }
2058
2059 Ok(moved)
2060}
2061
2062#[cfg(test)]
2063#[path = "text_tests.rs"]
2064mod tests;