1use 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
45pub 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
123pub 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
138pub 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
208pub 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
220pub 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
230pub 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 pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
274 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 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 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 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 "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
569fn 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
643fn 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
695fn 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 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 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 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 conditions.push(format!("name = ?{} COLLATE BINARY", params.len()));
858 }
859
860 if !filter.names_ci.is_empty() {
861 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 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
895fn 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
950fn 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
994fn 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 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 "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 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#[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 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 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
1523const 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;