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 e.into_storage_error(StorageCapability::Text, op)
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 current_writer_task(
371 &self,
372 operation: &'static str,
373 ) -> Result<Option<WriterTaskHandle>, StorageError> {
374 self.pool
375 .writer_task_for_write(self.writer_task.as_ref(), operation)
376 }
377
378 async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
386 where
387 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
388 R: Send + 'static,
389 {
390 if let Some(writer_task) = self.current_writer_task(op)? {
391 return writer_task
392 .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
393 .await;
394 }
395
396 self.pool
397 .record_direct_route(crate::timeout_sink::Site::DirectRouteFtsGeneralWrite);
398 self.with_writer_unmanaged(op, f).await
399 }
400
401 async fn with_writer_unmanaged<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
405 where
406 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
407 R: Send + 'static,
408 {
409 let pool = Arc::clone(&self.pool);
410 let result = tokio::task::spawn_blocking(move || {
411 pool.execute_direct_transaction(StorageCapability::Text, op, move |conn| {
412 f(conn).map_err(|error| map_err(error, op))
413 })
414 })
415 .await
416 .map_err(|e| StorageError::driver(StorageCapability::Text, op, e))?;
417 if let Err(err) = &result {
418 let msg = err.to_string();
419 if msg.contains("locked") || msg.contains("busy") {
420 if self.is_file_backed {
421 crate::timeout_sink::emit_timeout(
426 &crate::timeout_sink::db_label(&self.pool),
427 crate::timeout_sink::Site::StandaloneText,
428 &msg,
429 None,
430 );
431 }
432 let open: Vec<String> = khive_storage::tx_registry::snapshot()
433 .into_iter()
434 .map(|(age, label)| {
435 format!(
436 "{}@{}ms",
437 label.as_deref().unwrap_or("unlabeled"),
438 age.as_millis()
439 )
440 })
441 .collect();
442 tracing::warn!(
443 op,
444 open_tx_count = open.len(),
445 open_txs = %open.join(","),
446 "text write starved on the SQLite write lock; open registered \
447 transactions listed (empty means the holder issued no \
448 registered BEGIN IMMEDIATE in this process)"
449 );
450 }
451 }
452 result
453 }
454
455 async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
456 where
457 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
458 R: Send + 'static,
459 {
460 super::run_pooled_store_read(
461 Arc::clone(&self.pool),
462 StorageCapability::Text,
463 op,
464 move |conn| f(conn).map_err(|error| map_err(error, op)),
465 )
466 .await
467 }
468}
469
470fn tags_to_json(tags: &[String]) -> String {
473 serde_json::to_string(tags).unwrap_or_else(|_| "[]".to_string())
474}
475
476fn tags_from_json(s: &str) -> Vec<String> {
477 serde_json::from_str(s).unwrap_or_default()
478}
479
480fn dt_to_micros(dt: &DateTime<Utc>) -> i64 {
481 dt.timestamp_micros()
482}
483
484fn micros_to_dt(micros: i64) -> DateTime<Utc> {
485 Utc.timestamp_micros(micros)
486 .single()
487 .unwrap_or_else(Utc::now)
488}
489
490fn sanitize_fts5_query(query: &str) -> String {
497 let spaced: String = query
501 .chars()
502 .map(|c| {
503 if matches!(c, '(' | ')' | ',' | ':' | '-' | '.' | '/') {
504 ' '
505 } else {
506 c
507 }
508 })
509 .collect();
510
511 let sanitized: String = spaced
518 .chars()
519 .filter(|c| {
520 !matches!(c, '*' | '"' | '\'' | '+' | '^' | '~' | '!' | '$' | '\0') && !c.is_control()
521 })
522 .collect();
523
524 sanitized
526 .split_whitespace()
527 .filter(|t| {
528 !matches!(
529 t.to_ascii_uppercase().as_str(),
530 "AND" | "OR" | "NOT" | "NEAR"
531 )
532 })
533 .collect::<Vec<_>>()
534 .join(" ")
535}
536
537fn sanitize_fts5_query_legacy_merged(query: &str) -> String {
544 let spaced: String = query
545 .chars()
546 .map(|c| {
547 if matches!(c, '(' | ')' | ',' | ':' | '/') {
548 ' '
549 } else {
550 c
551 }
552 })
553 .collect();
554
555 let sanitized: String = spaced
556 .chars()
557 .filter(|c| {
558 !matches!(
559 c,
560 '*' | '"' | '\'' | '+' | '-' | '^' | '.' | '~' | '!' | '$' | '\0'
561 ) && !c.is_control()
562 })
563 .collect();
564
565 sanitized
566 .split_whitespace()
567 .filter(|t| {
568 !matches!(
569 t.to_ascii_uppercase().as_str(),
570 "AND" | "OR" | "NOT" | "NEAR"
571 )
572 })
573 .collect::<Vec<_>>()
574 .join(" ")
575}
576
577fn sanitize_fts5_phrase_literal(token: &str) -> Option<String> {
586 let mut literal = String::with_capacity(token.len());
587 for c in token
588 .chars()
589 .filter(|c| !matches!(c, '*' | '\0') && !c.is_control())
590 {
591 if c == '"' {
592 literal.push_str("\"\"");
593 } else {
594 literal.push(c);
595 }
596 }
597 if literal.is_empty() {
598 None
599 } else {
600 Some(literal)
601 }
602}
603
604const FTS5_TRIGRAM_MIN_SAFE_LEN: usize = 3;
612
613fn is_fts5_bareword_safe(s: &str) -> bool {
617 !s.is_empty() && s.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
618}
619
620fn sanitize_fts5_token_group(token: &str) -> Option<String> {
628 let split = sanitize_fts5_query(token);
629 let split_terms: Vec<&str> = split.split_whitespace().collect();
630 if split_terms.is_empty() {
631 return None;
632 }
633
634 let all_bareword_safe = split_terms.iter().all(|t| is_fts5_bareword_safe(t));
635 if split_terms.len() == 1 && is_fts5_bareword_safe(token) {
636 return Some(split_terms[0].to_string());
637 }
638
639 let has_trigram_unsafe_segment = split_terms
640 .iter()
641 .any(|t| t.chars().count() < FTS5_TRIGRAM_MIN_SAFE_LEN);
642
643 let mut alternatives = Vec::new();
644 if all_bareword_safe && !has_trigram_unsafe_segment {
645 alternatives.push(format!("({})", split_terms.join(" ")));
646 }
647
648 let merged = sanitize_fts5_query_legacy_merged(token);
659 let merged_terms: Vec<&str> = merged.split_whitespace().collect();
660 let merged_all_bareword_safe =
661 !merged_terms.is_empty() && merged_terms.iter().all(|t| is_fts5_bareword_safe(t));
662 let merged_has_unsafe_segment = merged_terms.len() > 1
663 && merged_terms
664 .iter()
665 .any(|t| t.chars().count() < FTS5_TRIGRAM_MIN_SAFE_LEN);
666 let phrase = sanitize_fts5_phrase_literal(token);
675 let duplicates_split = merged == split_terms.join(" ");
676 let duplicates_phrase = phrase.as_deref() == Some(merged.as_str());
677 if merged_all_bareword_safe
678 && !merged.is_empty()
679 && !duplicates_split
680 && !duplicates_phrase
681 && !merged_has_unsafe_segment
682 {
683 alternatives.push(merged);
684 }
685
686 if let Some(phrase) = phrase {
687 alternatives.push(format!("\"{}\"", phrase));
688 }
689
690 match alternatives.len() {
691 0 => None,
692 1 => alternatives.into_iter().next(),
693 _ => Some(format!("({})", alternatives.join(" OR "))),
694 }
695}
696
697fn join_plain_groups(groups: &[String]) -> String {
712 let mut expr = String::new();
713 for (i, group) in groups.iter().enumerate() {
714 if i > 0 {
715 let prev_compound = groups[i - 1].starts_with('(');
716 let this_compound = group.starts_with('(');
717 expr.push_str(if prev_compound || this_compound {
718 " AND "
719 } else {
720 " "
721 });
722 }
723 expr.push_str(group);
724 }
725 expr
726}
727
728fn build_match_expr(query: &str, mode: TextQueryMode) -> Option<String> {
745 match mode {
746 TextQueryMode::AnyTerm => {
747 let groups: Vec<String> = query
748 .split_whitespace()
749 .filter_map(sanitize_fts5_token_group)
750 .collect();
751 if groups.is_empty() {
752 None
753 } else {
754 Some(groups.join(" OR "))
755 }
756 }
757 TextQueryMode::Plain => {
758 let groups: Vec<String> = query
759 .split_whitespace()
760 .filter_map(sanitize_fts5_token_group)
761 .collect();
762 if groups.is_empty() {
763 None
764 } else {
765 Some(join_plain_groups(&groups))
766 }
767 }
768 TextQueryMode::Phrase => {
769 sanitize_fts5_phrase_literal(query).map(|literal| format!("\"{}\"", literal))
770 }
771 }
772}
773
774fn quote_fts5_phrase(value: &str) -> String {
775 format!("\"{}\"", value.replace('"', "\"\""))
776}
777
778fn record_kind_match_expr(filter: Option<&TextFilter>) -> Option<String> {
785 let record_kinds = &filter?.record_kinds;
786 if record_kinds.is_empty() || record_kinds.iter().any(|kind| kind.chars().count() < 3) {
787 return None;
788 }
789
790 let clauses: Vec<String> = record_kinds
791 .iter()
792 .map(|kind| format!("record_kind : {}", quote_fts5_phrase(kind)))
793 .collect();
794 if clauses.len() == 1 {
795 clauses.into_iter().next()
796 } else {
797 Some(format!("({})", clauses.join(" OR ")))
798 }
799}
800
801fn build_filtered_match_expr(
802 query: &str,
803 mode: TextQueryMode,
804 filter: Option<&TextFilter>,
805) -> Option<String> {
806 let query = format!("{{title body}} : ({})", build_match_expr(query, mode)?);
810 match record_kind_match_expr(filter) {
811 Some(classifier) => Some(format!("{classifier} AND ({query})")),
812 None => Some(query),
813 }
814}
815
816const LEXICAL_BM25_RANK: &str = "bm25(0.0, 0.0, 1.0, 1.0, 0.0, 0.0, 0.0, 0.0, 0.0)";
823
824fn build_filter_clause(
829 filter: &TextFilter,
830 table: &str,
831 start_idx: usize,
832) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
833 let mut conditions: Vec<String> = Vec::new();
834 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
835 let mut idx = start_idx;
836
837 if !filter.ids.is_empty() {
838 let placeholders: Vec<String> = filter
839 .ids
840 .iter()
841 .map(|_| {
842 let p = format!("?{}", idx);
843 idx += 1;
844 p
845 })
846 .collect();
847 conditions.push(format!(
848 "{}.subject_id IN ({})",
849 table,
850 placeholders.join(", ")
851 ));
852 for id in &filter.ids {
853 params.push(Box::new(id.to_string()));
854 }
855 }
856
857 if !filter.kinds.is_empty() {
858 let placeholders: Vec<String> = filter
859 .kinds
860 .iter()
861 .map(|_| {
862 let p = format!("?{}", idx);
863 idx += 1;
864 p
865 })
866 .collect();
867 conditions.push(format!("{}.kind IN ({})", table, placeholders.join(", ")));
868 for kind in &filter.kinds {
869 params.push(Box::new(kind.to_string()));
870 }
871 }
872
873 if !filter.record_kinds.is_empty() {
874 let placeholders: Vec<String> = filter
875 .record_kinds
876 .iter()
877 .map(|_| {
878 let p = format!("?{}", idx);
879 idx += 1;
880 p
881 })
882 .collect();
883 conditions.push(format!(
884 "{}.record_kind IN ({})",
885 table,
886 placeholders.join(", ")
887 ));
888 for kind in &filter.record_kinds {
889 params.push(Box::new(kind.clone()));
890 }
891 }
892
893 if !filter.namespaces.is_empty() {
894 let placeholders: Vec<String> = filter
895 .namespaces
896 .iter()
897 .map(|_| {
898 let p = format!("?{}", idx);
899 idx += 1;
900 p
901 })
902 .collect();
903 conditions.push(format!(
904 "{}.namespace IN ({})",
905 table,
906 placeholders.join(", ")
907 ));
908 for ns in &filter.namespaces {
909 params.push(Box::new(ns.clone()));
910 }
911 }
912
913 if conditions.is_empty() {
914 (String::new(), params)
915 } else {
916 (format!(" AND {}", conditions.join(" AND ")), params)
917 }
918}
919
920#[cfg(test)]
944pub(crate) mod delete_dml_test_seam {
945 use std::collections::HashSet;
946 use std::sync::Mutex;
947 use uuid::Uuid;
948
949 pub(crate) static TARGETS: Mutex<Option<HashSet<(String, Uuid)>>> = Mutex::new(None);
950
951 pub(crate) fn arm(namespace: &str, subject_id: Uuid) {
952 TARGETS
953 .lock()
954 .unwrap()
955 .get_or_insert_with(HashSet::new)
956 .insert((namespace.to_string(), subject_id));
957 }
958
959 pub(crate) fn disarm(namespace: &str, subject_id: Uuid) {
960 if let Some(targets) = TARGETS.lock().unwrap().as_mut() {
961 targets.remove(&(namespace.to_string(), subject_id));
962 }
963 }
964
965 pub(crate) fn matches(namespace: &str, subject_id: Uuid) -> bool {
966 TARGETS
967 .lock()
968 .unwrap()
969 .as_ref()
970 .is_some_and(|targets| targets.contains(&(namespace.to_string(), subject_id)))
971 }
972}
973
974fn delete_document_dml(
983 conn: &rusqlite::Connection,
984 table: &str,
985 namespace: &str,
986 subject_id: Uuid,
987) -> Result<bool, SqliteError> {
988 let [fts_statement, map_statement] = delete_document_statements(table, namespace, subject_id);
989
990 let mut stmt = conn.prepare(&fts_statement.sql)?;
991 bind_params(&mut stmt, &fts_statement.params)?;
992 let deleted = stmt.raw_execute()? > 0;
993 drop(stmt);
994
995 #[cfg(test)]
996 if delete_dml_test_seam::matches(namespace, subject_id) {
997 return Err(SqliteError::InvalidData(
998 "delete_document_dml test seam: forced failure between the FTS-row \
999 delete and the map-row delete"
1000 .to_string(),
1001 ));
1002 }
1003
1004 let mut map_stmt = conn.prepare(&map_statement.sql)?;
1005 bind_params(&mut map_stmt, &map_statement.params)?;
1006 map_stmt.raw_execute()?;
1007
1008 Ok(deleted)
1009}
1010
1011fn upsert_document_dml(
1017 conn: &rusqlite::Connection,
1018 table: &str,
1019 document: &TextDocument,
1020) -> Result<(), rusqlite::Error> {
1021 let statement = delete_document_statement(table, &document.namespace, document.subject_id);
1025 let mut stmt = conn.prepare(&statement.sql)?;
1026 bind_params(&mut stmt, &statement.params)?;
1027 stmt.raw_execute()?;
1028
1029 for statement in insert_document_statements(table, document) {
1030 let mut stmt = conn.prepare(&statement.sql)?;
1031 bind_params(&mut stmt, &statement.params)?;
1032 stmt.raw_execute()?;
1033 }
1034 Ok(())
1035}
1036
1037fn batch_upsert_documents_dml(
1046 conn: &rusqlite::Connection,
1047 table: &str,
1048 documents: &[TextDocument],
1049 attempted: u64,
1050) -> Result<BatchWriteSummary, rusqlite::Error> {
1051 let map = rowid_map_table(table);
1052 let del_sql = format!(
1059 "DELETE FROM {table} WHERE rowid IN \
1060 (SELECT rowid FROM {map} WHERE namespace = ?1 AND subject_id = ?2) \
1061 AND namespace = ?1 AND subject_id = ?2"
1062 );
1063 let ins_sql = format!(
1064 "INSERT INTO {} \
1065 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, record_kind) \
1066 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
1067 table
1068 );
1069 let map_ins_sql = format!(
1072 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
1073 VALUES (?1, ?2, last_insert_rowid())"
1074 );
1075
1076 let mut summary = BatchWriteSummary {
1077 attempted,
1078 ..BatchWriteSummary::default()
1079 };
1080
1081 for (index, doc) in documents.iter().enumerate() {
1082 conn.execute_batch("SAVEPOINT fts_upsert_doc")?;
1083 let id_str = doc.subject_id.to_string();
1084 let namespace = &doc.namespace;
1085 let result = (|| {
1086 conn.execute(&del_sql, rusqlite::params![namespace, &id_str])?;
1087
1088 let tags_json = tags_to_json(&doc.tags);
1089 let metadata_json: Option<String> = doc.metadata.as_ref().map(|v| v.to_string());
1090
1091 conn.execute(
1092 &ins_sql,
1093 rusqlite::params![
1094 &id_str,
1095 &doc.kind.to_string(),
1096 doc.title.as_deref().unwrap_or(""),
1097 &doc.body,
1098 &tags_json,
1099 namespace,
1100 &metadata_json,
1101 dt_to_micros(&doc.updated_at),
1102 &doc.record_kind,
1103 ],
1104 )?;
1105 conn.execute(&map_ins_sql, rusqlite::params![namespace, &id_str])?;
1106 Ok::<(), rusqlite::Error>(())
1107 })();
1108
1109 match result {
1110 Ok(()) => {
1111 conn.execute_batch("RELEASE SAVEPOINT fts_upsert_doc")?;
1112 summary.affected = summary.affected.saturating_add(1);
1113 }
1114 Err(e) => {
1115 let _ = conn.execute_batch("ROLLBACK TO SAVEPOINT fts_upsert_doc");
1116 let _ = conn.execute_batch("RELEASE SAVEPOINT fts_upsert_doc");
1117 let (class, retryability) = super::classify_batch_sqlite_error(&e);
1118 summary.record_failure(index, Some(id_str), class, retryability, e.to_string());
1119 }
1120 }
1121 }
1122
1123 Ok(summary)
1124}
1125
1126#[async_trait]
1127impl TextSearch for Fts5TextSearch {
1128 async fn upsert_document(&self, document: TextDocument) -> Result<(), StorageError> {
1129 let table = self.table_name.clone();
1130
1131 if let Some(writer_task) = self.current_writer_task("fts_upsert")? {
1139 let table2 = table.clone();
1140 return writer_task
1141 .send_bounded(move |conn| {
1142 upsert_document_dml(conn, &table2, &document)
1143 .map_err(|e| map_err(e, "fts_upsert"))
1144 })
1145 .await;
1146 }
1147
1148 self.with_writer("fts_upsert", move |conn| {
1150 upsert_document_dml(conn, &table, &document)
1151 })
1152 .await
1153 }
1154
1155 async fn upsert_documents(
1156 &self,
1157 documents: Vec<TextDocument>,
1158 ) -> Result<BatchWriteSummary, StorageError> {
1159 let table = self.table_name.clone();
1160 let attempted = documents.len() as u64;
1161
1162 if let Some(writer_task) = self.current_writer_task("fts_upsert_batch")? {
1170 let table2 = table.clone();
1171 return writer_task
1172 .send_bounded(move |conn| {
1173 batch_upsert_documents_dml(conn, &table2, &documents, attempted)
1174 .map_err(|e| map_err(e, "fts_upsert_batch"))
1175 })
1176 .await;
1177 }
1178
1179 self.with_writer("fts_upsert_batch", move |conn| {
1181 batch_upsert_documents_dml(conn, &table, &documents, attempted)
1182 })
1183 .await
1184 }
1185
1186 async fn delete_document(
1187 &self,
1188 namespace: &str,
1189 subject_id: Uuid,
1190 ) -> Result<bool, StorageError> {
1191 let table = self.table_name.clone();
1192 let namespace = namespace.to_string();
1193
1194 if self.scan_fallback {
1195 let statement = delete_document_statement_scan_fallback(&table, &namespace, subject_id);
1196 return self
1197 .with_writer("fts_delete_scan_fallback", move |conn| {
1198 let mut stmt = conn.prepare(&statement.sql)?;
1199 bind_params(&mut stmt, &statement.params)?;
1200 Ok(stmt.raw_execute()? > 0)
1201 })
1202 .await;
1203 }
1204
1205 if let Some(writer_task) = self.current_writer_task("fts_delete")? {
1213 let table2 = table.clone();
1214 let namespace2 = namespace.clone();
1215 return writer_task
1216 .send_bounded(move |conn| {
1217 delete_document_dml(conn, &table2, &namespace2, subject_id)
1218 .map_err(|e| map_sqlite_err(e, "fts_delete"))
1219 })
1220 .await;
1221 }
1222
1223 self.with_writer("fts_delete", move |conn| {
1224 delete_document_dml(conn, &table, &namespace, subject_id).map_err(|error| match error {
1225 SqliteError::Rusqlite(inner) => inner,
1226 other => rusqlite::Error::InvalidParameterName(other.to_string()),
1227 })
1228 })
1229 .await
1230 }
1231
1232 async fn get_document(
1233 &self,
1234 namespace: &str,
1235 subject_id: Uuid,
1236 ) -> Result<Option<TextDocument>, StorageError> {
1237 let namespace = namespace.to_string();
1238 let table = self.table_name.clone();
1239 let scan_fallback = self.scan_fallback;
1240
1241 self.with_reader("fts_get", move |conn| {
1242 let sql = if scan_fallback {
1243 format!(
1244 "SELECT subject_id, kind, title, body, tags, namespace, \
1245 metadata, updated_at, record_kind \
1246 FROM {table} WHERE namespace = ?1 AND subject_id = ?2"
1247 )
1248 } else {
1249 let map = rowid_map_table(&table);
1250 format!(
1251 "SELECT t.subject_id, t.kind, t.title, t.body, t.tags, t.namespace, \
1252 t.metadata, t.updated_at, t.record_kind \
1253 FROM {table} AS t JOIN {map} AS m ON m.rowid = t.rowid \
1254 WHERE m.namespace = ?1 AND m.subject_id = ?2 \
1255 AND t.namespace = ?1 AND t.subject_id = ?2"
1256 )
1257 };
1258 let mut stmt = conn.prepare(&sql)?;
1259 let mut rows = stmt.query(rusqlite::params![namespace, subject_id.to_string()])?;
1260
1261 match rows.next()? {
1262 Some(row) => {
1263 let id_str: String = row.get(0)?;
1264 let kind_str: String = row.get(1)?;
1265 let title: String = row.get(2)?;
1266 let body: String = row.get(3)?;
1267 let tags_json: String = row.get(4)?;
1268 let ns: String = row.get(5)?;
1269 let metadata_json: Option<String> = row.get(6)?;
1270 let updated_at_micros: i64 = row.get(7)?;
1271 let record_kind: Option<String> = row.get(8)?;
1272
1273 let sid = Uuid::parse_str(&id_str).map_err(|e| {
1274 rusqlite::Error::FromSqlConversionFailure(
1275 0,
1276 rusqlite::types::Type::Text,
1277 Box::new(e),
1278 )
1279 })?;
1280
1281 let kind = kind_str.parse::<SubstrateKind>().map_err(|e| {
1282 rusqlite::Error::FromSqlConversionFailure(
1283 1,
1284 rusqlite::types::Type::Text,
1285 Box::new(e),
1286 )
1287 })?;
1288
1289 Ok(Some(TextDocument {
1290 subject_id: sid,
1291 kind,
1292 record_kind,
1293 title: if title.is_empty() { None } else { Some(title) },
1294 body,
1295 tags: tags_from_json(&tags_json),
1296 namespace: ns,
1297 metadata: metadata_json.and_then(|s| serde_json::from_str(&s).ok()),
1298 updated_at: micros_to_dt(updated_at_micros),
1299 }))
1300 }
1301 None => Ok(None),
1302 }
1303 })
1304 .await
1305 }
1306
1307 async fn search(&self, request: TextSearchRequest) -> Result<Vec<TextSearchHit>, StorageError> {
1308 let table = self.table_name.clone();
1309 let usage = khive_storage::usage::current();
1310
1311 let result = self
1312 .with_reader("fts_search", move |conn| {
1313 let match_expr = match build_filtered_match_expr(
1314 &request.query,
1315 request.mode,
1316 request.filter.as_ref(),
1317 ) {
1318 Some(expr) => expr,
1319 None => return Ok(Vec::new()),
1320 };
1321
1322 let snippet_expr = if request.snippet_chars == 0 {
1329 "NULL AS snippet".to_string()
1330 } else {
1331 let chars = i32::try_from(request.snippet_chars).unwrap_or(i32::MAX);
1332 format!("snippet({table}, 3, '', '', '...', {chars})")
1333 };
1334
1335 let (filter_clause, filter_params) = if let Some(ref filter) = request.filter {
1336 build_filter_clause(filter, &table, 3)
1337 } else {
1338 (String::new(), Vec::new())
1339 };
1340
1341 let sql = format!(
1342 "SELECT subject_id, rank, title, {snippet_expr} \
1343 FROM {table} WHERE {table} MATCH ?1 \
1344 AND rank MATCH '{LEXICAL_BM25_RANK}'{filter_clause} \
1345 ORDER BY rank LIMIT ?2",
1346 );
1347
1348 let mut stmt = conn.prepare(&sql)?;
1349 count_fts_pass(usage.as_ref());
1350 stmt.raw_bind_parameter(1, &match_expr)?;
1351 stmt.raw_bind_parameter(2, request.top_k as i64)?;
1352
1353 for (i, param) in filter_params.iter().enumerate() {
1354 param
1355 .to_sql()
1356 .map(|val| stmt.raw_bind_parameter(3 + i, val))
1357 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1358 }
1359
1360 let mut hits = Vec::new();
1361 let mut rows = stmt.raw_query();
1362 let mut rank_idx = 0u32;
1363
1364 while let Some(row) = rows.next()? {
1365 let id_str: String = row.get(0)?;
1366 let fts_rank: f64 = row.get(1)?;
1367 let title: String = row.get(2)?;
1368 let snippet: Option<String> = row.get(3)?;
1369
1370 let subject_id = Uuid::parse_str(&id_str).map_err(|e| {
1371 rusqlite::Error::FromSqlConversionFailure(
1372 0,
1373 rusqlite::types::Type::Text,
1374 Box::new(e),
1375 )
1376 })?;
1377
1378 rank_idx += 1;
1379 hits.push((subject_id, fts_rank, rank_idx, title, snippet));
1380 }
1381
1382 let min_rank = hits.iter().map(|h| h.1).fold(f64::INFINITY, f64::min);
1385 let max_rank = hits.iter().map(|h| h.1).fold(f64::NEG_INFINITY, f64::max);
1386 let range = max_rank - min_rank;
1387
1388 let results = hits
1389 .into_iter()
1390 .map(|(subject_id, raw_rank, rank, title, snippet)| {
1391 let score = if range.abs() < 1e-12 {
1392 1.0
1393 } else {
1394 let t = (max_rank - raw_rank) / range;
1395 0.05 + 0.95 * t
1396 };
1397 TextSearchHit {
1398 subject_id,
1399 score: DeterministicScore::from_f64(score),
1400 rank,
1401 title: if title.is_empty() { None } else { Some(title) },
1402 snippet: snippet.filter(|s| !s.is_empty()),
1403 }
1404 })
1405 .collect();
1406
1407 Ok(results)
1408 })
1409 .await;
1410
1411 result
1412 }
1413
1414 async fn count(&self, filter: TextFilter) -> Result<u64, StorageError> {
1415 let table = self.table_name.clone();
1416
1417 self.with_reader("fts_count", move |conn| {
1418 let indexed_classifier = record_kind_match_expr(Some(&filter));
1419 let filter_start = if indexed_classifier.is_some() { 2 } else { 1 };
1420 let (filter_clause, filter_params) = build_filter_clause(&filter, &table, filter_start);
1421
1422 let sql = if indexed_classifier.is_some() {
1423 format!("SELECT COUNT(*) FROM {table} WHERE {table} MATCH ?1{filter_clause}")
1424 } else if filter_clause.is_empty() {
1425 format!("SELECT COUNT(*) FROM {table}")
1426 } else {
1427 let where_part = filter_clause.trim_start_matches(" AND ");
1428 format!("SELECT COUNT(*) FROM {table} WHERE {where_part}")
1429 };
1430
1431 let mut stmt = conn.prepare(&sql)?;
1432
1433 if let Some(classifier) = indexed_classifier {
1434 stmt.raw_bind_parameter(1, classifier)?;
1435 }
1436
1437 for (i, param) in filter_params.iter().enumerate() {
1438 param
1439 .to_sql()
1440 .map(|val| stmt.raw_bind_parameter(filter_start + i, val))
1441 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1442 }
1443
1444 let mut rows = stmt.raw_query();
1445 match rows.next()? {
1446 Some(row) => {
1447 let count: i64 = row.get(0)?;
1448 Ok(count as u64)
1449 }
1450 None => Ok(0),
1451 }
1452 })
1453 .await
1454 }
1455
1456 async fn stats(&self) -> Result<TextIndexStats, StorageError> {
1457 let table = self.table_name.clone();
1458
1459 self.with_reader("fts_stats", move |conn| {
1460 let sql = format!("SELECT COUNT(*) FROM {}", table);
1461 let count: i64 = conn.query_row(&sql, [], |row| row.get(0))?;
1462
1463 Ok(TextIndexStats {
1464 document_count: count as u64,
1465 needs_rebuild: false,
1466 last_rebuild_at: None,
1467 })
1468 })
1469 .await
1470 }
1471
1472 async fn search_with_options(
1473 &self,
1474 request: TextSearchRequest,
1475 options: TextSearchOptions,
1476 ) -> Result<Vec<TextSearchHit>, StorageError> {
1477 match options.gather_mode {
1478 TextGatherMode::Ranked => self.search(request).await,
1479 TextGatherMode::Unranked => self.search_unranked(request).await,
1480 TextGatherMode::RankWithinCap => {
1481 let gather_limit = options
1482 .gather_limit
1483 .unwrap_or(request.top_k)
1484 .max(request.top_k);
1485 self.search_rank_within_cap(request, gather_limit).await
1486 }
1487 }
1488 }
1489
1490 async fn term_stats(
1491 &self,
1492 request: TextTermStatsRequest,
1493 ) -> Result<Vec<TextTermStats>, StorageError> {
1494 let table = self.table_name.clone();
1495
1496 self.with_reader("fts_term_stats", move |conn| {
1497 let filter = request.filter.as_ref();
1498 let indexed_classifier = record_kind_match_expr(filter);
1499
1500 let count_filter_start = if indexed_classifier.is_some() { 2 } else { 1 };
1501 let (count_filter_clause, count_filter_params) = if let Some(f) = filter {
1502 build_filter_clause(f, &table, count_filter_start)
1503 } else {
1504 (String::new(), Vec::new())
1505 };
1506
1507 let document_count: u64 = {
1508 let count_sql = if indexed_classifier.is_some() {
1509 format!(
1510 "SELECT COUNT(*) FROM {table} \
1511 WHERE {table} MATCH ?1{count_filter_clause}"
1512 )
1513 } else if count_filter_clause.is_empty() {
1514 format!("SELECT COUNT(*) FROM {table}")
1515 } else {
1516 let where_part = count_filter_clause.trim_start_matches(" AND ");
1517 format!("SELECT COUNT(*) FROM {table} WHERE {where_part}")
1518 };
1519 let mut stmt = conn.prepare(&count_sql)?;
1520 if let Some(ref classifier) = indexed_classifier {
1521 stmt.raw_bind_parameter(1, classifier)?;
1522 }
1523 for (i, param) in count_filter_params.iter().enumerate() {
1524 param
1525 .to_sql()
1526 .map(|val| stmt.raw_bind_parameter(count_filter_start + i, val))
1527 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1528 }
1529 let mut rows = stmt.raw_query();
1530 match rows.next()? {
1531 Some(row) => {
1532 let c: i64 = row.get(0)?;
1533 c as u64
1534 }
1535 None => 0,
1536 }
1537 };
1538
1539 let mut results = Vec::with_capacity(request.terms.len());
1540 for term in &request.terms {
1541 let sanitized = sanitize_fts5_token_group(term).unwrap_or_default();
1543 if sanitized.is_empty() {
1544 results.push(TextTermStats {
1545 term: term.clone(),
1546 sanitized_term: sanitized,
1547 document_frequency: 0,
1548 document_count,
1549 inverse_document_frequency: 0.0,
1550 });
1551 continue;
1552 }
1553
1554 let (term_filter_clause, term_filter_params) = if let Some(f) = filter {
1556 build_filter_clause(f, &table, 2)
1557 } else {
1558 (String::new(), Vec::new())
1559 };
1560
1561 let text_term = format!("{{title body}} : ({sanitized})");
1562 let term_match = match &indexed_classifier {
1563 Some(classifier) => format!("{classifier} AND ({text_term})"),
1564 None => text_term,
1565 };
1566 let count_sql = format!(
1567 "SELECT COUNT(*) FROM {table} WHERE {table} MATCH ?1{term_filter_clause}"
1568 );
1569 let mut stmt = conn.prepare(&count_sql)?;
1570 stmt.raw_bind_parameter(1, &term_match)?;
1571 for (i, param) in term_filter_params.iter().enumerate() {
1572 param
1573 .to_sql()
1574 .map(|val| stmt.raw_bind_parameter(2 + i, val))
1575 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1576 }
1577
1578 let df: u64 = {
1579 let mut rows = stmt.raw_query();
1580 match rows.next()? {
1581 Some(row) => {
1582 let c: i64 = row.get(0)?;
1583 c as u64
1584 }
1585 None => 0,
1586 }
1587 };
1588
1589 let idf = Fts5TextSearch::bm25_idf(df, document_count);
1590 results.push(TextTermStats {
1591 term: term.clone(),
1592 sanitized_term: sanitized,
1593 document_frequency: df,
1594 document_count,
1595 inverse_document_frequency: idf,
1596 });
1597 }
1598
1599 Ok(results)
1600 })
1601 .await
1602 }
1603
1604 async fn rebuild(&self, _scope: IndexRebuildScope) -> Result<TextIndexStats, StorageError> {
1605 let table = self.table_name.clone();
1606
1607 self.with_writer("fts_rebuild", move |conn| {
1608 let sql = format!("INSERT INTO {}({}) VALUES('rebuild')", table, table);
1610 conn.execute(&sql, [])?;
1611
1612 let count_sql = format!("SELECT COUNT(*) FROM {}", table);
1613 let count: i64 = conn.query_row(&count_sql, [], |row| row.get(0))?;
1614
1615 Ok(TextIndexStats {
1616 document_count: count as u64,
1617 needs_rebuild: false,
1618 last_rebuild_at: Some(Utc::now()),
1619 })
1620 })
1621 .await
1622 }
1623}
1624
1625impl Fts5TextSearch {
1626 fn bm25_idf(df: u64, document_count: u64) -> f64 {
1628 let n = document_count as f64;
1629 let f = df as f64;
1630 ((n - f + 0.5) / (f + 0.5) + 1.0).ln()
1631 }
1632
1633 async fn search_unranked(
1635 &self,
1636 request: TextSearchRequest,
1637 ) -> Result<Vec<TextSearchHit>, StorageError> {
1638 let table = self.table_name.clone();
1639 let usage = khive_storage::usage::current();
1640
1641 self.with_reader("fts_search_unranked", move |conn| {
1642 let match_expr = match build_filtered_match_expr(
1643 &request.query,
1644 request.mode,
1645 request.filter.as_ref(),
1646 ) {
1647 Some(expr) => expr,
1648 None => return Ok(Vec::new()),
1649 };
1650
1651 let (filter_clause, filter_params) = if let Some(ref filter) = request.filter {
1652 build_filter_clause(filter, &table, 3)
1653 } else {
1654 (String::new(), Vec::new())
1655 };
1656
1657 let sql = format!(
1659 "SELECT subject_id, title \
1660 FROM {table} WHERE {table} MATCH ?1{filter_clause} \
1661 LIMIT ?2",
1662 );
1663
1664 let mut stmt = conn.prepare(&sql)?;
1665 count_fts_pass(usage.as_ref());
1666 stmt.raw_bind_parameter(1, &match_expr)?;
1667 stmt.raw_bind_parameter(2, request.top_k as i64)?;
1668
1669 for (i, param) in filter_params.iter().enumerate() {
1670 param
1671 .to_sql()
1672 .map(|val| stmt.raw_bind_parameter(3 + i, val))
1673 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1674 }
1675
1676 let mut results = Vec::new();
1677 let mut rows = stmt.raw_query();
1678 let mut rank_idx = 0u32;
1679
1680 while let Some(row) = rows.next()? {
1681 let id_str: String = row.get(0)?;
1682 let title: String = row.get(1)?;
1683
1684 let subject_id = Uuid::parse_str(&id_str).map_err(|e| {
1685 rusqlite::Error::FromSqlConversionFailure(
1686 0,
1687 rusqlite::types::Type::Text,
1688 Box::new(e),
1689 )
1690 })?;
1691
1692 rank_idx += 1;
1693 results.push(TextSearchHit {
1694 subject_id,
1695 score: DeterministicScore::from_f64(1.0),
1696 rank: rank_idx,
1697 title: if title.is_empty() { None } else { Some(title) },
1698 snippet: None,
1699 });
1700 }
1701
1702 Ok(results)
1703 })
1704 .await
1705 }
1706
1707 async fn search_rank_within_cap(
1709 &self,
1710 request: TextSearchRequest,
1711 gather_limit: u32,
1712 ) -> Result<Vec<TextSearchHit>, StorageError> {
1713 let table = self.table_name.clone();
1714 let usage = khive_storage::usage::current();
1715
1716 self.with_reader("fts_search_rank_within_cap", move |conn| {
1717 let match_expr = match build_filtered_match_expr(
1718 &request.query,
1719 request.mode,
1720 request.filter.as_ref(),
1721 ) {
1722 Some(expr) => expr,
1723 None => return Ok(Vec::new()),
1724 };
1725
1726 let (filter_clause, filter_params) = if let Some(ref filter) = request.filter {
1727 build_filter_clause(filter, &table, 3)
1728 } else {
1729 (String::new(), Vec::new())
1730 };
1731
1732 let gather_sql = format!(
1734 "SELECT subject_id FROM {table} WHERE {table} MATCH ?1{filter_clause} LIMIT ?2"
1735 );
1736
1737 let mut stmt = conn.prepare(&gather_sql)?;
1738 count_fts_pass(usage.as_ref());
1739 stmt.raw_bind_parameter(1, &match_expr)?;
1740 stmt.raw_bind_parameter(2, gather_limit as i64)?;
1741 for (i, param) in filter_params.iter().enumerate() {
1742 param
1743 .to_sql()
1744 .map(|val| stmt.raw_bind_parameter(3 + i, val))
1745 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))??;
1746 }
1747
1748 let mut gathered_ids: Vec<String> = Vec::new();
1749 let mut rows = stmt.raw_query();
1750 while let Some(row) = rows.next()? {
1751 gathered_ids.push(row.get::<_, String>(0)?);
1752 }
1753
1754 if gathered_ids.is_empty() {
1755 return Ok(Vec::new());
1756 }
1757
1758 let snippet_expr = if request.snippet_chars == 0 {
1760 "NULL AS snippet".to_string()
1761 } else {
1762 let chars = i32::try_from(request.snippet_chars).unwrap_or(i32::MAX);
1763 format!("snippet({table}, 3, '', '', '...', {chars})")
1764 };
1765
1766 let id_placeholders: Vec<String> = gathered_ids
1768 .iter()
1769 .enumerate()
1770 .map(|(i, _)| format!("?{}", 3 + i))
1771 .collect();
1772 let in_clause = id_placeholders.join(", ");
1773
1774 let rank_sql = format!(
1775 "SELECT subject_id, rank, title, {snippet_expr} \
1776 FROM {table} WHERE {table} MATCH ?1 \
1777 AND rank MATCH '{LEXICAL_BM25_RANK}' \
1778 AND subject_id IN ({in_clause}) \
1779 ORDER BY rank LIMIT ?2",
1780 );
1781
1782 let mut stmt2 = conn.prepare(&rank_sql)?;
1783 count_fts_pass(usage.as_ref());
1784 stmt2.raw_bind_parameter(1, &match_expr)?;
1785 stmt2.raw_bind_parameter(2, request.top_k as i64)?;
1786 for (i, id_str) in gathered_ids.iter().enumerate() {
1787 stmt2.raw_bind_parameter(3 + i, id_str.as_str())?;
1788 }
1789
1790 let mut hits = Vec::new();
1791 let mut rows2 = stmt2.raw_query();
1792 let mut rank_idx = 0u32;
1793
1794 while let Some(row) = rows2.next()? {
1795 let id_str: String = row.get(0)?;
1796 let fts_rank: f64 = row.get(1)?;
1797 let title: String = row.get(2)?;
1798 let snippet: Option<String> = row.get(3)?;
1799
1800 let subject_id = Uuid::parse_str(&id_str).map_err(|e| {
1801 rusqlite::Error::FromSqlConversionFailure(
1802 0,
1803 rusqlite::types::Type::Text,
1804 Box::new(e),
1805 )
1806 })?;
1807
1808 rank_idx += 1;
1809 hits.push((subject_id, fts_rank, rank_idx, title, snippet));
1810 }
1811
1812 let min_rank = hits.iter().map(|h| h.1).fold(f64::INFINITY, f64::min);
1814 let max_rank = hits.iter().map(|h| h.1).fold(f64::NEG_INFINITY, f64::max);
1815 let range = max_rank - min_rank;
1816
1817 let results = hits
1818 .into_iter()
1819 .map(|(subject_id, raw_rank, rank, title, snippet)| {
1820 let score = if range.abs() < 1e-12 {
1821 1.0
1822 } else {
1823 let t = (max_rank - raw_rank) / range;
1824 0.05 + 0.95 * t
1825 };
1826 TextSearchHit {
1827 subject_id,
1828 score: DeterministicScore::from_f64(score),
1829 rank,
1830 title: if title.is_empty() { None } else { Some(title) },
1831 snippet: snippet.filter(|s| !s.is_empty()),
1832 }
1833 })
1834 .collect();
1835
1836 Ok(results)
1837 })
1838 .await
1839 }
1840
1841 #[allow(dead_code)]
1852 pub(crate) async fn rename_namespace(
1853 &self,
1854 old_namespace: &str,
1855 new_namespace: &str,
1856 ) -> Result<u64, StorageError> {
1857 if old_namespace == new_namespace {
1858 return Ok(0);
1859 }
1860 let table = self.table_name.clone();
1861 let old_ns = old_namespace.to_string();
1862 let new_ns = new_namespace.to_string();
1863
1864 if let Some(writer_task) = self.current_writer_task("fts_rename_namespace")? {
1872 let table2 = table.clone();
1873 return writer_task
1874 .send_bounded(move |conn| {
1875 rename_namespace_dml(conn, &table2, &old_ns, &new_ns)
1876 .map_err(|e| map_err(e, "fts_rename_namespace"))
1877 })
1878 .await;
1879 }
1880
1881 self.pool
1882 .record_direct_route(crate::timeout_sink::Site::DirectRouteFtsRenameNamespace);
1883
1884 self.with_writer_unmanaged("fts_rename_namespace", move |conn| {
1885 rename_namespace_dml(conn, &table, &old_ns, &new_ns)
1887 })
1888 .await
1889 }
1890}
1891
1892struct FtsRenameRow {
1893 subject_id: String,
1894 kind: String,
1895 title: String,
1896 body: String,
1897 tags: String,
1898 metadata: Option<String>,
1899 updated_at: i64,
1900 record_kind: Option<String>,
1901}
1902
1903fn rename_namespace_dml(
1911 conn: &rusqlite::Connection,
1912 table: &str,
1913 old_ns: &str,
1914 new_ns: &str,
1915) -> Result<u64, rusqlite::Error> {
1916 let sel_sql = format!(
1917 "SELECT subject_id, kind, title, body, tags, metadata, updated_at, record_kind \
1918 FROM {} WHERE namespace = ?1",
1919 table
1920 );
1921 let rows: Vec<FtsRenameRow> = {
1922 let mut stmt = conn.prepare(&sel_sql)?;
1923 let iter = stmt.query_map(rusqlite::params![old_ns], |row| {
1924 Ok(FtsRenameRow {
1925 subject_id: row.get(0)?,
1926 kind: row.get(1)?,
1927 title: row.get(2)?,
1928 body: row.get(3)?,
1929 tags: row.get(4)?,
1930 metadata: row.get(5)?,
1931 updated_at: row.get(6)?,
1932 record_kind: row.get(7)?,
1933 })
1934 })?;
1935 iter.collect::<Result<Vec<_>, _>>()?
1936 };
1937 let moved = rows.len() as u64;
1938 if moved == 0 {
1939 return Ok(0);
1940 }
1941
1942 let del_sql = format!("DELETE FROM {} WHERE namespace = ?1", table);
1943 conn.execute(&del_sql, rusqlite::params![old_ns])?;
1944
1945 let map = rowid_map_table(table);
1949 let map_del_sql = format!("DELETE FROM {map} WHERE namespace = ?1");
1950 conn.execute(&map_del_sql, rusqlite::params![old_ns])?;
1951
1952 let ins_sql = format!(
1953 "INSERT INTO {} \
1954 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, record_kind) \
1955 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
1956 table
1957 );
1958 let map_ins_sql = format!(
1959 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
1960 VALUES (?1, ?2, last_insert_rowid())"
1961 );
1962 for row in &rows {
1963 conn.execute(
1964 &ins_sql,
1965 rusqlite::params![
1966 row.subject_id,
1967 row.kind,
1968 row.title,
1969 row.body,
1970 row.tags,
1971 new_ns,
1972 row.metadata,
1973 row.updated_at,
1974 row.record_kind,
1975 ],
1976 )?;
1977 conn.execute(&map_ins_sql, rusqlite::params![new_ns, row.subject_id])?;
1978 }
1979
1980 Ok(moved)
1981}
1982
1983#[cfg(test)]
1984#[path = "text_tests.rs"]
1985mod tests;
1986
1987#[cfg(test)]
1988#[path = "text_busy_tests.rs"]
1989mod direct_busy_tests;