Skip to main content

khive_db/stores/
entity.rs

1//! SQL-backed `EntityStore` implementation.
2
3use std::collections::HashSet;
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use rusqlite::OptionalExtension;
8use uuid::Uuid;
9
10use khive_storage::attachment::{Attachment, AttachmentSubstrate};
11use khive_storage::entity::{Entity, EntityFilter, EntityTypeCounts};
12use khive_storage::error::StorageError;
13use khive_storage::types::{
14    BatchWriteSummary, DeleteMode, Page, PageRequest, SeekCursor, SeekPage, SqlStatement, SqlValue,
15};
16use khive_storage::EntityStore;
17use khive_storage::StorageCapability;
18
19use crate::error::SqliteError;
20use crate::pool::ConnectionPool;
21use crate::sql_bridge::bind_params;
22use crate::stores::attachment::{attachment_upsert_statement, delete_record_attachments_statement};
23use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
24
25fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
26    StorageError::driver(StorageCapability::Entities, op, e)
27}
28
29fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
30    e.into_storage_error(StorageCapability::Entities, op)
31}
32
33const NAMESPACE_COUNT_CHUNK_SIZE: usize = 500;
34
35const ENTITIES_COUNT_BY_TYPE_SQL: &str = include_str!("../../sql/entities-count-by-type.sql");
36
37const ENTITY_SELECT_COLUMNS: &str =
38    "entities.id, entities.namespace, entities.kind, entities.entity_type, entities.name, \
39     entities.description, entities.properties, entities.tags, entities.created_at, \
40     entities.updated_at, entities.deleted_at, entities.merged_into, entities.merge_event_id, \
41     (SELECT attachment.content_ref FROM attachments AS attachment \
42      WHERE attachment.record_uuid = entities.id \
43        AND attachment.substrate = 'entity' AND attachment.role = 'content') AS content_ref, entities.version";
44
45// ---------------------------------------------------------------------------
46// Pure statement builders (ADR-099 B3 r6 structural cut)
47//
48// These carry NO I/O — they turn an already-computed `Entity` (or a bare id)
49// into the exact `SqlStatement` this store executes. `upsert_entity` and
50// `delete_entity` below call them and execute the result; ADR-099's atomic
51// prepare path (`khive-runtime`) calls them too, to build the same statement
52// for its own guarded, synchronous apply. One statement generator, two
53// execution mechanisms (async trait dispatch vs. synchronous atomic unit) —
54// per ADR-099's accepted "handler-logic-duplication objection" text, the
55// bulk-apply path reuses the handler's existing statement generation instead
56// of re-deriving it.
57// ---------------------------------------------------------------------------
58
59/// Insert at version one, or replace the fields and advance the existing row once.
60pub fn entity_upsert_statement(entity: &Entity) -> SqlStatement {
61    let mut statement = entity_write_statement(entity, "INSERT", "entity-upsert");
62    statement.sql.push_str(
63        " ON CONFLICT(id) DO UPDATE SET namespace=excluded.namespace, kind=excluded.kind, \
64         entity_type=excluded.entity_type, name=excluded.name, description=excluded.description, \
65         properties=excluded.properties, tags=excluded.tags, created_at=excluded.created_at, \
66         updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, \
67         merged_into=excluded.merged_into, merge_event_id=excluded.merge_event_id, \
68         version=entities.version+1",
69    );
70    statement
71}
72
73fn entity_write_statement(entity: &Entity, insert: &str, label: &str) -> SqlStatement {
74    let properties_str = entity
75        .properties
76        .as_ref()
77        .map(|v| serde_json::to_string(v).unwrap_or_default());
78    let tags_str = serde_json::to_string(&entity.tags).unwrap_or_else(|_| "[]".to_string());
79    SqlStatement {
80        sql: format!(
81            "{insert} INTO entities \
82              (id, namespace, kind, entity_type, name, description, properties, tags, \
83               created_at, updated_at, deleted_at, merged_into, merge_event_id) \
84              VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)"
85        ),
86        params: vec![
87            SqlValue::Text(entity.id.to_string()),
88            SqlValue::Text(entity.namespace.clone()),
89            SqlValue::Text(entity.kind.clone()),
90            match &entity.entity_type {
91                Some(t) => SqlValue::Text(t.clone()),
92                None => SqlValue::Null,
93            },
94            SqlValue::Text(entity.name.clone()),
95            match &entity.description {
96                Some(d) => SqlValue::Text(d.clone()),
97                None => SqlValue::Null,
98            },
99            match properties_str {
100                Some(p) => SqlValue::Text(p),
101                None => SqlValue::Null,
102            },
103            SqlValue::Text(tags_str),
104            SqlValue::Integer(entity.created_at),
105            SqlValue::Integer(entity.updated_at),
106            match entity.deleted_at {
107                Some(d) => SqlValue::Integer(d),
108                None => SqlValue::Null,
109            },
110            match entity.merged_into {
111                Some(u) => SqlValue::Text(u.to_string()),
112                None => SqlValue::Null,
113            },
114            match entity.merge_event_id {
115                Some(u) => SqlValue::Text(u.to_string()),
116                None => SqlValue::Null,
117            },
118        ],
119        label: Some(label.to_string()),
120    }
121}
122
123/// Conditional-insert companion to [`entity_upsert_statement`]. Every
124/// conflict leaves the existing row untouched so a caller can read the
125/// winner and explicitly reapply its intended delta.
126pub fn entity_insert_if_absent_statement(entity: &Entity) -> SqlStatement {
127    let mut statement = entity_upsert_statement(entity);
128    statement.sql = "INSERT INTO entities \
129              (id, namespace, kind, entity_type, name, description, properties, tags, \
130               created_at, updated_at, deleted_at, merged_into, merge_event_id) \
131              VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13) \
132              ON CONFLICT DO NOTHING"
133        .to_string();
134    statement.label = Some("entity-insert-if-absent".to_string());
135    statement
136}
137
138/// Full-entity compare-and-swap update used after caller-side normalization
139/// was derived from a read snapshot. Unlike [`entity_upsert_statement`], this
140/// never inserts and cannot overwrite a row whose revision or deletion
141/// marker moved after the snapshot was read. The replacement revision must
142/// also be strictly greater than the persisted snapshot timestamp. The
143/// replacement's `version` is the expected persisted snapshot revision; the
144/// UPDATE advances it exactly once. Timestamp equality
145/// is a refused CAS, never a successful write with an unchanged concurrency
146/// token. Mirrors `note_replace_if_unchanged_statement`
147/// (`crates/khive-db/src/stores/note.rs`).
148pub fn entity_replace_if_unchanged_statement(
149    entity: &Entity,
150    expected_updated_at: i64,
151    expected_deleted_at: Option<i64>,
152) -> SqlStatement {
153    let properties_str = entity
154        .properties
155        .as_ref()
156        .map(|v| serde_json::to_string(v).unwrap_or_default());
157    let tags_str = serde_json::to_string(&entity.tags).unwrap_or_else(|_| "[]".to_string());
158    SqlStatement {
159        sql: "UPDATE entities SET \
160                namespace = ?1, kind = ?2, entity_type = ?3, name = ?4, description = ?5, \
161                properties = ?6, tags = ?7, updated_at = ?8, deleted_at = ?9, \
162                merged_into = ?10, merge_event_id = ?11, version = version + 1 \
163              WHERE id = ?12 AND updated_at = ?13 AND deleted_at IS ?14 \
164                AND ?8 > updated_at AND version = ?15"
165            .to_string(),
166        params: vec![
167            SqlValue::Text(entity.namespace.clone()),
168            SqlValue::Text(entity.kind.clone()),
169            match &entity.entity_type {
170                Some(t) => SqlValue::Text(t.clone()),
171                None => SqlValue::Null,
172            },
173            SqlValue::Text(entity.name.clone()),
174            match &entity.description {
175                Some(d) => SqlValue::Text(d.clone()),
176                None => SqlValue::Null,
177            },
178            match properties_str {
179                Some(p) => SqlValue::Text(p),
180                None => SqlValue::Null,
181            },
182            SqlValue::Text(tags_str),
183            SqlValue::Integer(entity.updated_at),
184            match entity.deleted_at {
185                Some(d) => SqlValue::Integer(d),
186                None => SqlValue::Null,
187            },
188            match entity.merged_into {
189                Some(u) => SqlValue::Text(u.to_string()),
190                None => SqlValue::Null,
191            },
192            match entity.merge_event_id {
193                Some(u) => SqlValue::Text(u.to_string()),
194                None => SqlValue::Null,
195            },
196            SqlValue::Text(entity.id.to_string()),
197            SqlValue::Integer(expected_updated_at),
198            match expected_deleted_at {
199                Some(value) => SqlValue::Integer(value),
200                None => SqlValue::Null,
201            },
202            SqlValue::Integer(entity.version),
203        ],
204        label: Some("entity-replace-if-unchanged".to_string()),
205    }
206}
207
208/// The exact soft-delete `UPDATE` this store's `delete_entity(Soft)` issues.
209pub fn entity_soft_delete_statement(id: Uuid, deleted_at: i64) -> SqlStatement {
210    SqlStatement {
211        sql: "UPDATE entities SET deleted_at = ?1, version = version + 1 WHERE id = ?2 AND deleted_at IS NULL".to_string(),
212        params: vec![
213            SqlValue::Integer(deleted_at),
214            SqlValue::Text(id.to_string()),
215        ],
216        label: Some("entity-delete-soft".to_string()),
217    }
218}
219
220/// The exact hard-delete `DELETE` this store's `delete_entity(Hard)` issues
221/// (no `deleted_at` predicate — purges live and already-tombstoned rows).
222pub fn entity_hard_delete_statement(id: Uuid) -> SqlStatement {
223    SqlStatement {
224        sql: "DELETE FROM entities WHERE id = ?1".to_string(),
225        params: vec![SqlValue::Text(id.to_string())],
226        label: Some("entity-delete-hard".to_string()),
227    }
228}
229
230/// An EntityStore backed by SQLite. Namespace is the caller's responsibility.
231///
232/// UUID is globally unique — get/delete by ID alone. Query/count use the
233/// namespace parameter as passed. Read routing is always pool-backed; the
234/// constructor's legacy file-backed flag is retained for API compatibility.
235pub struct SqlEntityStore {
236    pool: Arc<ConnectionPool>,
237    writer_task: Option<WriterTaskHandle>,
238}
239
240#[derive(Clone, Copy, PartialEq, Eq)]
241enum EntityPageMode {
242    ExactTotal,
243    CountFree,
244}
245
246impl EntityPageMode {
247    fn operation(self) -> &'static str {
248        match self {
249            Self::ExactTotal => "query_entities",
250            Self::CountFree => "query_entities_count_free",
251        }
252    }
253}
254
255impl SqlEntityStore {
256    /// Create a new store.
257    ///
258    /// When `KHIVE_WRITE_QUEUE=1` (`PoolConfig::write_queue_enabled`), every
259    /// write path on this store — the batch `upsert_entities` (its own
260    /// explicit flag check) AND every single-row write routed through the
261    /// shared `with_writer` helper (`upsert_entity`, `delete_entity`) —
262    /// routes through the pool-wide `WriterTask`
263    /// (`ConnectionPool::writer_task_handle`) instead of the legacy
264    /// pool-mutex path. The handle is a clone of the ONE writer task owned
265    /// by `pool` — constructing multiple stores (or multiple namespaces)
266    /// over the same pool never spawns more than one writer task; see
267    /// `ConnectionPool::writer_task_handle`'s doc comment for why that
268    /// matters. `None` (falling back to the legacy path for every write)
269    /// if the resolved flag is off, or if the writer task failed to spawn
270    /// (for example, an in-memory pool, which has no standalone-connection
271    /// support) — enabled by default for file-backed pools; explicit
272    /// off/degraded fallback remains possible.
273    pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
274        // Enabled by default for file-backed pools; explicit off/degraded
275        // fallback remains possible: a missing writer task — whether
276        // explicitly disabled, spawn degraded (e.g. in-memory pool), or no
277        // Tokio runtime was available at this first access (ADR-067
278        // Component A runtime-handle guard) — is cached without failing
279        // construction. Every write re-resolves it; strict mode refuses a
280        // remaining miss and compatibility mode may use the legacy path.
281        let writer_task = pool.writer_task_handle().ok().flatten();
282
283        Self { pool, writer_task }
284    }
285
286    fn current_writer_task(
287        &self,
288        operation: &'static str,
289    ) -> Result<Option<WriterTaskHandle>, StorageError> {
290        self.pool
291            .writer_task_for_write(self.writer_task.as_ref(), operation)
292    }
293
294    /// Route a single-row write through the pool-wide `WriterTask` when
295    /// the write queue is enabled and a handle is available. Strict mode
296    /// refuses a missing handle; compatibility mode falls back to the legacy
297    /// pool-mutex path.
298    ///
299    /// ADR-067 Component A (Fork C slice 2): this is the ONE routing point
300    /// for every `with_writer` caller in this store — `upsert_entity`,
301    /// `delete_entity` (soft/hard) all reach the WriterTask through this
302    /// helper rather than each duplicating the flag check. `f` must be
303    /// DML-only (a single statement, no bare `BEGIN IMMEDIATE`): on the
304    /// flag-on path it runs inside the WriterTask's own transaction, and a
305    /// nested `BEGIN IMMEDIATE` would violate SQLite's nested-transaction
306    /// rule. `upsert_entities` (the batch method) performs the same write-time
307    /// lookup first; a non-strict `None` then falls through this helper, which
308    /// records the actual compatibility fallback. Strict mode returns before
309    /// either direct-writer seam is reached.
310    async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
311    where
312        F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
313        R: Send + 'static,
314    {
315        if let Some(writer_task) = self.current_writer_task(op)? {
316            return writer_task
317                .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
318                .await;
319        }
320
321        self.pool
322            .record_direct_route(crate::timeout_sink::Site::DirectRouteEntity);
323        let pool = Arc::clone(&self.pool);
324        tokio::task::spawn_blocking(move || {
325            let guard = pool
326                .autocommit_write_unit()
327                .map_err(|e| map_sqlite_err(e, op))?;
328            f(guard.conn())
329                .map_err(|e| map_err(e, op))
330                .inspect_err(|error| pool.record_direct_writer_error(error))
331        })
332        .await
333        .map_err(|e| StorageError::driver(StorageCapability::Entities, op, e))?
334    }
335
336    /// Route multi-statement entity mutations through one write transaction.
337    /// The writer task already supplies that transaction; the direct fallback
338    /// opens and closes its own fail-closed `BEGIN IMMEDIATE` unit.
339    async fn with_writer_tx<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
340    where
341        F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
342        R: Send + 'static,
343    {
344        if let Some(writer_task) = self.current_writer_task(op)? {
345            return writer_task
346                .send_bounded(move |conn| f(conn).map_err(|error| map_err(error, op)))
347                .await;
348        }
349
350        self.pool
351            .record_direct_route(crate::timeout_sink::Site::DirectRouteEntity);
352        let pool = Arc::clone(&self.pool);
353        tokio::task::spawn_blocking(move || {
354            let guard = pool
355                .transaction_write_unit()
356                .map_err(|error| map_sqlite_err(error, op))
357                .inspect_err(|error| pool.record_direct_writer_error(error))?;
358            let conn = guard.conn();
359            let (result, terminal_state) = execute_wrapped_transaction(conn, op, move |conn| {
360                f(conn).map_err(|error| map_err(error, op))
361            });
362            if terminal_state.is_some() {
363                pool.retire_pooled_writer(conn);
364            }
365            result.inspect_err(|error| pool.record_direct_writer_error(error))
366        })
367        .await
368        .map_err(|error| StorageError::driver(StorageCapability::Entities, op, error))?
369    }
370
371    async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
372    where
373        F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
374        R: Send + 'static,
375    {
376        super::run_pooled_store_read(
377            Arc::clone(&self.pool),
378            StorageCapability::Entities,
379            op,
380            move |conn| f(conn).map_err(|error| map_err(error, op)),
381        )
382        .await
383    }
384
385    async fn with_list_reader<F, R>(
386        &self,
387        op: &'static str,
388        index: Option<&'static str>,
389        mut read: F,
390    ) -> Result<R, StorageError>
391    where
392        F: FnMut(&rusqlite::Connection, bool) -> Result<R, rusqlite::Error> + Send + 'static,
393        R: Send + 'static,
394    {
395        let dispatch = tracing::dispatcher::get_default(Clone::clone);
396        self.with_reader(op, move |conn| match read(conn, true) {
397            Ok(value) => Ok(value),
398            Err(error) => {
399                let Some(index) = index.filter(|index| missing_list_index(&error, index)) else {
400                    return Err(error);
401                };
402                tracing::dispatcher::with_default(&dispatch, || {
403                    tracing::warn!(
404                        index,
405                        operation = op,
406                        "entity list index missing; retrying without forced index"
407                    );
408                });
409                read(conn, false)
410            }
411        })
412        .await
413    }
414
415    async fn query_entities_page(
416        &self,
417        namespace: &str,
418        filter: EntityFilter,
419        page: PageRequest,
420        mode: EntityPageMode,
421    ) -> Result<Page<Entity>, StorageError> {
422        let operation = mode.operation();
423        let namespace = namespace.to_string();
424        let skip_total = mode == EntityPageMode::CountFree || is_complete_id_lookup(&filter, &page);
425        let limit_i64 = i64::from(page.limit);
426        let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
427            capability: StorageCapability::Entities,
428            operation: operation.into(),
429            message: format!(
430                "PageRequest: offset must be <= i64::MAX, got {}",
431                page.offset
432            ),
433        })?;
434
435        let streaming = mode == EntityPageMode::CountFree && is_entity_streaming_page(&filter);
436        let index = if streaming {
437            entity_count_free_index(&filter)
438        } else {
439            entity_list_index(&filter)
440        };
441        let read = move |conn: &rusqlite::Connection, force_index| {
442            let total = if filter.names_ci.is_empty() && !skip_total {
443                let (count_sql, count_params) = build_entity_where(&namespace, &filter);
444                let sql = entity_list_query(
445                    &filter,
446                    build_entity_count_query(&filter, &count_sql),
447                    force_index,
448                );
449                let mut stmt = conn.prepare(&sql)?;
450                let param_refs: Vec<&dyn rusqlite::types::ToSql> =
451                    count_params.iter().map(|p| p.as_ref()).collect();
452                Some(stmt.query_row(param_refs.as_slice(), |row| row.get::<_, i64>(0))? as u64)
453            } else {
454                None
455            };
456
457            let mut lookup_filter = filter.clone();
458            lookup_filter.names_ci.clear();
459            let effective_filter = if filter.names_ci.is_empty() {
460                &filter
461            } else {
462                &lookup_filter
463            };
464            let (where_sql, mut data_params) = if streaming {
465                build_entity_streaming_where(&namespace, effective_filter)
466            } else {
467                build_entity_where(&namespace, effective_filter)
468            };
469
470            let candidate_param_indices = if filter.names_ci.is_empty() {
471                Vec::new()
472            } else {
473                let mut candidates: Vec<String> = filter
474                    .names_ci
475                    .iter()
476                    .map(|name| name.to_ascii_lowercase())
477                    .collect();
478                candidates.sort_unstable();
479                candidates.dedup();
480                candidates
481                    .into_iter()
482                    .map(|candidate| {
483                        data_params.push(Box::new(candidate));
484                        data_params.len()
485                    })
486                    .collect()
487            };
488
489            // #818: when a name_prefix filter is active, an exact
490            // ASCII-case-insensitive match must never be pushed out of the page by
491            // pattern candidates that merely share the prefix. Rank exact
492            // matches first (deterministic tiebreak via created_at) so page
493            // truncation can never hide the record a caller resolved by name.
494            let order_by = if let Some(ref prefix) = filter.name_prefix {
495                data_params.push(Box::new(prefix.to_ascii_lowercase()));
496                format!(
497                    "CASE WHEN LOWER(name) = ?{} THEN 0 ELSE 1 END, created_at DESC, id DESC",
498                    data_params.len()
499                )
500            } else {
501                // #1671: append `id` as the final tiebreak in the primary
502                // key's direction so equal-`created_at` rows keep a fixed
503                // order across page boundaries. The deterministic total order
504                // removes tie-order instability only — offset paging can still
505                // duplicate or skip rows under concurrent inserts/deletes or
506                // sort-key updates (that would need snapshot isolation or
507                // keyset pagination).
508                "created_at DESC, id DESC".to_string()
509            };
510
511            data_params.push(Box::new(limit_i64));
512            data_params.push(Box::new(offset_i64));
513
514            let limit_idx = data_params.len() - 1;
515            let offset_idx = data_params.len();
516
517            let columns = ENTITY_SELECT_COLUMNS;
518            let data_sql = if streaming {
519                build_entity_count_free_page_query(
520                    columns,
521                    effective_filter,
522                    &where_sql,
523                    &order_by,
524                    limit_idx,
525                    offset_idx,
526                )
527            } else if filter.names_ci.is_empty() {
528                build_entity_page_query(
529                    columns,
530                    effective_filter,
531                    &where_sql,
532                    &order_by,
533                    limit_idx,
534                    offset_idx,
535                )
536            } else {
537                build_candidate_entity_query(
538                    columns,
539                    effective_filter,
540                    &where_sql,
541                    &candidate_param_indices,
542                    &order_by,
543                    limit_idx,
544                    offset_idx,
545                )
546            };
547
548            let data_sql = if streaming {
549                entity_count_free_query(effective_filter, data_sql, force_index)
550            } else {
551                entity_list_query(effective_filter, data_sql, force_index)
552            };
553            let mut stmt = conn.prepare(&data_sql)?;
554            let param_refs: Vec<&dyn rusqlite::types::ToSql> =
555                data_params.iter().map(|p| p.as_ref()).collect();
556            let rows = stmt.query_map(param_refs.as_slice(), read_entity)?;
557
558            let mut items = Vec::new();
559            for row in rows {
560                items.push(row?);
561            }
562
563            Ok(Page { items, total })
564        };
565        self.with_list_reader(operation, index, read).await
566    }
567}
568
569// =============================================================================
570// Helpers
571// =============================================================================
572
573fn read_entity(row: &rusqlite::Row<'_>) -> Result<Entity, rusqlite::Error> {
574    let id_str: String = row.get(0)?;
575    let namespace: String = row.get(1)?;
576    let kind: String = row.get(2)?;
577    let entity_type: Option<String> = row.get(3)?;
578    let name: String = row.get(4)?;
579    let description: Option<String> = row.get(5)?;
580    let properties_str: Option<String> = row.get(6)?;
581    let tags_str: String = row.get(7)?;
582    let created_at: i64 = row.get(8)?;
583    let updated_at: i64 = row.get(9)?;
584    let deleted_at: Option<i64> = row.get(10)?;
585    let merged_into_str: Option<String> = row.get(11)?;
586    let merge_event_id_str: Option<String> = row.get(12)?;
587    let content_ref: Option<String> = row.get(13)?;
588    let version: i64 = row.get(14)?;
589
590    let id = parse_uuid(&id_str)?;
591
592    let properties = properties_str
593        .map(|s| {
594            serde_json::from_str(&s).map_err(|e| {
595                rusqlite::Error::FromSqlConversionFailure(
596                    6,
597                    rusqlite::types::Type::Text,
598                    Box::new(e),
599                )
600            })
601        })
602        .transpose()?;
603
604    let tags: Vec<String> = serde_json::from_str(&tags_str).map_err(|e| {
605        rusqlite::Error::FromSqlConversionFailure(7, rusqlite::types::Type::Text, Box::new(e))
606    })?;
607
608    let merged_into = merged_into_str
609        .as_deref()
610        .map(Uuid::parse_str)
611        .transpose()
612        .map_err(|e| {
613            rusqlite::Error::FromSqlConversionFailure(10, rusqlite::types::Type::Text, Box::new(e))
614        })?;
615
616    let merge_event_id = merge_event_id_str
617        .as_deref()
618        .map(Uuid::parse_str)
619        .transpose()
620        .map_err(|e| {
621            rusqlite::Error::FromSqlConversionFailure(11, rusqlite::types::Type::Text, Box::new(e))
622        })?;
623
624    Ok(Entity {
625        id,
626        namespace,
627        kind,
628        entity_type,
629        name,
630        description,
631        properties,
632        tags,
633        created_at,
634        updated_at,
635        version,
636        deleted_at,
637        merged_into,
638        merge_event_id,
639        content_ref,
640    })
641}
642
643/// DML-only batch upsert loop shared by both the legacy (flag-off) and
644/// WriterTask-routed (flag-on) `upsert_entities` paths (ADR-067 slice 1).
645///
646/// Issues no `BEGIN` / `COMMIT` / `ROLLBACK` itself — the caller owns the
647/// enclosing transaction. Per-row failures are captured into
648/// `BatchWriteSummary::failed`/`first_error` rather than aborting the loop,
649/// matching the existing partial-success contract: this function's own
650/// `Result` is `Ok` unless a caller bug is present, since no branch here
651/// returns `Err`.
652fn batch_upsert_entities(
653    conn: &rusqlite::Connection,
654    entities: &[Entity],
655    attempted: u64,
656) -> Result<BatchWriteSummary, rusqlite::Error> {
657    let mut summary = BatchWriteSummary {
658        attempted,
659        ..BatchWriteSummary::default()
660    };
661
662    for (index, entity) in entities.iter().enumerate() {
663        let id_str = entity.id.to_string();
664        let statement = entity_upsert_statement(entity);
665        let result = (|| {
666            let mut prepared = conn.prepare(&statement.sql)?;
667            bind_params(&mut prepared, &statement.params)?;
668            prepared.raw_execute()
669        })();
670        match result {
671            Ok(_) => summary.affected = summary.affected.saturating_add(1),
672            Err(e) => {
673                let (class, retryability) = super::classify_batch_sqlite_error(&e);
674                summary.record_failure(index, Some(id_str), class, retryability, e.to_string());
675            }
676        }
677    }
678
679    Ok(summary)
680}
681
682fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
683    Uuid::parse_str(s).map_err(|e| {
684        rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
685    })
686}
687
688fn build_entity_where(
689    namespace: &str,
690    filter: &EntityFilter,
691) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
692    build_entity_where_with_mode(namespace, filter, false)
693}
694
695/// Evaluate legacy type membership on the current row instead of first
696/// materializing every matching entity ID. Deduplication also makes each
697/// single-value ordered-index prefix a single SQL bound value.
698fn build_entity_streaming_where(
699    namespace: &str,
700    filter: &EntityFilter,
701) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
702    let mut filter = filter.clone();
703    for values in [
704        &mut filter.namespaces,
705        &mut filter.kinds,
706        &mut filter.entity_types,
707    ] {
708        values.sort_unstable();
709        values.dedup();
710    }
711    for values in filter.entity_types_by_kind.values_mut() {
712        values.sort_unstable();
713        values.dedup();
714    }
715    build_entity_where_with_mode(namespace, &filter, true)
716}
717
718fn build_entity_where_with_mode(
719    namespace: &str,
720    filter: &EntityFilter,
721    row_local: bool,
722) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
723    // When filter.namespaces is non-empty use `namespace IN (...)` so that
724    // multi-namespace read visibility works.  Otherwise fall back to the
725    // single-namespace equality check for backward compatibility.
726    let (ns_condition, ns_params): (String, Vec<Box<dyn rusqlite::types::ToSql>>) =
727        if !filter.namespaces.is_empty() {
728            let placeholders: Vec<String> = (1..=filter.namespaces.len())
729                .map(|i| format!("?{i}"))
730                .collect();
731            let params: Vec<Box<dyn rusqlite::types::ToSql>> = filter
732                .namespaces
733                .iter()
734                .map(|ns| -> Box<dyn rusqlite::types::ToSql> { Box::new(ns.clone()) })
735                .collect();
736            (
737                format!("namespace IN ({})", placeholders.join(", ")),
738                params,
739            )
740        } else {
741            (
742                "namespace = ?1".to_string(),
743                vec![Box::new(namespace.to_string())],
744            )
745        };
746
747    let mut conditions: Vec<String> = vec![ns_condition, "deleted_at IS NULL".to_string()];
748    let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = ns_params;
749
750    if !filter.ids.is_empty() {
751        let placeholders: Vec<String> = filter
752            .ids
753            .iter()
754            .map(|id| {
755                params.push(Box::new(id.to_string()));
756                format!("?{}", params.len())
757            })
758            .collect();
759        conditions.push(format!("id IN ({})", placeholders.join(", ")));
760    }
761
762    if !filter.kinds.is_empty() {
763        let placeholders: Vec<String> = filter
764            .kinds
765            .iter()
766            .map(|k| {
767                params.push(Box::new(k.clone()));
768                format!("?{}", params.len())
769            })
770            .collect();
771        conditions.push(format!("kind IN ({})", placeholders.join(", ")));
772    }
773
774    let type_scope = conditions.join(" AND ");
775    let type_predicate = |scope: &str, placeholders: &str| {
776        if filter.legacy_entity_type_fallback && row_local {
777            // CASE is lazy: malformed JSON never reaches json_type/extract.
778            // A present canonical type always wins over legacy properties.
779            format!(
780                "CASE WHEN entity_type IS NOT NULL THEN entity_type IN ({placeholders}) \
781                 WHEN json_valid(properties) THEN \
782                     CASE WHEN json_type(properties, '$.type') = 'text' \
783                          THEN json_extract(properties, '$.type') IN ({placeholders}) \
784                          ELSE 0 END \
785                 ELSE 0 END"
786            )
787        } else if filter.legacy_entity_type_fallback {
788            // Legacy properties can contain invalid JSON. Exclude those rows
789            // from type fallback without rewriting them. Keep json_valid as
790            // an explicit term matching the partial legacy-type index;
791            // json_type alone does not imply its validity predicate to SQLite.
792            format!(
793                "id IN (SELECT id FROM entities WHERE {scope} \
794                 AND entity_type IN ({placeholders}) \
795                 UNION ALL SELECT id FROM entities WHERE {scope} \
796                 AND entity_type IS NULL AND json_valid(properties) \
797                 AND json_type(properties, '$.type') = 'text' \
798                 AND json_extract(properties, '$.type') IN ({placeholders}))"
799            )
800        } else {
801            format!("entity_type IN ({placeholders})")
802        }
803    };
804    if !filter.entity_types.is_empty() {
805        let placeholders: Vec<String> = filter
806            .entity_types
807            .iter()
808            .map(|t| {
809                params.push(Box::new(t.clone()));
810                format!("?{}", params.len())
811            })
812            .collect();
813        conditions.push(type_predicate(&type_scope, &placeholders.join(", ")));
814    }
815
816    if !filter.entity_types_by_kind.is_empty() {
817        let mut groups = Vec::new();
818        for (kind, types) in &filter.entity_types_by_kind {
819            if types.is_empty() {
820                continue;
821            }
822            params.push(Box::new(kind.clone()));
823            let kind_param = params.len();
824            let placeholders = types
825                .iter()
826                .map(|value| {
827                    params.push(Box::new(value.clone()));
828                    format!("?{}", params.len())
829                })
830                .collect::<Vec<_>>()
831                .join(", ");
832            let scope = format!("{type_scope} AND kind = ?{kind_param}");
833            let predicate = type_predicate(&scope, &placeholders);
834            groups.push(format!("(kind = ?{kind_param} AND {predicate})"));
835        }
836        conditions.push(if groups.is_empty() {
837            "0".to_string()
838        } else {
839            format!("({})", groups.join(" OR "))
840        });
841    }
842
843    if let Some(ref prefix) = filter.name_prefix {
844        params.push(Box::new(format!(
845            "{}%",
846            khive_types::escape_like_literal(prefix)
847        )));
848        conditions.push(format!("name LIKE ?{} ESCAPE '\\'", params.len()));
849    }
850
851    if let Some(ref exact) = filter.name_exact {
852        params.push(Box::new(exact.clone()));
853        // `entities.name` has no `COLLATE NOCASE` (see sql/schema.sql), so
854        // `=` is already SQLite's default case-sensitive BINARY comparison.
855        // `COLLATE BINARY` is spelled out here so this predicate stays
856        // correct even if the column's default collation ever changes.
857        conditions.push(format!("name = ?{} COLLATE BINARY", params.len()));
858    }
859
860    if !filter.names_ci.is_empty() {
861        // ADR-104 Stage C, R1: one batched `LOWER(name) IN (...)` predicate,
862        // served by `idx_entities_namespace_name_ci (namespace, LOWER(name))`.
863        let placeholders: Vec<String> = filter
864            .names_ci
865            .iter()
866            .map(|n| {
867                params.push(Box::new(n.to_ascii_lowercase()));
868                format!("?{}", params.len())
869            })
870            .collect();
871        conditions.push(format!("LOWER(name) IN ({})", placeholders.join(", ")));
872    }
873
874    if !filter.tags_any.is_empty() {
875        let placeholders: Vec<String> = filter
876            .tags_any
877            .iter()
878            .map(|t| {
879                // Normalise to lowercase so the comparison is case-insensitive
880                // domain filter must be case-insensitive.
881                params.push(Box::new(t.to_lowercase()));
882                format!("?{}", params.len())
883            })
884            .collect();
885        conditions.push(format!(
886            "EXISTS (SELECT 1 FROM json_each(tags) WHERE LOWER(json_each.value) IN ({}))",
887            placeholders.join(", ")
888        ));
889    }
890
891    let clause = format!(" WHERE {}", conditions.join(" AND "));
892    (clause, params)
893}
894
895/// An ID-filtered read is bounded by the caller's IDs, even on a fresh
896/// database without sqlite_stat1. The other predicates can match a large
897/// namespace; letting one of their indexes drive turns a bounded lookup into
898/// a scan of every row in that namespace.
899fn entity_read_source(filter: &EntityFilter) -> &'static str {
900    if filter.ids.is_empty() {
901        "entities"
902    } else {
903        "entities INDEXED BY sqlite_autoindex_entities_1"
904    }
905}
906
907fn build_entity_count_query(filter: &EntityFilter, where_sql: &str) -> String {
908    let source = entity_list_source(filter);
909    format!("SELECT COUNT(*) FROM {source}{where_sql}")
910}
911
912fn build_entity_page_query(
913    columns: &str,
914    filter: &EntityFilter,
915    where_sql: &str,
916    order_by: &str,
917    limit_idx: usize,
918    offset_idx: usize,
919) -> String {
920    let source = entity_list_source(filter);
921    build_entity_page_query_from_source(columns, source, where_sql, order_by, limit_idx, offset_idx)
922}
923
924fn build_entity_count_free_page_query(
925    columns: &str,
926    filter: &EntityFilter,
927    where_sql: &str,
928    order_by: &str,
929    limit_idx: usize,
930    offset_idx: usize,
931) -> String {
932    let source = entity_count_free_source(filter);
933    build_entity_page_query_from_source(columns, source, where_sql, order_by, limit_idx, offset_idx)
934}
935
936fn build_entity_page_query_from_source(
937    columns: &str,
938    source: &str,
939    where_sql: &str,
940    order_by: &str,
941    limit_idx: usize,
942    offset_idx: usize,
943) -> String {
944    format!(
945        "SELECT {columns} FROM {source}{where_sql} \
946         ORDER BY {order_by} LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
947    )
948}
949
950/// Match only ordinary list shapes whose leading predicates fit these indexes.
951/// Keep selective names/kinds/tags and legacy JSON paths free to choose their
952/// own access paths. Explicit IDs always keep the bounded primary-key plan.
953fn entity_list_source(filter: &EntityFilter) -> &'static str {
954    if !filter.ids.is_empty()
955        || !filter.kinds.is_empty()
956        || !filter.entity_types_by_kind.is_empty()
957        || filter.legacy_entity_type_fallback
958        || filter.name_prefix.is_some()
959        || filter.name_exact.is_some()
960        || !filter.names_ci.is_empty()
961        || !filter.tags_any.is_empty()
962    {
963        return entity_read_source(filter);
964    }
965    if !filter.entity_types.is_empty() {
966        "entities INDEXED BY idx_entities_live_namespace_type_order"
967    } else if filter.namespaces.len() <= 1 {
968        "entities INDEXED BY idx_entities_live_namespace_order"
969    } else {
970        "entities"
971    }
972}
973
974fn entity_list_index(filter: &EntityFilter) -> Option<&'static str> {
975    if !filter.ids.is_empty() {
976        return None;
977    }
978    entity_list_source(filter).strip_prefix("entities INDEXED BY ")
979}
980
981fn is_entity_streaming_page(filter: &EntityFilter) -> bool {
982    filter.ids.is_empty()
983        && filter.name_prefix.is_none()
984        && filter.name_exact.is_none()
985        && filter.names_ci.is_empty()
986}
987
988fn has_one_distinct_value(values: &[String]) -> bool {
989    values
990        .first()
991        .is_some_and(|first| values.iter().all(|value| value == first))
992}
993
994/// Pin creation-order walks only when every leading index column is fixed.
995/// Other predicates remain row-local residual filters, so the page can stop
996/// after its limit without materializing the complete matching ID set.
997fn entity_count_free_source(filter: &EntityFilter) -> &'static str {
998    if !is_entity_streaming_page(filter) {
999        return entity_list_source(filter);
1000    }
1001    if !filter.namespaces.is_empty() && !has_one_distinct_value(&filter.namespaces) {
1002        return "entities";
1003    }
1004    // A grouped-type map whose groups are all empty matches nothing and becomes
1005    // a constant-false predicate. SQLite folds it into the WHERE clause and can
1006    // then no longer prove the `deleted_at IS NULL` term of a partial order
1007    // index, so forcing one fails with "no query solution".
1008    if !filter.entity_types_by_kind.is_empty()
1009        && filter.entity_types_by_kind.values().all(Vec::is_empty)
1010    {
1011        return "entities";
1012    }
1013
1014    let mut kind_groups = filter
1015        .entity_types_by_kind
1016        .iter()
1017        .filter(|(_, types)| !types.is_empty());
1018    let one_kind_group = kind_groups.next().is_some() && kind_groups.next().is_none();
1019    if has_one_distinct_value(&filter.kinds) || one_kind_group {
1020        "entities INDEXED BY idx_entities_live_namespace_kind_order"
1021    } else if !filter.legacy_entity_type_fallback && has_one_distinct_value(&filter.entity_types) {
1022        // Multiple distinct types give multiple ordered runs, not one global
1023        // creation order. Use the namespace index for that residual predicate.
1024        "entities INDEXED BY idx_entities_live_namespace_type_order"
1025    } else {
1026        "entities INDEXED BY idx_entities_live_namespace_order"
1027    }
1028}
1029
1030fn entity_count_free_index(filter: &EntityFilter) -> Option<&'static str> {
1031    if !filter.ids.is_empty() {
1032        return None;
1033    }
1034    entity_count_free_source(filter).strip_prefix("entities INDEXED BY ")
1035}
1036
1037fn missing_list_index(error: &rusqlite::Error, index: &str) -> bool {
1038    if error
1039        .sqlite_error()
1040        .is_none_or(|detail| detail.extended_code != rusqlite::ffi::SQLITE_ERROR)
1041    {
1042        return false;
1043    }
1044    let message = match error {
1045        rusqlite::Error::SqliteFailure(_, Some(message)) => message.as_str(),
1046        rusqlite::Error::SqlInputError { msg, .. } => msg.as_str(),
1047        _ => return false,
1048    };
1049    message.strip_prefix("no such index: ") == Some(index)
1050}
1051
1052fn entity_list_query(filter: &EntityFilter, sql: String, force_index: bool) -> String {
1053    entity_query_with_index(
1054        sql,
1055        entity_list_source(filter),
1056        entity_list_index(filter),
1057        force_index,
1058    )
1059}
1060
1061fn entity_count_free_query(filter: &EntityFilter, sql: String, force_index: bool) -> String {
1062    entity_query_with_index(
1063        sql,
1064        entity_count_free_source(filter),
1065        entity_count_free_index(filter),
1066        force_index,
1067    )
1068}
1069
1070fn entity_query_with_index(
1071    sql: String,
1072    source: &str,
1073    index: Option<&str>,
1074    force_index: bool,
1075) -> String {
1076    if force_index || index.is_none() {
1077        sql
1078    } else {
1079        sql.replacen(source, "entities", 1)
1080    }
1081}
1082
1083fn build_candidate_entity_query(
1084    columns: &str,
1085    filter: &EntityFilter,
1086    where_sql: &str,
1087    candidate_param_indices: &[usize],
1088    order_by: &str,
1089    limit_idx: usize,
1090    offset_idx: usize,
1091) -> String {
1092    let source = entity_read_source(filter);
1093    let candidate_rows = candidate_param_indices
1094        .iter()
1095        .map(|idx| format!("(?{idx})"))
1096        .collect::<Vec<_>>()
1097        .join(", ");
1098
1099    format!(
1100        "WITH candidates(folded_name) AS (VALUES {candidate_rows}), \
1101         matched_entities(entity_id) AS (\
1102             SELECT (\
1103                 SELECT id FROM {source}{where_sql} \
1104                 AND LOWER(name) = candidates.folded_name LIMIT 1\
1105             ) FROM candidates\
1106         ) \
1107         SELECT {columns} FROM entities \
1108         JOIN matched_entities ON entities.id = matched_entities.entity_id \
1109         ORDER BY {order_by} LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
1110    )
1111}
1112
1113fn build_entity_cursor_query(
1114    columns: &str,
1115    filter: &EntityFilter,
1116    where_sql: &str,
1117    limit_idx: usize,
1118) -> String {
1119    // CROSS JOIN fixes the loop order. An explicit ID set drives entities by
1120    // primary key, then looks up each sequence. Every non-ID walk drives
1121    // the sequence range first and checks the entity through its primary key;
1122    // kind/type predicates must not force a full matching-set sort.
1123    let source = entity_read_source(filter);
1124    let from_clause = if !filter.ids.is_empty() {
1125        format!("{source} CROSS JOIN entities_seq ON entities.id = entities_seq.entity_id")
1126    } else {
1127        "entities_seq CROSS JOIN entities INDEXED BY sqlite_autoindex_entities_1 \
1128         ON entities.id = entities_seq.entity_id"
1129            .to_string()
1130    };
1131    format!(
1132        "SELECT {columns}, entities_seq.seq FROM {from_clause}{where_sql} \
1133         ORDER BY entities_seq.seq ASC LIMIT ?{limit_idx}"
1134    )
1135}
1136
1137fn is_complete_id_lookup(filter: &EntityFilter, page: &PageRequest) -> bool {
1138    !filter.ids.is_empty()
1139        && filter.kinds.is_empty()
1140        && filter.entity_types.is_empty()
1141        && filter.entity_types_by_kind.is_empty()
1142        && filter.name_prefix.is_none()
1143        && filter.name_exact.is_none()
1144        && filter.tags_any.is_empty()
1145        && filter.names_ci.is_empty()
1146        && page.offset == 0
1147        && usize::try_from(page.limit).ok() == Some(filter.ids.len())
1148}
1149
1150// =============================================================================
1151// EntityStore implementation
1152// =============================================================================
1153
1154#[async_trait]
1155impl EntityStore for SqlEntityStore {
1156    async fn upsert_entity(&self, entity: Entity) -> Result<(), StorageError> {
1157        let statement = entity_upsert_statement(&entity);
1158        self.with_writer("upsert_entity", move |conn| {
1159            let mut stmt = conn.prepare(&statement.sql)?;
1160            bind_params(&mut stmt, &statement.params)?;
1161            stmt.raw_execute()?;
1162            Ok(())
1163        })
1164        .await
1165    }
1166
1167    async fn insert_entity_if_absent(&self, entity: Entity) -> Result<bool, StorageError> {
1168        let statement = entity_insert_if_absent_statement(&entity);
1169        self.with_writer("insert_entity_if_absent", move |conn| {
1170            let mut stmt = conn.prepare(&statement.sql)?;
1171            bind_params(&mut stmt, &statement.params)?;
1172            Ok(stmt.raw_execute()? > 0)
1173        })
1174        .await
1175    }
1176
1177    async fn upsert_entity_with_attachments(
1178        &self,
1179        entity: Entity,
1180        attachments: Vec<Attachment>,
1181    ) -> Result<(), StorageError> {
1182        let entity_id = entity.id;
1183        let entity_statement = entity_upsert_statement(&entity);
1184        let mut attachment_statements = Vec::with_capacity(attachments.len());
1185        for attachment in attachments {
1186            attachment.validate()?;
1187            if attachment.record_uuid != entity_id
1188                || attachment.substrate != AttachmentSubstrate::Entity
1189            {
1190                return Err(StorageError::InvalidInput {
1191                    capability: StorageCapability::Attachments,
1192                    operation: "upsert_entity_with_attachments".into(),
1193                    message: format!(
1194                        "attachment {} must target entity {entity_id}",
1195                        attachment.role
1196                    ),
1197                });
1198            }
1199            attachment_statements.push(attachment_upsert_statement(&attachment)?);
1200        }
1201
1202        self.with_writer_tx("upsert_entity_with_attachments", move |conn| {
1203            let mut entity_stmt = conn.prepare(&entity_statement.sql)?;
1204            bind_params(&mut entity_stmt, &entity_statement.params)?;
1205            entity_stmt.raw_execute()?;
1206            drop(entity_stmt);
1207
1208            for statement in attachment_statements {
1209                let mut stmt = conn.prepare(&statement.sql)?;
1210                bind_params(&mut stmt, &statement.params)?;
1211                stmt.raw_execute()?;
1212            }
1213            Ok(())
1214        })
1215        .await
1216    }
1217
1218    async fn upsert_entities(
1219        &self,
1220        entities: Vec<Entity>,
1221    ) -> Result<BatchWriteSummary, StorageError> {
1222        let attempted = entities.len() as u64;
1223
1224        // ADR-067 slice 1: when the write queue is enabled, route through
1225        // the WriterTask channel. The closure is DML-only — no BEGIN
1226        // IMMEDIATE/COMMIT/ROLLBACK here, since the WriterTask's run loop
1227        // owns the transaction and `WriteRequest::execute_and_reply` owns
1228        // the commit/rollback decision (a bare BEGIN IMMEDIATE inside this
1229        // closure would violate SQLite's nested-transaction rule).
1230        if let Some(writer_task) = self.current_writer_task("upsert_entities")? {
1231            return writer_task
1232                .send_bounded(move |conn| {
1233                    batch_upsert_entities(conn, &entities, attempted)
1234                        .map_err(|e| map_err(e, "upsert_entities"))
1235                })
1236                .await;
1237        }
1238
1239        // Explicitly disabled or degraded fallback path: byte-for-byte unchanged from pre-ADR-067
1240        // behavior — the closure owns its own BEGIN IMMEDIATE/COMMIT/ROLLBACK
1241        // via the pool-mutex writer.
1242        let origin = self.pool.origin();
1243        self.with_writer("upsert_entities", move |conn| {
1244            conn.execute_batch("BEGIN IMMEDIATE")?;
1245            let _tx_handle = khive_storage::tx_registry::register_scoped(
1246                Some("entity_upsert_batch".to_string()),
1247                origin,
1248            );
1249
1250            let summary = batch_upsert_entities(conn, &entities, attempted)?;
1251
1252            if let Err(e) = conn.execute_batch("COMMIT") {
1253                let _ = conn.execute_batch("ROLLBACK");
1254                return Err(e);
1255            }
1256            Ok(summary)
1257        })
1258        .await
1259    }
1260
1261    async fn replace_entity_if_unchanged(
1262        &self,
1263        entity: Entity,
1264        expected_updated_at: i64,
1265        expected_deleted_at: Option<i64>,
1266    ) -> Result<bool, StorageError> {
1267        let statement = entity_replace_if_unchanged_statement(
1268            &entity,
1269            expected_updated_at,
1270            expected_deleted_at,
1271        );
1272        self.with_writer("replace_entity_if_unchanged", move |conn| {
1273            let mut stmt = conn.prepare(&statement.sql)?;
1274            bind_params(&mut stmt, &statement.params)?;
1275            Ok(stmt.raw_execute()? > 0)
1276        })
1277        .await
1278    }
1279
1280    async fn get_entity(&self, id: Uuid) -> Result<Option<Entity>, StorageError> {
1281        let id_str = id.to_string();
1282
1283        self.with_reader("get_entity", move |conn| {
1284            let sql = format!(
1285                "SELECT {ENTITY_SELECT_COLUMNS} FROM entities \
1286                 WHERE entities.id = ?1 AND entities.deleted_at IS NULL"
1287            );
1288            let mut stmt = conn.prepare(&sql)?;
1289            let mut rows = stmt.query(rusqlite::params![id_str])?;
1290            match rows.next()? {
1291                Some(row) => Ok(Some(read_entity(row)?)),
1292                None => Ok(None),
1293            }
1294        })
1295        .await
1296    }
1297
1298    async fn entity_sequence(&self, id: Uuid) -> Result<Option<i64>, StorageError> {
1299        let id = id.to_string();
1300        self.with_reader("entity_sequence", move |conn| {
1301            conn.query_row(
1302                "SELECT seq FROM entities_seq WHERE entity_id = ?1",
1303                rusqlite::params![id],
1304                |row| row.get(0),
1305            )
1306            .optional()
1307        })
1308        .await
1309    }
1310
1311    async fn delete_entity(&self, id: Uuid, mode: DeleteMode) -> Result<bool, StorageError> {
1312        match mode {
1313            DeleteMode::Soft => {
1314                let now = chrono::Utc::now().timestamp_micros();
1315                let statement = entity_soft_delete_statement(id, now);
1316                self.with_writer("delete_entity_soft", move |conn| {
1317                    let mut stmt = conn.prepare(&statement.sql)?;
1318                    bind_params(&mut stmt, &statement.params)?;
1319                    Ok(stmt.raw_execute()? > 0)
1320                })
1321                .await
1322            }
1323            DeleteMode::Hard => {
1324                let entity_statement = entity_hard_delete_statement(id);
1325                let attachment_statement =
1326                    delete_record_attachments_statement(id, AttachmentSubstrate::Entity);
1327                self.with_writer_tx("delete_entity_hard", move |conn| {
1328                    let mut entity_stmt = conn.prepare(&entity_statement.sql)?;
1329                    bind_params(&mut entity_stmt, &entity_statement.params)?;
1330                    let deleted = entity_stmt.raw_execute()? > 0;
1331                    drop(entity_stmt);
1332                    if deleted {
1333                        let mut attachment_stmt = conn.prepare(&attachment_statement.sql)?;
1334                        bind_params(&mut attachment_stmt, &attachment_statement.params)?;
1335                        attachment_stmt.raw_execute()?;
1336                    }
1337                    Ok(deleted)
1338                })
1339                .await
1340            }
1341        }
1342    }
1343
1344    async fn query_entities(
1345        &self,
1346        namespace: &str,
1347        filter: EntityFilter,
1348        page: PageRequest,
1349    ) -> Result<Page<Entity>, StorageError> {
1350        self.query_entities_page(namespace, filter, page, EntityPageMode::ExactTotal)
1351            .await
1352    }
1353
1354    async fn query_entities_count_free(
1355        &self,
1356        namespace: &str,
1357        filter: EntityFilter,
1358        page: PageRequest,
1359    ) -> Result<Page<Entity>, StorageError> {
1360        self.query_entities_page(namespace, filter, page, EntityPageMode::CountFree)
1361            .await
1362    }
1363
1364    async fn query_entities_after(
1365        &self,
1366        namespace: &str,
1367        filter: EntityFilter,
1368        after: Option<SeekCursor>,
1369        limit: u32,
1370    ) -> Result<SeekPage<Entity>, StorageError> {
1371        if limit == 0 {
1372            return Ok(SeekPage::default());
1373        }
1374        if !filter.names_ci.is_empty() {
1375            return Err(StorageError::InvalidInput {
1376                capability: StorageCapability::Entities,
1377                operation: "query_entities_after".into(),
1378                message: "names_ci candidate folding is not compatible with seek pagination".into(),
1379            });
1380        }
1381
1382        let namespace = namespace.to_string();
1383        let limit_usize = limit as usize;
1384        let probe_limit_i64 = i64::from(limit) + 1;
1385        self.with_reader("query_entities_after", move |conn| {
1386            let (mut where_sql, mut params) = if filter.ids.is_empty() {
1387                build_entity_streaming_where(&namespace, &filter)
1388            } else {
1389                build_entity_where(&namespace, &filter)
1390            };
1391            if let Some(cursor) = after {
1392                params.push(Box::new(cursor.sequence));
1393                where_sql.push_str(&format!(" AND entities_seq.seq > ?{}", params.len()));
1394            }
1395            params.push(Box::new(probe_limit_i64));
1396            let limit_idx = params.len();
1397
1398            let columns = ENTITY_SELECT_COLUMNS;
1399            let sql = build_entity_cursor_query(columns, &filter, &where_sql, limit_idx);
1400            let mut stmt = conn.prepare(&sql)?;
1401            let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1402                params.iter().map(|param| param.as_ref()).collect();
1403            let rows = stmt.query_map(param_refs.as_slice(), |row| {
1404                Ok((read_entity(row)?, row.get::<_, i64>(15)?))
1405            })?;
1406            let mut entries = rows.collect::<Result<Vec<_>, _>>()?;
1407            let has_more = entries.len() > limit_usize;
1408            if has_more {
1409                entries.truncate(limit_usize);
1410            }
1411            let next_after = if has_more {
1412                entries.last().map(|(entity, sequence)| SeekCursor {
1413                    sequence: *sequence,
1414                    id: entity.id,
1415                })
1416            } else {
1417                None
1418            };
1419            let items = entries.into_iter().map(|(entity, _)| entity).collect();
1420            Ok(SeekPage { items, next_after })
1421        })
1422        .await
1423    }
1424
1425    async fn get_entity_including_deleted(&self, id: Uuid) -> Result<Option<Entity>, StorageError> {
1426        let id_str = id.to_string();
1427
1428        self.with_reader("get_entity_including_deleted", move |conn| {
1429            let sql =
1430                format!("SELECT {ENTITY_SELECT_COLUMNS} FROM entities WHERE entities.id = ?1");
1431            let mut stmt = conn.prepare(&sql)?;
1432            let mut rows = stmt.query(rusqlite::params![id_str])?;
1433            match rows.next()? {
1434                Some(row) => Ok(Some(read_entity(row)?)),
1435                None => Ok(None),
1436            }
1437        })
1438        .await
1439    }
1440
1441    async fn count_entities_by_type(
1442        &self,
1443        namespaces: &[String],
1444    ) -> Result<Option<EntityTypeCounts>, StorageError> {
1445        let namespaces =
1446            serde_json::to_string(namespaces).map_err(|error| StorageError::Serialization {
1447                capability: StorageCapability::Entities,
1448                message: error.to_string(),
1449            })?;
1450        self.with_reader("count_entities_by_type", move |conn| {
1451            let mut statement = conn.prepare(ENTITIES_COUNT_BY_TYPE_SQL)?;
1452            let rows = statement.query_map(rusqlite::params![namespaces], |row| {
1453                let count: i64 = row.get(1)?;
1454                let count = u64::try_from(count).map_err(|error| {
1455                    rusqlite::Error::FromSqlConversionFailure(
1456                        1,
1457                        rusqlite::types::Type::Integer,
1458                        Box::new(error),
1459                    )
1460                })?;
1461                Ok((row.get::<_, Option<String>>(0)?, count))
1462            })?;
1463            rows.collect::<Result<Vec<_>, _>>().map(Some)
1464        })
1465        .await
1466    }
1467
1468    async fn count_entities(
1469        &self,
1470        namespace: &str,
1471        filter: EntityFilter,
1472    ) -> Result<u64, StorageError> {
1473        let namespace = namespace.to_string();
1474
1475        let index = entity_list_index(&filter);
1476        self.with_list_reader("count_entities", index, move |conn, force_index| {
1477            if filter.namespaces.is_empty() {
1478                let (where_sql, params) = build_entity_where(&namespace, &filter);
1479                let sql = entity_list_query(
1480                    &filter,
1481                    build_entity_count_query(&filter, &where_sql),
1482                    force_index,
1483                );
1484                let mut stmt = conn.prepare(&sql)?;
1485                let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1486                    params.iter().map(|p| p.as_ref()).collect();
1487                let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
1488                return Ok(count as u64);
1489            }
1490
1491            let deduped_namespaces: Vec<String> = filter
1492                .namespaces
1493                .iter()
1494                .cloned()
1495                .collect::<HashSet<_>>()
1496                .into_iter()
1497                .collect();
1498
1499            let mut total = 0;
1500            for chunk in deduped_namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
1501                let chunk_filter = EntityFilter {
1502                    namespaces: chunk.to_vec(),
1503                    ..filter.clone()
1504                };
1505                let (where_sql, params) = build_entity_where(&namespace, &chunk_filter);
1506                let sql = entity_list_query(
1507                    &chunk_filter,
1508                    build_entity_count_query(&chunk_filter, &where_sql),
1509                    force_index && index.is_some(),
1510                );
1511                let mut stmt = conn.prepare(&sql)?;
1512                let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1513                    params.iter().map(|p| p.as_ref()).collect();
1514                let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
1515                total += count as u64;
1516            }
1517            Ok(total)
1518        })
1519        .await
1520    }
1521}
1522
1523// =============================================================================
1524// DDL
1525// =============================================================================
1526
1527const ENTITIES_DDL: &str = include_str!("../../sql/entities-ddl.sql");
1528
1529pub(crate) fn ensure_entities_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
1530    conn.execute_batch(ENTITIES_DDL)
1531}
1532
1533#[cfg(test)]
1534#[path = "entity_tests.rs"]
1535mod tests;
1536
1537#[cfg(test)]
1538#[path = "entity_type_counts_tests.rs"]
1539mod entity_type_counts_tests;
1540
1541#[cfg(test)]
1542#[path = "entity_busy_tests.rs"]
1543mod direct_busy_tests;
1544
1545#[cfg(test)]
1546#[path = "entity_list_plan_tests.rs"]
1547mod entity_list_plan_tests;