Skip to main content

khive_db/stores/
text.rs

1//! FTS5-backed `TextSearch`: one virtual table per model, scores normalized to `(0.05, 1.0]`.
2
3use 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
26/// Name of the sidecar rowid map for a given FTS table: `{namespace,
27/// subject_id} -> rowid`, `WITHOUT ROWID`, `PRIMARY KEY (namespace,
28/// subject_id)`.
29///
30/// `namespace` and `subject_id` are declared `UNINDEXED` in every FTS5 DDL
31/// (trigram-tokenized production tables cannot index them without polluting
32/// text queries with UUID fragments), so a point lookup on either column is a
33/// full virtual-table scan of `table` — this sidecar turns it into a
34/// primary-key lookup instead. `table` must already be a trusted, sanitized
35/// table name.
36pub fn rowid_map_table(table: &str) -> String {
37    format!("{table}_rowids")
38}
39
40/// Name of [`rowid_map_table`]'s own state sidecar: a `key -> value` table
41/// recording whether the map's one-time legacy backfill has completed, so
42/// completion is a durable, transactionally-written fact instead of being
43/// inferred from the map's row count.
44pub fn rowid_map_state_table(table: &str) -> String {
45    format!("{}_state", rowid_map_table(table))
46}
47
48/// Value of the `backfill` key in [`rowid_map_state_table`] once
49/// `ensure_fts_rowid_map_backfilled` has completed a full reconciliation
50/// pass over `table`.
51pub const ROWID_MAP_BACKFILL_COMPLETE: &str = "complete";
52
53/// DDL for [`rowid_map_table`]'s sidecar and its [`rowid_map_state_table`].
54/// `IF NOT EXISTS` so every creation site (backend ensure-path, migrations,
55/// test helpers) can call it unconditionally. `table` must already be a
56/// trusted, sanitized table name.
57pub 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
74/// The exact `DELETE` this store's `delete_document` issues, for a given
75/// FTS table (ADR-099 B3 r6 structural cut — see `entity.rs`'s sibling
76/// block). Looks the target row's rowid up in [`rowid_map_table`] first (a
77/// primary-key lookup) instead of scanning `table` for `namespace`/
78/// `subject_id`, which are `UNINDEXED` FTS5 columns. `table` must already be
79/// a trusted, sanitized table name (this mirrors `delete_document`'s own
80/// pre-existing lack of a placeholder for table names — `format!` is
81/// required since table identifiers cannot be bound as SQL parameters).
82///
83/// The trailing `AND namespace = ?1 AND subject_id = ?2` re-checks the key
84/// on the (already rowid-narrowed) candidate row itself: if the map ever
85/// held a stale entry — e.g. a crash between an earlier delete's FTS-row
86/// removal and its own map-row removal, before FTS5 reused that rowid for a
87/// different document — this statement would otherwise delete whatever
88/// unrelated document now lives at that rowid instead of affecting zero
89/// rows. Cheap: the `rowid IN (...)` subquery already narrows to at most one
90/// candidate row via the map's primary key, so this is a single-row
91/// post-filter, not a second table scan.
92///
93/// Callers that also need the sidecar's own row removed (a real delete, not
94/// half of an upsert) must additionally run [`delete_document_map_statement`]
95/// — see [`delete_document_statements`] for the paired, order-safe form.
96pub 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
112/// Pre-map fallback: the literal delete this store issued before the rowid
113/// map existed, scanning `table`'s `UNINDEXED` `namespace`/`subject_id`
114/// columns directly. Used only by [`Fts5TextSearch`] instances constructed
115/// in scan-fallback mode (a read-only snapshot whose FTS table predates the
116/// sidecar map — see `StorageBackend::text_with_tokenizer`'s read-only
117/// branch), where the map table does not exist to route through.
118fn 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
133/// Remove `(namespace, subject_id)`'s row from `table`'s [`rowid_map_table`].
134///
135/// Order matters relative to [`delete_document_statement`]: that statement's
136/// subquery still needs this row, so run it first, then this one, in the same
137/// transaction. Use [`delete_document_statements`] to get both in the correct
138/// order.
139pub 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
155/// The full, order-safe deletion of one FTS document: index 0 (the FTS row,
156/// looked up through the map) MUST execute before index 1 (the map row
157/// itself) in the same transaction/connection — index 0's subquery reads the
158/// row index 1 removes. Callers that only need the FTS row gone as half of a
159/// delete-then-insert upsert (the map row gets overwritten by the following
160/// insert's `INSERT OR REPLACE`, so removing it first is redundant) can use
161/// [`delete_document_statement`] alone.
162pub 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
173/// Build the `INSERT` half of the FTS delete-then-insert upsert.
174///
175/// `table` must be a trusted, sanitized table name because SQL identifiers
176/// cannot be bound as parameters.
177///
178/// This statement alone leaves [`rowid_map_table`] out of date — pair it with
179/// [`insert_document_map_statement`] (or use [`insert_document_statements`]
180/// for the order-safe combination) so the sidecar tracks the new rowid.
181pub 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
211/// Upsert `(namespace, subject_id)`'s [`rowid_map_table`] row to point at the
212/// rowid the connection's most recent successful `INSERT` produced.
213///
214/// Must run immediately after the FTS `INSERT` it tracks, on the same
215/// connection, with nothing else written in between: `last_insert_rowid()` is
216/// connection-scoped and reports whichever `INSERT` ran last. `INSERT OR
217/// REPLACE` means this is safe to call whether or not a map row for this key
218/// already exists (the ordinary upsert case does: the old row is overwritten
219/// with the new rowid rather than requiring a prior delete).
220pub 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
239/// The full, order-safe insertion of one FTS document: index 0 (the FTS
240/// `INSERT`) MUST execute immediately before index 1 (the rowid-map upsert)
241/// on the same connection, with no other write between them — see
242/// [`insert_document_map_statement`]'s adjacency contract. Every insert path
243/// (create, and the insert half of an upsert) must go through this pair, or
244/// the sidecar map silently falls out of sync with the FTS table it indexes.
245pub 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/// Ensure the FTS5 virtual table for `table_key` (and its [`rowid_map_table`]
253/// sidecar) exist.
254///
255/// Used in tests to set up an in-memory FTS5 table without the full `StorageBackend`.
256#[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
294/// A TextSearch backed by SQLite FTS5 virtual tables.
295///
296/// Each instance manages one table: `fts_{table_key}`. Documents are stored
297/// with their metadata in UNINDEXED columns. `title`, `body`, and the optional
298/// granular `record_kind` classifier are indexed; corpus filters use the
299/// classifier inside MATCH and retain exact row predicates.
300pub struct Fts5TextSearch {
301    pool: Arc<ConnectionPool>,
302    is_file_backed: bool,
303    table_name: String,
304    writer_task: Option<WriterTaskHandle>,
305    /// Set only for a read-only snapshot whose FTS table predates the
306    /// [`rowid_map_table`] sidecar (see `StorageBackend::text_with_tokenizer`'s
307    /// read-only branch, which cannot create/backfill the map on a read-only
308    /// connection). `get_document`/`delete_document` fall back to the
309    /// pre-map `namespace`/`subject_id` scan predicates in this mode instead
310    /// of joining/subquerying a map table that does not exist.
311    scan_fallback: bool,
312}
313
314impl Fts5TextSearch {
315    /// Create a new FTS5 text search instance.
316    ///
317    /// The FTS5 virtual table (and its [`rowid_map_table`] sidecar) must
318    /// already exist (created by `StorageBackend::text()`).
319    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    /// Create an instance in scan-fallback mode: the sidecar map is assumed
324    /// absent, so every point lookup/delete falls back to the pre-map
325    /// `namespace`/`subject_id` scan predicates. Reserved for a read-only
326    /// snapshot whose FTS table predates the map (see the [`scan_fallback`]
327    /// field doc and `StorageBackend::text_with_tokenizer`'s read-only
328    /// branch).
329    ///
330    /// [`scan_fallback`]: Self::scan_fallback
331    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        // Enabled by default for file-backed pools; explicit off/degraded
347        // construction remains synchronous (ADR-067 Component A, mirrors
348        // entity.rs policy): a missing writer task is cached without failing
349        // construction. Every write re-resolves it and applies
350        // strict/compatibility policy then.
351        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    /// Re-derive writer-task availability at write time instead of trusting
362    /// only the field cached at construction (ADR-136 D1 gate 3 amendment).
363    /// `self.writer_task` permanently caches `None` when this store was
364    /// constructed outside a Tokio runtime (`writer_task_handle()` returns
365    /// `Err(WriterTaskNoRuntime)`, which construction collapses via
366    /// `.ok().flatten()`) — every later write, even ones running inside a
367    /// runtime, would otherwise silently keep bypassing an enabled queue.
368    /// The pool helper also enforces strict fail-closed routing and preserves
369    /// a typed `WriterTaskNoRuntime` error when no runtime is available.
370    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    /// Route a single-row write through the pool-wide `WriterTask` when
379    /// the write queue is enabled and a handle is available. Strict mode
380    /// refuses a missing handle; compatibility mode falls back to the legacy
381    /// standalone-connection / pool-mutex path (ADR-067 Component A, Fork C
382    /// slice 2). See `crates/khive-db/docs/api/text.md`
383    /// for the per-caller routing rules (which methods bypass this via
384    /// `with_writer_unmanaged` and why).
385    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    /// Direct standalone or pooled typed transaction path, bypassing the
402    /// WriterTask channel after routing has selected a compatibility fallback.
403    /// Callbacks contain DML only; this seam owns BEGIN and settlement.
404    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                    // Only the standalone-connection path: the pool-mutex
422                    // (`try_writer`) path's own timeout is already recorded
423                    // at `Site::PoolAdmission` inside `ConnectionPool::writer`
424                    // — emitting here too would double-count the same event.
425                    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
470// -- Helper functions --
471
472fn 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
490/// Sanitize an FTS5 query string to prevent driver errors from special
491/// chars: replace grouping/separator chars with spaces (splitting
492/// punctuated identifiers into terms), strip remaining FTS5 operator
493/// characters, then drop FTS5 keyword tokens (AND, OR, NOT, NEAR). See
494/// `crates/khive-db/docs/api/text.md` for the exact char sets and the
495/// issues (#388 and others) that grew them.
496fn sanitize_fts5_query(query: &str) -> String {
497    // Pass 1: replace grouping/separator chars with spaces to isolate tokens.
498    // Colon, hyphen, and dot are included here (not in Pass 2) so punctuated
499    // identifiers become separate terms rather than merged tokens.
500    let spaced: String = query
501        .chars()
502        .map(|c| {
503            if matches!(c, '(' | ')' | ',' | ':' | '-' | '.' | '/') {
504                ' '
505            } else {
506                c
507            }
508        })
509        .collect();
510
511    // Pass 2: remove remaining FTS5 special chars and control characters.
512    // Single quote (apostrophe) is included because FTS5 Plain-mode queries treat
513    // it as a string-literal delimiter causing "syntax error near '''".
514    // Dollar sign is included (#388) because FTS5's MATCH parser rejects it
515    // unconditionally — `syntax error near "$"` — wherever it appears in the
516    // expression, e.g. "$prev.id" (a common agent query for DSL docs).
517    let sanitized: String = spaced
518        .chars()
519        .filter(|c| {
520            !matches!(c, '*' | '"' | '\'' | '+' | '^' | '~' | '!' | '$' | '\0') && !c.is_control()
521        })
522        .collect();
523
524    // Pass 3: filter FTS5 operator keywords.
525    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
537/// Legacy (pre-#397) sanitization: hyphen and dot are stripped outright
538/// (not space-split), so `khive-pack-memory` normalizes to the single merged
539/// bareword `khivepackmemory`. Used only to build the merged-form
540/// OR-alternative in [`sanitize_fts5_token_group`] — never as the sole
541/// sanitized query — so pre-#397-indexed content stays reachable. See
542/// `crates/khive-db/docs/api/text.md`.
543fn 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
577/// Escape a raw token for use inside a double-quoted FTS5 phrase. Embedded
578/// double quotes are doubled; the trailing-`*` prefix-query trigger and
579/// control characters are removed. Everything else, including punctuation,
580/// passes through literally, since FTS5 phrase text is matched
581/// against the column's own tokenization of that literal text (word-exact
582/// under `unicode61`, substring-exact under `trigram`), not re-sanitized.
583///
584/// Returns `None` if nothing survives the filter.
585fn 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
604/// Below this length, FTS5's built-in `trigram` tokenizer (the production
605/// default — `backend.rs::StorageBackend::text()`) produces zero tokens for
606/// a bareword MATCH term. Verified against a live `tokenize='trigram'`
607/// table: a term this short silently drops out of its AND clause instead of
608/// constraining the match, e.g. `2026 07 10` matches any row containing
609/// `2026` regardless of month/day. `sanitize_fts5_token_group` treats any
610/// split segment at or below this length as trigram-unsafe.
611const FTS5_TRIGRAM_MIN_SAFE_LEN: usize = 3;
612
613/// FTS5 reserves punctuation in barewords for query syntax, including some
614/// forms that change query meaning without producing an error. Keep bareword
615/// alternatives to the conservative ASCII set and quote everything else.
616fn is_fts5_bareword_safe(s: &str) -> bool {
617    !s.is_empty() && s.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
618}
619
620/// Sanitize a single whitespace-isolated raw token into an FTS5
621/// match-expression fragment, emitting additive OR-alternatives (split
622/// AND-group, legacy-merged form, quoted literal phrase) so the result is
623/// never a narrower match than any single safe form alone. Returns `None` if
624/// the token sanitizes to nothing. See
625/// `crates/khive-db/docs/api/text.md` for the full #397 trigram-safety
626/// rationale behind which alternatives get emitted and when.
627fn 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    // An operator-bearing token (e.g. `NEAR(alpha-beta,5)`) can make the
649    // legacy merge itself collapse to multiple space-separated terms rather
650    // than one bareword: pass 1 of `sanitize_fts5_query_legacy_merged` spaces
651    // out `(`, `)`, and `,` while pass 2 removes `-`/`.` outright, so
652    // `NEAR(alpha-beta,5)` merges to `"alphabeta 5"`, not one word. Pushed
653    // unguarded, that multi-term fragment carries the same trigram-unsafe
654    // `5` the split-group check above exists to exclude, and under FTS5's
655    // implicit-AND adjacency it silently drops, broadening the OR-alternative
656    // to any row containing `alphabeta`. Apply the same trigram-safety gate
657    // here whenever the merge is multi-term.
658    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    // Compare against the split AND-group's own space-joined content — the
667    // form it actually contributes to the expression — not the bare
668    // concatenation of its terms. For an ordinary punctuated identifier the
669    // concatenation always matches the merged bareword (both simply drop the
670    // separators), which suppressed this alternative unconditionally and cut
671    // off legacy-indexed content; a real duplicate only exists when the
672    // merge itself produced the same space-separated content the split group
673    // already emits (e.g. an operator token whose merge stays multi-term).
674    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
697/// Join Plain-mode per-token groups into one MATCH expression.
698///
699/// FTS5's implicit-AND adjacency rule (a bare space between two terms)
700/// applies only between two *plain* terms — verified against a live table:
701/// `GQA KV cache` (all bare) matches, but `GQA AND KV AND cache` does not,
702/// because FTS5's explicit `AND` treats a trigram-unsafe short bareword
703/// (`KV`, 2 chars, zero trigrams) as an unsatisfiable operand, while
704/// implicit adjacency instead lets it drop out harmlessly. So: use a bare
705/// space between two plain (unparenthesized) groups to preserve that
706/// existing leniency, and fall back to explicit `AND` only where at least
707/// one side is a parenthesized OR-group (`sanitize_fts5_token_group`'s
708/// output for punctuated tokens) — adjacency there without the operator is
709/// a MATCH-expression syntax error (`(a OR b) c` fails; `(a OR b) AND c`
710/// does not).
711fn 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
728/// Build the FTS5 MATCH expression for a query string under a given mode.
729///
730/// Centralizes the AnyTerm/Plain/Phrase branching previously duplicated
731/// across `search()`, `search_unranked()`, and `search_rank_within_cap()`.
732///
733/// AnyTerm and Plain both process the query token-by-token (split on
734/// whitespace) through [`sanitize_fts5_token_group`], so a punctuated
735/// identifier anywhere in the query gets its split/merged OR-alternative;
736/// AnyTerm joins the per-token groups with `OR`; Plain joins them via
737/// [`join_plain_groups`] (implicit-AND space where safe, explicit `AND`
738/// where a group is parenthesized). Phrase mode keeps the single-string
739/// literal behavior — a double-quoted FTS5 phrase cannot contain
740/// `OR`/parenthesized groups.
741///
742/// Returns `None` when the sanitized query is empty (caller short-circuits
743/// to an empty result set rather than sending an invalid MATCH expression).
744fn 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
778/// Build an indexed classifier predicate for FTS5's MATCH expression.
779///
780/// The production tokenizer is trigram, so values shorter than three
781/// characters have no postings. In that uncommon case the caller keeps only
782/// the exact row predicate from `build_filter_clause`; correctness is
783/// preserved while the indexed optimization deliberately declines.
784fn 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    // Keep lexical terms confined to the two historically indexed text
807    // columns. Once record_kind became indexed, an unqualified query would
808    // otherwise make `memory` match every memory row solely via its classifier.
809    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
816/// Custom FTS5 rank configuration that keeps the classifier out of lexical
817/// relevance while retaining the optimized hidden-`rank` ORDER BY path.
818///
819/// `record_kind` participates in MATCH for candidate pruning, but its weight is
820/// zero so the classifier cannot become a relevance signal. Subject metadata
821/// columns were already UNINDEXED; title/body retain their default unit weight.
822const LEXICAL_BM25_RANK: &str = "bm25(0.0, 0.0, 1.0, 1.0, 0.0, 0.0, 0.0, 0.0, 0.0)";
823
824/// Build a WHERE clause fragment and params for a `TextFilter`.
825///
826/// Returns `(clause, params)` where clause is empty if no filters are active.
827/// Parameter indices start at `?{start_idx}`.
828fn 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/// Test-only seam for [`delete_document_dml`]'s unmanaged (no-writer-task)
921/// path: lets a test force the map-row delete to fail after the FTS-row
922/// delete has already run, so the enclosing `BEGIN IMMEDIATE` / `ROLLBACK`
923/// block's atomicity can be exercised without a real process crash
924/// mid-transaction. Compiles to nothing outside `#[cfg(test)]` — production
925/// code never reads this flag.
926///
927/// A process-wide `Mutex<HashSet<(String, Uuid)>>` rather than a
928/// thread-local: the unmanaged path runs the DML closure inside
929/// `tokio::task::spawn_blocking` (`with_writer_unmanaged`), on a worker
930/// thread distinct from the test's own thread, so a thread-local set by the
931/// test would never be observed by the closure. Each `(namespace,
932/// subject_id)` target is armed/disarmed independently in the shared set —
933/// `arm` inserts its own tuple and `disarm` removes only that same tuple —
934/// so two tests running concurrently, each targeting a different key, cannot
935/// clobber each other's armed state the way a single shared `Option` slot
936/// would (one test's `disarm` would otherwise clear a target a second test
937/// armed after the first but before the first's own disarm ran).
938/// `delete_document_dml` only forces the failure when its own arguments are
939/// a member of the set, so an unrelated delete running concurrently (e.g.
940/// `test_delete_document`, which is not `#[serial]` against this seam) on a
941/// key that was never armed is unaffected regardless of what else is armed
942/// at the time.
943#[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
974/// DML-only single-document delete shared by both the legacy (flag-off) and
975/// WriterTask-routed (flag-on) `delete_document` paths.
976///
977/// Issues no `BEGIN` / `COMMIT` / `ROLLBACK` itself — the caller owns the
978/// enclosing transaction. Both statements (the FTS row, looked up through
979/// the map, then the map row itself) must commit atomically: a crash
980/// between them would leave the map pointing at a rowid FTS5 may later
981/// reuse for a different document.
982fn 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
1011/// DML-only single-document upsert shared by both the legacy (flag-off) and
1012/// WriterTask-routed (flag-on) `upsert_document` paths (ADR-067 Component A).
1013///
1014/// Issues no `BEGIN` / `COMMIT` / `ROLLBACK` itself — the caller owns the
1015/// enclosing transaction.
1016fn upsert_document_dml(
1017    conn: &rusqlite::Connection,
1018    table: &str,
1019    document: &TextDocument,
1020) -> Result<(), rusqlite::Error> {
1021    // Delete (old rowid, via the map) -> insert (new rowid) -> map upsert
1022    // (new rowid). The map's own row for this key needs no separate delete:
1023    // `INSERT OR REPLACE` in the third statement overwrites it in place.
1024    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
1037/// DML-only batch upsert loop shared by both the legacy (flag-off) and
1038/// WriterTask-routed (flag-on) `upsert_documents` paths (ADR-067 Component A).
1039///
1040/// Issues no OUTER `BEGIN` / `COMMIT` / `ROLLBACK` — the caller owns the
1041/// enclosing transaction. The per-row named `SAVEPOINT fts_upsert_doc` is
1042/// preserved unchanged: it is what gives this loop its partial-success
1043/// semantics (one bad document does not abort the whole batch) independent
1044/// of which outer transaction wraps the loop.
1045fn 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    // Delete via the map's primary key instead of a `namespace`/`subject_id`
1053    // scan of the (UNINDEXED on those columns) FTS table itself. The trailing
1054    // `AND namespace = ?1 AND subject_id = ?2` re-checks the key on the
1055    // rowid-narrowed candidate row: a stale map entry pointing at a rowid
1056    // FTS5 has since reused for a different document must delete zero rows,
1057    // not that unrelated document (mirrors `delete_document_statement`).
1058    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    // `INSERT OR REPLACE` — no separate delete of the map's own row needed;
1070    // see `upsert_document_dml`'s matching comment.
1071    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        // ADR-067 Component A: when the write queue is enabled, route
1132        // through the pool-wide WriterTask. DML-only closure — no BEGIN
1133        // IMMEDIATE/COMMIT/ROLLBACK here, since the WriterTask's run loop
1134        // owns the transaction. `current_writer_task("fts_upsert")`
1135        // (ADR-136 D1 gate 3
1136        // amendment) re-checks past a construction-time `None` cache so a
1137        // handle that only became available later is still used here.
1138        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        // The direct typed unit owns the transaction around this DML body.
1149        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        // ADR-067 Component A: when the write queue is enabled, route
1163        // through the pool-wide WriterTask. DML-only closure (the per-row
1164        // `SAVEPOINT fts_upsert_doc` is preserved unchanged — only the OUTER
1165        // BEGIN IMMEDIATE/COMMIT is removed, since the WriterTask's run loop
1166        // owns the enclosing transaction). `current_writer_task` re-checks
1167        // past a construction-time `None` cache — see `upsert_document`'s
1168        // matching note.
1169        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        // The direct typed unit owns the transaction around the complete batch.
1180        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        // Route through the shared WriterTask (DML-only closure — the
1206        // queue's run loop owns the transaction) when available, exactly
1207        // like `upsert_document`'s matching branch. The fallback below wraps
1208        // both statements in one explicit transaction: a crash between the
1209        // FTS-row delete and the map-row delete must never persist only one
1210        // of the two, or a later insert could reuse the freed rowid while
1211        // the map still points at it.
1212        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                // Snippet column index 3 = body in the FTS5 schema.
1323                // snippet_chars == 0 is the sentinel for "no snippet" — skip the
1324                // snippet(...) call entirely and return NULL instead.  This avoids
1325                // the ~12ms BM25 snippet computation on the hot recall path where
1326                // snippets are unused.  Callers that need snippets (diagnostics) pass
1327                // snippet_chars > 0 and get the same behaviour as before.
1328                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                // Normalize scores within the result set to (0.05, 1.0].
1383                // Best rank (most negative) maps to 1.0, worst to 0.05.
1384                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                // Keep document frequency aligned with the per-token search expression.
1542                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                // Per-term count: MATCH is ?1, so filter params start at ?2.
1555                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            // FTS5 rebuild command: repopulates the internal index structures.
1609            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    /// Robertson-Walker BM25 IDF: ln(((N - df + 0.5) / (df + 0.5)) + 1)
1627    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    /// Gather candidates without BM25 ranking; return with uniform score 1.0.
1634    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            // No rank column, no ORDER BY — avoids BM25 computation entirely.
1658            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    /// Two-stage gather: cheap unranked LIMIT gather_limit, then BM25-rank the subset.
1708    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            // Stage 1: cheap unranked gather of rowids.
1733            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            // Stage 2: BM25-rank only the gathered subset via subject_id IN (...).
1759            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            // Build IN clause for the gathered IDs.
1767            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            // Normalize scores within the ranked subset (same formula as search()).
1813            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    /// Move all FTS5 documents from `old_namespace` to `new_namespace` in a
1842    /// single transaction.
1843    ///
1844    /// FTS5 virtual tables do not support updating indexed columns (`title`,
1845    /// `body`) via UPDATE. The correct approach is read-then-delete-then-reinsert.
1846    ///
1847    /// Callers must invoke this after any SQL-level namespace change on the
1848    /// backing entity table so that FTS5 keyword search stays consistent with
1849    /// the entity store.
1850    // REASON: reserved for namespace migration operations
1851    #[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        // ADR-136 D1 gate 2: queue-first, DML-only closure. `run_writer_task`
1865        // already owns the enclosing `BEGIN IMMEDIATE`/`COMMIT`/`ROLLBACK` for
1866        // this request, so the SELECT enumerating rows-to-move runs INSIDE
1867        // that same transaction as the DELETE+INSERT that consumes it —
1868        // closing the TOCTOU window the pre-migration path had (its SELECT
1869        // ran before its own `BEGIN IMMEDIATE`, so a writer landing between
1870        // the two could resurrect a stale row or lose one mid-rename).
1871        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            // The SELECT and writes share the typed unit's transaction.
1886            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
1903/// Move every FTS5 document row from `old_ns` to `new_ns` in `table`.
1904/// DML-only — issues no `BEGIN`/`COMMIT`/`ROLLBACK` of its own, so it is safe
1905/// to call both inside an already-open transaction (the `WriterTask` queue
1906/// path) and wrapped by a caller-managed one (the legacy standalone path).
1907/// The `SELECT` enumerating rows-to-move and the `DELETE`+`INSERT` that
1908/// consumes it always run inside the SAME transaction the caller opened —
1909/// see `Fts5TextSearch::rename_namespace`'s TOCTOU note.
1910fn 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    // The rows just deleted above are gone from `table` entirely (whole
1946    // namespace, not a single subject), so their map entries are pure
1947    // leftovers now — clear them before reinserting under `new_ns`.
1948    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;