1use std::collections::HashSet;
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use rusqlite::OptionalExtension;
8use uuid::Uuid;
9
10use khive_storage::error::StorageError;
11use khive_storage::note::{FilterOp, Note, NoteFilter, SortDir};
12use khive_storage::types::{
13 BatchWriteSummary, DeleteMode, Page, PageRequest, SeekCursor, SeekPage, SqlStatement, SqlValue,
14};
15use khive_storage::NoteStore;
16use khive_storage::StorageCapability;
17
18use crate::error::SqliteError;
19use crate::pool::ConnectionPool;
20use crate::sql_bridge::bind_params;
21use crate::writer_task::WriterTaskHandle;
22
23fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
24 StorageError::driver(StorageCapability::Notes, op, e)
25}
26
27fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
28 StorageError::driver(StorageCapability::Notes, op, e)
29}
30
31const NAMESPACE_COUNT_CHUNK_SIZE: usize = 500;
32
33pub const NOTE_UPSERT_SQL: &str = "INSERT INTO notes \
51 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
52 properties, created_at, updated_at, deleted_at) \
53 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13) \
54 ON CONFLICT(id) DO UPDATE SET \
55 namespace = excluded.namespace, \
56 kind = excluded.kind, \
57 status = excluded.status, \
58 name = excluded.name, \
59 content = excluded.content, \
60 salience = excluded.salience, \
61 decay_factor = excluded.decay_factor, \
62 expires_at = excluded.expires_at, \
63 properties = excluded.properties, \
64 updated_at = excluded.updated_at, \
65 deleted_at = excluded.deleted_at";
66
67pub fn note_upsert_statement(note: &Note) -> SqlStatement {
69 let properties_str = note
70 .properties
71 .as_ref()
72 .map(|v| serde_json::to_string(v).unwrap_or_default());
73 SqlStatement {
74 sql: NOTE_UPSERT_SQL.to_string(),
75 params: vec![
76 SqlValue::Text(note.id.to_string()),
77 SqlValue::Text(note.namespace.clone()),
78 SqlValue::Text(note.kind.to_string()),
79 SqlValue::Text(note.status.clone()),
80 match ¬e.name {
81 Some(n) => SqlValue::Text(n.clone()),
82 None => SqlValue::Null,
83 },
84 SqlValue::Text(note.content.clone()),
85 match note.salience {
86 Some(s) => SqlValue::Float(s),
87 None => SqlValue::Null,
88 },
89 match note.decay_factor {
90 Some(d) => SqlValue::Float(d),
91 None => SqlValue::Null,
92 },
93 match note.expires_at {
94 Some(e) => SqlValue::Integer(e),
95 None => SqlValue::Null,
96 },
97 match properties_str {
98 Some(p) => SqlValue::Text(p),
99 None => SqlValue::Null,
100 },
101 SqlValue::Integer(note.created_at),
102 SqlValue::Integer(note.updated_at),
103 match note.deleted_at {
104 Some(d) => SqlValue::Integer(d),
105 None => SqlValue::Null,
106 },
107 ],
108 label: Some("note-upsert".to_string()),
109 }
110}
111
112pub fn note_update_properties_statement(
120 id: Uuid,
121 properties: &Option<serde_json::Value>,
122 updated_at: i64,
123) -> SqlStatement {
124 let properties_str = properties
125 .as_ref()
126 .map(|v| serde_json::to_string(v).unwrap_or_default());
127 SqlStatement {
128 sql: "UPDATE notes SET properties = ?1, updated_at = ?2 \
129 WHERE id = ?3 AND deleted_at IS NULL"
130 .to_string(),
131 params: vec![
132 match properties_str {
133 Some(p) => SqlValue::Text(p),
134 None => SqlValue::Null,
135 },
136 SqlValue::Integer(updated_at),
137 SqlValue::Text(id.to_string()),
138 ],
139 label: Some("note-update-properties".to_string()),
140 }
141}
142
143pub fn note_set_property_statement(
152 id: Uuid,
153 key: &str,
154 value: &serde_json::Value,
155 updated_at: i64,
156) -> Result<SqlStatement, StorageError> {
157 if key.contains('\0') {
161 return Err(StorageError::InvalidInput {
162 capability: StorageCapability::Notes,
163 operation: "set_note_property".into(),
164 message: "property key must not contain U+0000".to_string(),
165 });
166 }
167 let path = format!("$.{}", serde_json::Value::String(key.to_string()));
168 Ok(SqlStatement {
169 sql: "UPDATE notes \
170 SET properties = json_set(COALESCE(properties, '{}'), ?1, json(?2)), \
171 updated_at = ?3 \
172 WHERE id = ?4 AND deleted_at IS NULL \
173 AND (properties IS NULL OR json_type(properties) = 'object')"
174 .to_string(),
175 params: vec![
176 SqlValue::Text(path),
177 SqlValue::Text(value.to_string()),
178 SqlValue::Integer(updated_at),
179 SqlValue::Text(id.to_string()),
180 ],
181 label: Some("note-set-property".to_string()),
182 })
183}
184
185pub fn note_soft_delete_statement(id: Uuid, deleted_at: i64) -> SqlStatement {
187 SqlStatement {
188 sql: "UPDATE notes SET status = 'deleted', deleted_at = ?1 \
189 WHERE id = ?2 AND deleted_at IS NULL"
190 .to_string(),
191 params: vec![
192 SqlValue::Integer(deleted_at),
193 SqlValue::Text(id.to_string()),
194 ],
195 label: Some("note-delete-soft".to_string()),
196 }
197}
198
199pub fn note_hard_delete_statement(id: Uuid) -> SqlStatement {
201 SqlStatement {
202 sql: "DELETE FROM notes WHERE id = ?1".to_string(),
203 params: vec![SqlValue::Text(id.to_string())],
204 label: Some("note-delete-hard".to_string()),
205 }
206}
207
208pub struct SqlNoteStore {
213 pool: Arc<ConnectionPool>,
214 is_file_backed: bool,
215 writer_task: Option<WriterTaskHandle>,
216}
217
218impl SqlNoteStore {
219 pub fn new(pool: Arc<ConnectionPool>, is_file_backed: bool) -> Self {
221 let writer_task = pool.writer_task_handle().ok().flatten();
226
227 Self {
228 pool,
229 is_file_backed,
230 writer_task,
231 }
232 }
233
234 fn open_standalone_reader(&self) -> Result<rusqlite::Connection, StorageError> {
235 let config = self.pool.config();
236 let path = config.path.as_ref().ok_or_else(|| StorageError::Pool {
237 operation: "note_reader".into(),
238 message: "in-memory databases do not support standalone connections".into(),
239 })?;
240
241 let conn = rusqlite::Connection::open_with_flags(
242 path,
243 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY
244 | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX
245 | rusqlite::OpenFlags::SQLITE_OPEN_URI,
246 )
247 .map_err(|e| map_err(e, "open_note_reader"))?;
248
249 conn.busy_timeout(config.busy_timeout)
250 .map_err(|e| map_err(e, "open_note_reader"))?;
251 conn.pragma_update(None, "foreign_keys", "ON")
252 .map_err(|e| map_err(e, "open_note_reader"))?;
253 conn.pragma_update(None, "synchronous", "NORMAL")
254 .map_err(|e| map_err(e, "open_note_reader"))?;
255
256 Ok(conn)
257 }
258
259 async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
277 where
278 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
279 R: Send + 'static,
280 {
281 if let Some(writer_task) = &self.writer_task {
282 return writer_task
283 .send(move |conn| f(conn).map_err(|e| map_err(e, op)))
284 .await;
285 }
286
287 let pool = Arc::clone(&self.pool);
288 tokio::task::spawn_blocking(move || {
289 let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
290 f(guard.conn()).map_err(|e| map_err(e, op))
291 })
292 .await
293 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
294 }
295
296 async fn with_writer_tx<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
309 where
310 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
311 R: Send + 'static,
312 {
313 if let Some(writer_task) = &self.writer_task {
314 return writer_task
315 .send(move |conn| f(conn).map_err(|e| map_err(e, op)))
316 .await;
317 }
318
319 let pool = Arc::clone(&self.pool);
320 tokio::task::spawn_blocking(move || {
321 let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
322 let conn = guard.conn();
323 conn.execute_batch("BEGIN IMMEDIATE")
324 .map_err(|e| map_err(e, op))?;
325
326 match f(conn) {
327 Ok(value) => match conn.execute_batch("COMMIT") {
328 Ok(()) => Ok(value),
329 Err(e) => {
330 let _ = conn.execute_batch("ROLLBACK");
331 Err(map_err(e, op))
332 }
333 },
334 Err(e) => {
335 let _ = conn.execute_batch("ROLLBACK");
336 Err(map_err(e, op))
337 }
338 }
339 })
340 .await
341 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
342 }
343
344 async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
345 where
346 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
347 R: Send + 'static,
348 {
349 if self.is_file_backed {
350 let conn = self.open_standalone_reader()?;
351 tokio::task::spawn_blocking(move || f(&conn).map_err(|e| map_err(e, op)))
352 .await
353 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
354 } else {
355 let pool = Arc::clone(&self.pool);
356 tokio::task::spawn_blocking(move || {
357 let guard = pool.reader().map_err(|e| map_sqlite_err(e, op))?;
358 f(guard.conn()).map_err(|e| map_err(e, op))
359 })
360 .await
361 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
362 }
363 }
364}
365
366fn read_note(row: &rusqlite::Row<'_>) -> Result<Note, rusqlite::Error> {
371 let id_str: String = row.get(0)?;
372 let namespace: String = row.get(1)?;
373 let kind: String = row.get(2)?;
374 let status: String = row.get(3)?;
375 let name: Option<String> = row.get(4)?;
376 let content: String = row.get(5)?;
377 let salience: Option<f64> = row.get(6)?;
378 let decay_factor: Option<f64> = row.get(7)?;
379 let expires_at: Option<i64> = row.get(8)?;
380 let properties_str: Option<String> = row.get(9)?;
381 let created_at: i64 = row.get(10)?;
382 let updated_at: i64 = row.get(11)?;
383 let deleted_at: Option<i64> = row.get(12)?;
384
385 let id = parse_uuid(&id_str)?;
386
387 let properties = properties_str
388 .map(|s| {
389 serde_json::from_str(&s).map_err(|e| {
390 rusqlite::Error::FromSqlConversionFailure(
391 9,
392 rusqlite::types::Type::Text,
393 Box::new(e),
394 )
395 })
396 })
397 .transpose()?;
398
399 Ok(Note {
400 id,
401 namespace,
402 kind,
403 status,
404 name,
405 content,
406 salience,
407 decay_factor,
408 expires_at,
409 properties,
410 created_at,
411 updated_at,
412 deleted_at,
413 })
414}
415
416fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
417 Uuid::parse_str(s).map_err(|e| {
418 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
419 })
420}
421
422fn query_note_page_snapshot(
423 conn: &rusqlite::Connection,
424 operation: &'static str,
425 namespace: &str,
426 count_sql: &str,
427 count_params: &[Box<dyn rusqlite::types::ToSql>],
428 data_sql: &str,
429 data_params: &[Box<dyn rusqlite::types::ToSql>],
430) -> Result<Page<Note>, rusqlite::Error> {
431 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Deferred)?;
432
433 let total: i64 = {
434 let mut stmt = tx.prepare(count_sql)?;
435 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
436 count_params.iter().map(|param| param.as_ref()).collect();
437 stmt.query_row(param_refs.as_slice(), |row| row.get(0))?
438 };
439
440 #[cfg(test)]
441 tests::page_snapshot_seam::hook(operation, namespace);
442 #[cfg(not(test))]
443 let _ = (operation, namespace);
444
445 let items = {
446 let mut stmt = tx.prepare(data_sql)?;
447 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
448 data_params.iter().map(|param| param.as_ref()).collect();
449 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
450 rows.collect::<Result<Vec<_>, _>>()?
451 };
452
453 tx.commit()?;
454 Ok(Page {
455 items,
456 total: Some(total as u64),
457 })
458}
459
460fn batch_upsert_notes(
468 conn: &rusqlite::Connection,
469 notes: &[Note],
470 attempted: u64,
471) -> Result<BatchWriteSummary, rusqlite::Error> {
472 let mut affected = 0u64;
473 let mut failed = 0u64;
474 let mut first_error = String::new();
475
476 let mut stmt = conn.prepare_cached(NOTE_UPSERT_SQL)?;
481
482 for note in notes {
483 let id_str = note.id.to_string();
484 let kind_str = note.kind.to_string();
485 let status_str = note.status.clone();
486 let properties_str = note
487 .properties
488 .as_ref()
489 .map(|v| serde_json::to_string(v).unwrap_or_default());
490
491 match stmt.execute(rusqlite::params![
492 id_str,
493 ¬e.namespace,
494 kind_str,
495 status_str,
496 ¬e.name,
497 note.content,
498 note.salience,
499 note.decay_factor,
500 note.expires_at,
501 properties_str,
502 note.created_at,
503 note.updated_at,
504 note.deleted_at,
505 ]) {
506 Ok(_) => {
507 assign_note_seq(conn, &id_str)?;
508 affected += 1;
509 }
510 Err(e) => {
511 if first_error.is_empty() {
512 first_error = e.to_string();
513 }
514 failed += 1;
515 }
516 }
517 }
518
519 Ok(BatchWriteSummary {
520 attempted,
521 affected,
522 failed,
523 first_error,
524 })
525}
526
527fn assign_note_seq(conn: &rusqlite::Connection, note_id: &str) -> Result<(), rusqlite::Error> {
533 conn.execute(
534 "INSERT OR IGNORE INTO notes_seq (note_id) VALUES (?1)",
535 rusqlite::params![note_id],
536 )?;
537 Ok(())
538}
539
540fn build_note_where(
541 namespace: &str,
542 kind: Option<&str>,
543) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
544 let mut conditions: Vec<String> = vec![
545 "namespace = ?1".to_string(),
546 "deleted_at IS NULL".to_string(),
547 ];
548 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![Box::new(namespace.to_string())];
549
550 if let Some(k) = kind {
551 params.push(Box::new(k.to_string()));
552 conditions.push(format!("kind = ?{}", params.len()));
553 }
554
555 let clause = format!(" WHERE {}", conditions.join(" AND "));
556 (clause, params)
557}
558
559fn build_note_where_for_namespaces(
560 namespaces: &[String],
561 kind: Option<&str>,
562) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
563 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = namespaces
564 .iter()
565 .map(|namespace| -> Box<dyn rusqlite::types::ToSql> { Box::new(namespace.clone()) })
566 .collect();
567 let namespace_condition = match namespaces.len() {
568 0 => "0".to_string(),
569 1 => "namespace = ?1".to_string(),
570 _ => {
571 let placeholders: Vec<String> =
572 (1..=namespaces.len()).map(|i| format!("?{i}")).collect();
573 format!("namespace IN ({})", placeholders.join(", "))
574 }
575 };
576 let mut conditions = vec![namespace_condition, "deleted_at IS NULL".to_string()];
577
578 if let Some(kind) = kind {
579 params.push(Box::new(kind.to_string()));
580 conditions.push(format!("kind = ?{}", params.len()));
581 }
582
583 let clause = format!(" WHERE {}", conditions.join(" AND "));
584 (clause, params)
585}
586
587fn validate_json_path(path: &str) -> Result<(), StorageError> {
590 let valid = path.starts_with("$.")
591 && path[2..].split('.').all(|part| {
592 !part.is_empty() && part.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
593 });
594 if valid {
595 Ok(())
596 } else {
597 Err(StorageError::InvalidInput {
598 capability: StorageCapability::Notes,
599 operation: "query_notes_filtered".into(),
600 message: format!("invalid JSON path for note filter: {path:?}"),
601 })
602 }
603}
604
605fn json_extract_expr(path: &str) -> String {
606 format!("json_extract(properties, '{path}')")
607}
608
609fn json_type_expr(path: &str) -> String {
610 format!("json_type(properties, '{path}')")
611}
612
613fn sql_value_param(value: &SqlValue) -> Result<Box<dyn rusqlite::types::ToSql>, rusqlite::Error> {
614 Ok(match value {
615 SqlValue::Null => Box::new(Option::<String>::None),
616 SqlValue::Bool(v) => Box::new(*v as i64),
617 SqlValue::Integer(v) => Box::new(*v),
618 SqlValue::Float(v) => Box::new(*v),
619 SqlValue::Text(v) => Box::new(v.clone()),
620 SqlValue::Blob(v) => Box::new(v.clone()),
621 SqlValue::Json(v) => Box::new(
622 serde_json::to_string(v)
623 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?,
624 ),
625 SqlValue::Uuid(v) => Box::new(v.to_string()),
626 SqlValue::Timestamp(v) => Box::new(v.timestamp_micros()),
627 })
628}
629
630fn build_note_filter_where(
631 namespace: &str,
632 filter: &NoteFilter,
633) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
634 let (ns_condition, ns_params): (String, Vec<Box<dyn rusqlite::types::ToSql>>) =
637 if !filter.namespaces.is_empty() {
638 let placeholders: Vec<String> = (1..=filter.namespaces.len())
639 .map(|i| format!("?{i}"))
640 .collect();
641 let params: Vec<Box<dyn rusqlite::types::ToSql>> = filter
642 .namespaces
643 .iter()
644 .map(|ns| -> Box<dyn rusqlite::types::ToSql> { Box::new(ns.clone()) })
645 .collect();
646 (
647 format!("namespace IN ({})", placeholders.join(", ")),
648 params,
649 )
650 } else {
651 (
652 "namespace = ?1".to_string(),
653 vec![Box::new(namespace.to_string())],
654 )
655 };
656
657 let mut conditions = vec![ns_condition, "deleted_at IS NULL".to_string()];
658 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = ns_params;
659
660 if let Some(kind) = &filter.kind {
661 params.push(Box::new(kind.clone()));
662 conditions.push(format!("kind = ?{}", params.len()));
663 }
664
665 for pf in &filter.property_filters {
666 match &pf.op {
667 FilterOp::EqOrMissing => {
668 let expr = json_extract_expr(&pf.json_path);
669 params.push(sql_value_param(&pf.value)?);
670 conditions.push(format!(
671 "({expr} = ?{n} OR {expr} IS NULL)",
672 n = params.len()
673 ));
674 }
675 FilterOp::JsonTypeEq => {
676 let type_expr = json_type_expr(&pf.json_path);
677 params.push(sql_value_param(&pf.value)?);
678 conditions.push(format!("{type_expr} = ?{}", params.len()));
679 }
680 FilterOp::JsonTypeNeMissing => {
681 let type_expr = json_type_expr(&pf.json_path);
682 params.push(sql_value_param(&pf.value)?);
683 let n = params.len();
684 conditions.push(format!("({type_expr} IS NULL OR {type_expr} != ?{n})"));
685 }
686 FilterOp::In(values) => {
687 let expr = json_extract_expr(&pf.json_path);
688 if values.is_empty() {
689 conditions.push("0".to_string());
691 continue;
692 }
693 let mut placeholders = Vec::with_capacity(values.len());
694 for v in values {
695 params.push(sql_value_param(v)?);
696 placeholders.push(format!("?{}", params.len()));
697 }
698 conditions.push(format!("{expr} IN ({})", placeholders.join(", ")));
699 }
700 FilterOp::NotInOrMissing(values) => {
701 let expr = json_extract_expr(&pf.json_path);
702 if values.is_empty() {
703 continue;
705 }
706 let mut placeholders = Vec::with_capacity(values.len());
707 for v in values {
708 params.push(sql_value_param(v)?);
709 placeholders.push(format!("?{}", params.len()));
710 }
711 conditions.push(format!(
712 "({expr} IS NULL OR {expr} NOT IN ({}))",
713 placeholders.join(", ")
714 ));
715 }
716 _ => {
717 let expr = json_extract_expr(&pf.json_path);
718 let op = match pf.op {
719 FilterOp::Eq => "=",
720 FilterOp::Ne => "!=",
721 FilterOp::Lt => "<",
722 FilterOp::Lte => "<=",
723 FilterOp::Gt => ">",
724 FilterOp::Gte => ">=",
725 FilterOp::EqOrMissing
726 | FilterOp::JsonTypeEq
727 | FilterOp::JsonTypeNeMissing
728 | FilterOp::In(_)
729 | FilterOp::NotInOrMissing(_) => {
730 unreachable!()
731 }
732 };
733 params.push(sql_value_param(&pf.value)?);
734 conditions.push(format!("{expr} {op} ?{}", params.len()));
735 }
736 }
737 }
738
739 if let Some(min_ts) = filter.min_created_at {
740 params.push(Box::new(min_ts));
741 conditions.push(format!("created_at >= ?{}", params.len()));
742 }
743
744 Ok((format!(" WHERE {}", conditions.join(" AND ")), params))
745}
746
747#[async_trait]
752impl NoteStore for SqlNoteStore {
753 async fn upsert_note(&self, note: Note) -> Result<(), StorageError> {
754 let id_str = note.id.to_string();
755 let statement = note_upsert_statement(¬e);
756 self.with_writer_tx("upsert_note", move |conn| {
757 let mut stmt = conn.prepare_cached(&statement.sql)?;
758 bind_params(&mut stmt, &statement.params)?;
759 stmt.raw_execute()?;
760 assign_note_seq(conn, &id_str)?;
761 Ok(())
762 })
763 .await
764 }
765
766 async fn update_note_properties(
767 &self,
768 id: Uuid,
769 properties: Option<serde_json::Value>,
770 updated_at: i64,
771 ) -> Result<bool, StorageError> {
772 let statement = note_update_properties_statement(id, &properties, updated_at);
773 self.with_writer("update_note_properties", move |conn| {
774 let mut stmt = conn.prepare(&statement.sql)?;
775 bind_params(&mut stmt, &statement.params)?;
776 Ok(stmt.raw_execute()? > 0)
777 })
778 .await
779 }
780
781 async fn set_note_property(
782 &self,
783 id: Uuid,
784 key: &str,
785 value: serde_json::Value,
786 updated_at: i64,
787 ) -> Result<bool, StorageError> {
788 let statement = note_set_property_statement(id, key, &value, updated_at)?;
789 self.with_writer("set_note_property", move |conn| {
790 let mut stmt = conn.prepare(&statement.sql)?;
791 bind_params(&mut stmt, &statement.params)?;
792 Ok(stmt.raw_execute()? > 0)
793 })
794 .await
795 }
796
797 async fn try_patch_note_property(
798 &self,
799 id: Uuid,
800 namespace: &str,
801 filter: &NoteFilter,
802 json_path: &str,
803 value: serde_json::Value,
804 updated_at: i64,
805 ) -> Result<bool, StorageError> {
806 let namespace = namespace.to_string();
807 let filter = filter.clone();
808 let value_json = serde_json::to_string(&value).map_err(|e| {
809 StorageError::driver(StorageCapability::Notes, "try_patch_note_property", e)
810 })?;
811 let json_path = json_path.to_string();
812 let id_str = id.to_string();
813
814 self.with_writer("try_patch_note_property", move |conn| {
819 let (where_clause, mut params) = build_note_filter_where(&namespace, &filter)?;
820
821 let base = params.len();
827 let sql = format!(
828 "UPDATE notes SET properties = json_set(COALESCE(properties, '{{}}'), ?{p1}, json(?{p2})), \
829 updated_at = ?{p3} {where_clause} \
830 AND (properties IS NULL OR json_type(properties) = 'object') AND id = ?{p4}",
831 p1 = base + 1,
832 p2 = base + 2,
833 p3 = base + 3,
834 p4 = base + 4,
835 );
836 params.push(Box::new(json_path));
837 params.push(Box::new(value_json));
838 params.push(Box::new(updated_at));
839 params.push(Box::new(id_str));
840
841 let mut stmt = conn.prepare(&sql)?;
842 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
843 params.iter().map(|p| p.as_ref()).collect();
844 let rows = stmt.execute(param_refs.as_slice())?;
845 Ok(rows > 0)
846 })
847 .await
848 }
849
850 async fn try_insert_note(&self, note: Note) -> Result<bool, StorageError> {
851 let namespace = note.namespace.clone();
852 let id_str = note.id.to_string();
853 let kind_str = note.kind.to_string();
854 let status_str = note.status.clone();
855 let properties_str = note
856 .properties
857 .as_ref()
858 .map(|v| serde_json::to_string(v).unwrap_or_default());
859
860 let ext_id_opt: Option<String> = note
862 .properties
863 .as_ref()
864 .and_then(|v| v.get("external_id"))
865 .and_then(|v| v.as_str())
866 .filter(|s| !s.is_empty())
867 .map(|s| s.to_string());
868
869 self.with_writer_tx("try_insert_note", move |conn| {
870 let rows = conn.execute(
871 "INSERT OR IGNORE INTO notes \
872 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
873 properties, created_at, updated_at, deleted_at) \
874 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
875 rusqlite::params![
876 id_str,
877 namespace,
878 kind_str,
879 status_str,
880 note.name,
881 note.content,
882 note.salience,
883 note.decay_factor,
884 note.expires_at,
885 properties_str,
886 note.created_at,
887 note.updated_at,
888 note.deleted_at,
889 ],
890 )?;
891
892 if rows > 0 {
893 assign_note_seq(conn, &id_str)?;
894 return Ok(true);
895 }
896
897 if let Some(ref ext_id) = ext_id_opt {
903 let is_dedup: bool = conn.query_row(
904 "SELECT COUNT(*) > 0 FROM notes \
905 WHERE namespace = ?1 \
906 AND kind = ?2 \
907 AND json_extract(properties, '$.external_id') = ?3 \
908 AND deleted_at IS NULL",
909 rusqlite::params![namespace, kind_str, ext_id],
910 |row| row.get(0),
911 )?;
912 if is_dedup {
913 return Ok(false);
914 }
915 }
916
917 Err(rusqlite::Error::SqliteFailure(
920 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
921 Some(
922 "try_insert_note: INSERT ignored for a constraint other than \
923 external_id dedup; not masking as deduplication"
924 .to_string(),
925 ),
926 ))
927 })
928 .await
929 }
930
931 async fn upsert_notes(&self, notes: Vec<Note>) -> Result<BatchWriteSummary, StorageError> {
932 let attempted = notes.len() as u64;
933
934 let origin = self.pool.origin();
945 self.with_writer_tx("upsert_notes", move |conn| {
946 let _tx_handle = khive_storage::tx_registry::register_scoped(
947 Some("note_upsert_batch".to_string()),
948 origin,
949 );
950 batch_upsert_notes(conn, ¬es, attempted)
951 })
952 .await
953 }
954
955 async fn get_note(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
956 let id_str = id.to_string();
957
958 self.with_reader("get_note", move |conn| {
959 let mut stmt = conn.prepare(
960 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
961 properties, created_at, updated_at, deleted_at \
962 FROM notes WHERE id = ?1 AND deleted_at IS NULL",
963 )?;
964 let mut rows = stmt.query(rusqlite::params![id_str])?;
965 match rows.next()? {
966 Some(row) => Ok(Some(read_note(row)?)),
967 None => Ok(None),
968 }
969 })
970 .await
971 }
972
973 async fn get_note_including_deleted(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
974 let id_str = id.to_string();
975
976 self.with_reader("get_note_including_deleted", move |conn| {
977 let mut stmt = conn.prepare(
978 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
979 properties, created_at, updated_at, deleted_at \
980 FROM notes WHERE id = ?1",
981 )?;
982 let mut rows = stmt.query(rusqlite::params![id_str])?;
983 match rows.next()? {
984 Some(row) => Ok(Some(read_note(row)?)),
985 None => Ok(None),
986 }
987 })
988 .await
989 }
990
991 async fn note_sequence(&self, id: Uuid) -> Result<Option<i64>, StorageError> {
992 let id = id.to_string();
993 self.with_reader("note_sequence", move |conn| {
994 conn.query_row(
995 "SELECT seq FROM notes_seq WHERE note_id = ?1",
996 rusqlite::params![id],
997 |row| row.get(0),
998 )
999 .optional()
1000 })
1001 .await
1002 }
1003
1004 async fn get_notes_batch(&self, ids: &[Uuid]) -> Result<Vec<Note>, StorageError> {
1005 if ids.is_empty() {
1006 return Ok(vec![]);
1007 }
1008 const CHUNK: usize = 900;
1011 let id_strings: Vec<String> = ids.iter().map(|id| id.to_string()).collect();
1012
1013 let mut result = Vec::with_capacity(ids.len());
1014 for chunk in id_strings.chunks(CHUNK) {
1015 let chunk_owned = chunk.to_vec();
1016 let notes = self
1017 .with_reader("get_notes_batch", move |conn| {
1018 let placeholders: String = (1..=chunk_owned.len())
1019 .map(|i| format!("?{i}"))
1020 .collect::<Vec<_>>()
1021 .join(", ");
1022 let sql = format!(
1023 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1024 properties, created_at, updated_at, deleted_at \
1025 FROM notes WHERE id IN ({placeholders}) AND deleted_at IS NULL"
1026 );
1027 let mut stmt = conn.prepare(&sql)?;
1028 let params: Vec<&dyn rusqlite::types::ToSql> = chunk_owned
1029 .iter()
1030 .map(|s| s as &dyn rusqlite::types::ToSql)
1031 .collect();
1032 let rows = stmt.query_map(params.as_slice(), read_note)?;
1033 let mut notes = Vec::new();
1034 for row in rows {
1035 notes.push(row?);
1036 }
1037 Ok(notes)
1038 })
1039 .await?;
1040 result.extend(notes);
1041 }
1042 Ok(result)
1043 }
1044
1045 async fn delete_note(&self, id: Uuid, mode: DeleteMode) -> Result<bool, StorageError> {
1046 match mode {
1047 DeleteMode::Soft => {
1048 let now = chrono::Utc::now().timestamp_micros();
1049 let statement = note_soft_delete_statement(id, now);
1050 self.with_writer("delete_note_soft", move |conn| {
1051 let mut stmt = conn.prepare(&statement.sql)?;
1052 bind_params(&mut stmt, &statement.params)?;
1053 Ok(stmt.raw_execute()? > 0)
1054 })
1055 .await
1056 }
1057 DeleteMode::Hard => {
1058 let statement = note_hard_delete_statement(id);
1059 self.with_writer("delete_note_hard", move |conn| {
1060 let mut stmt = conn.prepare(&statement.sql)?;
1061 bind_params(&mut stmt, &statement.params)?;
1062 Ok(stmt.raw_execute()? > 0)
1063 })
1064 .await
1065 }
1066 }
1067 }
1068
1069 async fn query_notes(
1070 &self,
1071 namespace: &str,
1072 kind: Option<&str>,
1073 page: PageRequest,
1074 ) -> Result<Page<Note>, StorageError> {
1075 let namespace = namespace.to_string();
1076 let kind = kind.map(|k| k.to_string());
1077 let limit_i64 = i64::from(page.limit);
1078 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1079 capability: StorageCapability::Notes,
1080 operation: "query_notes".into(),
1081 message: format!(
1082 "PageRequest: offset must be <= i64::MAX, got {}",
1083 page.offset
1084 ),
1085 })?;
1086
1087 self.with_reader("query_notes", move |conn| {
1088 let (count_sql, count_params) = build_note_where(&namespace, kind.as_deref());
1089 let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
1090
1091 let (where_sql, mut data_params) = build_note_where(&namespace, kind.as_deref());
1092 data_params.push(Box::new(limit_i64));
1093 data_params.push(Box::new(offset_i64));
1094
1095 let limit_idx = data_params.len() - 1;
1096 let offset_idx = data_params.len();
1097
1098 let data_sql = format!(
1099 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1100 properties, created_at, updated_at, deleted_at \
1101 FROM notes{} ORDER BY created_at DESC LIMIT ?{} OFFSET ?{}",
1102 where_sql, limit_idx, offset_idx,
1103 );
1104
1105 query_note_page_snapshot(
1106 conn,
1107 "query_notes",
1108 &namespace,
1109 &count_sql,
1110 &count_params,
1111 &data_sql,
1112 &data_params,
1113 )
1114 })
1115 .await
1116 }
1117
1118 async fn query_notes_filtered(
1119 &self,
1120 namespace: &str,
1121 filter: &NoteFilter,
1122 page: PageRequest,
1123 ) -> Result<Page<Note>, StorageError> {
1124 for pf in &filter.property_filters {
1126 validate_json_path(&pf.json_path)?;
1127 }
1128 if let Some((path, _)) = &filter.order_by {
1129 validate_json_path(path)?;
1130 }
1131
1132 let namespace = namespace.to_string();
1133 let filter = filter.clone();
1134 let limit_i64 = i64::from(page.limit);
1135 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1136 capability: StorageCapability::Notes,
1137 operation: "query_notes_filtered".into(),
1138 message: format!(
1139 "PageRequest: offset must be <= i64::MAX, got {}",
1140 page.offset
1141 ),
1142 })?;
1143
1144 self.with_reader("query_notes_filtered", move |conn| {
1145 let (count_sql, count_params) = build_note_filter_where(&namespace, &filter)?;
1146 let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
1147
1148 let (where_sql, mut data_params) = build_note_filter_where(&namespace, &filter)?;
1149 data_params.push(Box::new(limit_i64));
1150 data_params.push(Box::new(offset_i64));
1151
1152 let order_clause = match &filter.order_by {
1153 Some((path, dir)) => {
1154 let dir_str = match dir {
1155 SortDir::Asc => "ASC",
1156 SortDir::Desc => "DESC",
1157 };
1158 format!(" ORDER BY {} {dir_str}", json_extract_expr(path))
1159 }
1160 None => " ORDER BY created_at DESC, id ASC".to_string(),
1161 };
1162
1163 let limit_idx = data_params.len() - 1;
1164 let offset_idx = data_params.len();
1165 let data_sql = format!(
1166 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
1167 expires_at, properties, created_at, updated_at, deleted_at \
1168 FROM notes{}{order_clause} LIMIT ?{} OFFSET ?{}",
1169 where_sql, limit_idx, offset_idx,
1170 );
1171
1172 query_note_page_snapshot(
1173 conn,
1174 "query_notes_filtered",
1175 &namespace,
1176 &count_sql,
1177 &count_params,
1178 &data_sql,
1179 &data_params,
1180 )
1181 })
1182 .await
1183 }
1184
1185 async fn query_notes_filtered_after(
1186 &self,
1187 namespace: &str,
1188 filter: &NoteFilter,
1189 after: Option<SeekCursor>,
1190 limit: u32,
1191 ) -> Result<SeekPage<Note>, StorageError> {
1192 if limit == 0 {
1193 return Ok(SeekPage::default());
1194 }
1195 if filter.order_by.is_some() {
1196 return Err(StorageError::InvalidInput {
1197 capability: StorageCapability::Notes,
1198 operation: "query_notes_filtered_after".into(),
1199 message: "custom order_by is not compatible with insertion-sequence pagination"
1200 .into(),
1201 });
1202 }
1203 for property_filter in &filter.property_filters {
1204 validate_json_path(&property_filter.json_path)?;
1205 }
1206
1207 let namespace = namespace.to_string();
1208 let filter = filter.clone();
1209 let limit_usize = limit as usize;
1210 let probe_limit_i64 = i64::from(limit) + 1;
1211 self.with_reader("query_notes_filtered_after", move |conn| {
1212 let (mut where_sql, mut params) = build_note_filter_where(&namespace, &filter)?;
1213 if let Some(cursor) = after {
1214 params.push(Box::new(cursor.sequence));
1215 where_sql.push_str(&format!(" AND notes_seq.seq > ?{}", params.len()));
1216 }
1217 params.push(Box::new(probe_limit_i64));
1218 let limit_idx = params.len();
1219 let sql = format!(
1222 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
1223 expires_at, properties, created_at, updated_at, deleted_at, notes_seq.seq \
1224 FROM notes_seq CROSS JOIN notes ON notes.id = notes_seq.note_id{where_sql} \
1225 ORDER BY notes_seq.seq ASC LIMIT ?{limit_idx}"
1226 );
1227 let mut stmt = conn.prepare(&sql)?;
1228 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1229 params.iter().map(|param| param.as_ref()).collect();
1230 let rows = stmt.query_map(param_refs.as_slice(), |row| {
1231 Ok((read_note(row)?, row.get::<_, i64>(13)?))
1232 })?;
1233 let mut entries = rows.collect::<Result<Vec<_>, _>>()?;
1234 let has_more = entries.len() > limit_usize;
1235 if has_more {
1236 entries.truncate(limit_usize);
1237 }
1238 let next_after = if has_more {
1239 entries.last().map(|(note, sequence)| SeekCursor {
1240 sequence: *sequence,
1241 id: note.id,
1242 })
1243 } else {
1244 None
1245 };
1246 let items = entries.into_iter().map(|(note, _)| note).collect();
1247 Ok(SeekPage { items, next_after })
1248 })
1249 .await
1250 }
1251
1252 async fn query_notes_filtered_bounded(
1253 &self,
1254 namespace: &str,
1255 filter: &NoteFilter,
1256 max_rows: u32,
1257 ) -> Result<Vec<Note>, StorageError> {
1258 for pf in &filter.property_filters {
1259 validate_json_path(&pf.json_path)?;
1260 }
1261 if let Some((path, _)) = &filter.order_by {
1262 validate_json_path(path)?;
1263 }
1264
1265 let namespace = namespace.to_string();
1266 let filter = filter.clone();
1267 let limit_i64 = i64::from(max_rows) + 1;
1268
1269 self.with_reader("query_notes_filtered_bounded", move |conn| {
1270 let (where_sql, mut data_params) = build_note_filter_where(&namespace, &filter)?;
1271 data_params.push(Box::new(limit_i64));
1272 let limit_idx = data_params.len();
1273
1274 let order_clause = match &filter.order_by {
1278 Some((path, dir)) => {
1279 let dir_str = match dir {
1280 SortDir::Asc => "ASC",
1281 SortDir::Desc => "DESC",
1282 };
1283 format!(" ORDER BY {} {dir_str}, id ASC", json_extract_expr(path))
1284 }
1285 None => " ORDER BY created_at DESC, id ASC".to_string(),
1286 };
1287
1288 let data_sql = format!(
1289 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
1290 expires_at, properties, created_at, updated_at, deleted_at \
1291 FROM notes{where_sql}{order_clause} LIMIT ?{limit_idx}",
1292 );
1293
1294 let mut stmt = conn.prepare(&data_sql)?;
1295 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1296 data_params.iter().map(|p| p.as_ref()).collect();
1297 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1298
1299 let mut items = Vec::new();
1300 for row in rows {
1301 items.push(row?);
1302 }
1303 Ok(items)
1304 })
1305 .await
1306 }
1307
1308 async fn count_notes(&self, namespace: &str, kind: Option<&str>) -> Result<u64, StorageError> {
1309 let namespace = namespace.to_string();
1310 let kind = kind.map(|k| k.to_string());
1311
1312 self.with_reader("count_notes", move |conn| {
1313 let (where_sql, params) = build_note_where(&namespace, kind.as_deref());
1314 let sql = format!("SELECT COUNT(*) FROM notes{}", where_sql);
1315 let mut stmt = conn.prepare(&sql)?;
1316 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1317 params.iter().map(|p| p.as_ref()).collect();
1318 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
1319 Ok(count as u64)
1320 })
1321 .await
1322 }
1323
1324 async fn count_notes_in_namespaces(
1325 &self,
1326 namespaces: &[String],
1327 kind: Option<&str>,
1328 ) -> Result<u64, StorageError> {
1329 let namespaces: Vec<String> = namespaces
1330 .iter()
1331 .cloned()
1332 .collect::<HashSet<_>>()
1333 .into_iter()
1334 .collect();
1335 let kind = kind.map(str::to_string);
1336
1337 self.with_reader("count_notes_in_namespaces", move |conn| {
1338 let mut total = 0;
1339 for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
1340 let (where_sql, params) = build_note_where_for_namespaces(chunk, kind.as_deref());
1341 let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
1342 let mut stmt = conn.prepare(&sql)?;
1343 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1344 params.iter().map(|p| p.as_ref()).collect();
1345 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
1346 total += count as u64;
1347 }
1348 Ok(total)
1349 })
1350 .await
1351 }
1352}
1353
1354const NOTES_DDL: &str = include_str!("../../sql/notes-ddl.sql");
1359
1360const NOTES_SEQ_REPAIR_DDL: &str = include_str!("../../sql/008-notes-seq-repair.sql");
1367
1368pub(crate) fn ensure_notes_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
1369 conn.execute_batch(NOTES_DDL)
1370}
1371
1372pub(crate) fn repair_notes_seq(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
1378 conn.execute_batch(NOTES_SEQ_REPAIR_DDL)
1379}
1380
1381#[cfg(test)]
1382#[path = "note_tests.rs"]
1383mod tests;