1pub mod transport;
4
5use std::collections::HashSet;
6use std::sync::Arc;
7
8use async_trait::async_trait;
9use rusqlite::OptionalExtension;
10use uuid::Uuid;
11
12use khive_storage::attachment::AttachmentSubstrate;
13use khive_storage::error::{StorageError, WriterTaskRequestState};
14use khive_storage::note::{
15 FilterOp, Note, NoteFilter, NoteInstantSeekAfter, NoteKeyCursor, NoteSeekAfter, NoteTagMode,
16 SortDir,
17};
18use khive_storage::types::{
19 BatchWriteSummary, BoundedCount, DeleteMode, Page, PageRequest, SeekCursor, SeekPage,
20 SqlStatement, SqlValue,
21};
22use khive_storage::NoteStore;
23use khive_storage::{StorageCapability, StorageResult};
24
25use crate::error::SqliteError;
26use crate::pool::ConnectionPool;
27use crate::sql_bridge::bind_params;
28use crate::stores::attachment::delete_record_attachments_statement;
29use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
30
31fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
32 StorageError::driver(StorageCapability::Notes, op, e)
33}
34
35fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
36 StorageError::driver(StorageCapability::Notes, op, e)
37}
38
39const NAMESPACE_COUNT_CHUNK_SIZE: usize = 500;
40
41pub fn note_key_prefix_successor(prefix: &str) -> Option<String> {
42 let mut chars: Vec<char> = prefix.chars().collect();
43 while let Some(last) = chars.pop() {
44 if last == char::MAX {
45 continue;
46 }
47 let next = if last == '\u{d7ff}' {
48 '\u{e000}'
49 } else {
50 char::from_u32(u32::from(last) + 1).expect("incremented non-max scalar")
51 };
52 chars.push(next);
53 return Some(chars.into_iter().collect());
54 }
55 None
56}
57
58pub const NOTE_UPSERT_SQL: &str = "INSERT INTO notes \
77 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
78 properties, created_at, updated_at, deleted_at, key) \
79 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) \
80 ON CONFLICT(id) DO UPDATE SET \
81 namespace = excluded.namespace, \
82 kind = excluded.kind, \
83 status = excluded.status, \
84 name = excluded.name, \
85 content = excluded.content, \
86 salience = excluded.salience, \
87 decay_factor = excluded.decay_factor, \
88 expires_at = excluded.expires_at, \
89 properties = excluded.properties, \
90 updated_at = excluded.updated_at, \
91 deleted_at = excluded.deleted_at";
92
93pub const NOTE_INSERT_IF_ABSENT_SQL: &str = "INSERT INTO notes \
102 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
103 properties, created_at, updated_at, deleted_at, key) \
104 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) \
105 ON CONFLICT(id) DO NOTHING";
106
107pub fn note_insert_if_absent_statement(note: &Note) -> SqlStatement {
109 let mut statement = note_upsert_statement(note);
110 statement.sql = NOTE_INSERT_IF_ABSENT_SQL.to_string();
111 statement.label = Some("note-insert-if-absent".to_string());
112 statement
113}
114
115pub fn note_insert_keyed_statement(note: &Note) -> SqlStatement {
117 let mut statement = note_upsert_statement(note);
118 statement.sql = "INSERT INTO notes \
119 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
120 properties, created_at, updated_at, deleted_at, key) \
121 VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14) \
122 ON CONFLICT(namespace,kind,key) WHERE key IS NOT NULL AND deleted_at IS NULL DO NOTHING"
123 .into();
124 statement.label = Some("note-keyed-create".into());
125 statement
126}
127
128pub fn note_upsert_statement(note: &Note) -> SqlStatement {
130 let properties_str = note
131 .properties
132 .as_ref()
133 .map(|v| serde_json::to_string(v).unwrap_or_default());
134 SqlStatement {
135 sql: NOTE_UPSERT_SQL.to_string(),
136 params: vec![
137 SqlValue::Text(note.id.to_string()),
138 SqlValue::Text(note.namespace.clone()),
139 SqlValue::Text(note.kind.to_string()),
140 SqlValue::Text(note.status.clone()),
141 match ¬e.name {
142 Some(n) => SqlValue::Text(n.clone()),
143 None => SqlValue::Null,
144 },
145 SqlValue::Text(note.content.clone()),
146 match note.salience {
147 Some(s) => SqlValue::Float(s),
148 None => SqlValue::Null,
149 },
150 match note.decay_factor {
151 Some(d) => SqlValue::Float(d),
152 None => SqlValue::Null,
153 },
154 match note.expires_at {
155 Some(e) => SqlValue::Integer(e),
156 None => SqlValue::Null,
157 },
158 match properties_str {
159 Some(p) => SqlValue::Text(p),
160 None => SqlValue::Null,
161 },
162 SqlValue::Integer(note.created_at),
163 SqlValue::Integer(note.updated_at),
164 match note.deleted_at {
165 Some(d) => SqlValue::Integer(d),
166 None => SqlValue::Null,
167 },
168 match ¬e.key {
169 Some(key) => SqlValue::Text(key.clone()),
170 None => SqlValue::Null,
171 },
172 ],
173 label: Some("note-upsert".to_string()),
174 }
175}
176
177pub fn note_replace_if_unchanged_statement(
184 note: &Note,
185 expected_updated_at: i64,
186 expected_deleted_at: Option<i64>,
187) -> SqlStatement {
188 let properties_str = note
189 .properties
190 .as_ref()
191 .map(|v| serde_json::to_string(v).unwrap_or_default());
192 SqlStatement {
193 sql: "UPDATE notes SET \
194 namespace = ?1, kind = ?2, status = ?3, name = ?4, content = ?5, \
195 salience = ?6, decay_factor = ?7, expires_at = ?8, properties = ?9, \
196 updated_at = ?10, deleted_at = ?11 \
197 WHERE id = ?12 AND updated_at = ?13 AND deleted_at IS ?14 \
198 AND ?10 > updated_at"
199 .to_string(),
200 params: vec![
201 SqlValue::Text(note.namespace.clone()),
202 SqlValue::Text(note.kind.to_string()),
203 SqlValue::Text(note.status.clone()),
204 match ¬e.name {
205 Some(name) => SqlValue::Text(name.clone()),
206 None => SqlValue::Null,
207 },
208 SqlValue::Text(note.content.clone()),
209 match note.salience {
210 Some(value) => SqlValue::Float(value),
211 None => SqlValue::Null,
212 },
213 match note.decay_factor {
214 Some(value) => SqlValue::Float(value),
215 None => SqlValue::Null,
216 },
217 match note.expires_at {
218 Some(value) => SqlValue::Integer(value),
219 None => SqlValue::Null,
220 },
221 match properties_str {
222 Some(value) => SqlValue::Text(value),
223 None => SqlValue::Null,
224 },
225 SqlValue::Integer(note.updated_at),
226 match note.deleted_at {
227 Some(value) => SqlValue::Integer(value),
228 None => SqlValue::Null,
229 },
230 SqlValue::Text(note.id.to_string()),
231 SqlValue::Integer(expected_updated_at),
232 match expected_deleted_at {
233 Some(value) => SqlValue::Integer(value),
234 None => SqlValue::Null,
235 },
236 ],
237 label: Some("note-replace-if-unchanged".to_string()),
238 }
239}
240
241pub fn note_metadata_replace_if_unchanged_statement(
245 note: &Note,
246 expected_updated_at: i64,
247 expected_deleted_at: Option<i64>,
248) -> SqlStatement {
249 let mut statement =
250 note_replace_if_unchanged_statement(note, expected_updated_at, expected_deleted_at);
251 statement.sql = "UPDATE notes SET status=?3, name=?4, salience=?6, decay_factor=?7, expires_at=?8, updated_at=?10 \
252 WHERE id=?12 AND updated_at=?13 AND deleted_at IS ?14 AND ?10 > updated_at \
253 AND namespace=?1 AND kind=?2 AND content=?5 AND properties IS ?9 AND deleted_at IS ?11".into();
254 statement.label = Some("stream-note-metadata-cas".into());
255 statement
256}
257
258pub fn note_update_properties_statement(
266 id: Uuid,
267 properties: &Option<serde_json::Value>,
268 updated_at: i64,
269) -> SqlStatement {
270 let properties_str = properties
271 .as_ref()
272 .map(|v| serde_json::to_string(v).unwrap_or_default());
273 SqlStatement {
274 sql: "UPDATE notes SET properties = ?1, updated_at = ?2 \
275 WHERE id = ?3 AND deleted_at IS NULL"
276 .to_string(),
277 params: vec![
278 match properties_str {
279 Some(p) => SqlValue::Text(p),
280 None => SqlValue::Null,
281 },
282 SqlValue::Integer(updated_at),
283 SqlValue::Text(id.to_string()),
284 ],
285 label: Some("note-update-properties".to_string()),
286 }
287}
288
289pub fn note_set_property_statement(
298 id: Uuid,
299 key: &str,
300 value: &serde_json::Value,
301 updated_at: i64,
302) -> Result<SqlStatement, StorageError> {
303 if key.contains('\0') {
307 return Err(StorageError::InvalidInput {
308 capability: StorageCapability::Notes,
309 operation: "set_note_property".into(),
310 message: "property key must not contain U+0000".to_string(),
311 });
312 }
313 let path = format!("$.{}", serde_json::Value::String(key.to_string()));
314 Ok(SqlStatement {
315 sql: "UPDATE notes \
316 SET properties = json_set(COALESCE(properties, '{}'), ?1, json(?2)), \
317 updated_at = ?3 \
318 WHERE id = ?4 AND deleted_at IS NULL \
319 AND (properties IS NULL OR json_type(properties) = 'object')"
320 .to_string(),
321 params: vec![
322 SqlValue::Text(path),
323 SqlValue::Text(value.to_string()),
324 SqlValue::Integer(updated_at),
325 SqlValue::Text(id.to_string()),
326 ],
327 label: Some("note-set-property".to_string()),
328 })
329}
330
331pub fn note_soft_delete_statement(id: Uuid, deleted_at: i64) -> SqlStatement {
333 SqlStatement {
334 sql: "UPDATE notes SET status = 'deleted', deleted_at = ?1 \
335 WHERE id = ?2 AND deleted_at IS NULL"
336 .to_string(),
337 params: vec![
338 SqlValue::Integer(deleted_at),
339 SqlValue::Text(id.to_string()),
340 ],
341 label: Some("note-delete-soft".to_string()),
342 }
343}
344
345pub fn note_hard_delete_statement(id: Uuid) -> SqlStatement {
347 SqlStatement {
348 sql: "DELETE FROM notes WHERE id = ?1".to_string(),
349 params: vec![SqlValue::Text(id.to_string())],
350 label: Some("note-delete-hard".to_string()),
351 }
352}
353
354pub struct SqlNoteStore {
360 pool: Arc<ConnectionPool>,
361 writer_task: Option<WriterTaskHandle>,
362}
363
364impl SqlNoteStore {
365 pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
367 let writer_task = pool.writer_task_handle().ok().flatten();
374
375 Self { pool, writer_task }
376 }
377
378 fn current_writer_task(
379 &self,
380 operation: &'static str,
381 ) -> Result<Option<WriterTaskHandle>, StorageError> {
382 self.pool
383 .writer_task_for_write(self.writer_task.as_ref(), operation)
384 }
385
386 async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
405 where
406 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
407 R: Send + 'static,
408 {
409 if let Some(writer_task) = self.current_writer_task(op)? {
410 return writer_task
411 .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
412 .await;
413 }
414
415 self.pool
416 .record_direct_route(crate::timeout_sink::Site::DirectRouteNote);
417 let pool = Arc::clone(&self.pool);
418 tokio::task::spawn_blocking(move || {
419 let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
420 f(guard.conn()).map_err(|e| map_err(e, op))
421 })
422 .await
423 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
424 }
425
426 async fn with_writer_tx<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
438 where
439 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
440 R: Send + 'static,
441 {
442 self.with_writer_tx_storage(op, move |conn| f(conn).map_err(|error| map_err(error, op)))
443 .await
444 }
445
446 async fn with_writer_tx_storage<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
447 where
448 F: FnOnce(&rusqlite::Connection) -> Result<R, StorageError> + Send + 'static,
449 R: Send + 'static,
450 {
451 if let Some(writer_task) = self.current_writer_task(op)? {
452 return writer_task.send_bounded(f).await;
453 }
454
455 self.pool
456 .record_direct_route(crate::timeout_sink::Site::DirectRouteNote);
457 let pool = Arc::clone(&self.pool);
458 tokio::task::spawn_blocking(move || {
459 let guard = pool.try_writer().map_err(|e| map_sqlite_err(e, op))?;
460 let conn = guard.conn();
461 if !conn.is_autocommit() {
462 pool.retire_pooled_writer(conn);
463 return Err(StorageError::WriterTaskTerminated {
464 request_state: WriterTaskRequestState::SideEffectsUnknown,
465 });
466 }
467 if let Err(begin_error) = conn.execute_batch("BEGIN IMMEDIATE") {
468 if !conn.is_autocommit() {
469 pool.retire_pooled_writer(conn);
470 return Err(StorageError::WriterTaskTerminated {
471 request_state: WriterTaskRequestState::SideEffectsUnknown,
472 });
473 }
474 return Err(map_err(begin_error, op));
475 }
476
477 let (result, terminal_state) = execute_wrapped_transaction(conn, op, f);
478 if terminal_state.is_some() {
479 pool.retire_pooled_writer(conn);
480 }
481 result
482 })
483 .await
484 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
485 }
486
487 async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
488 where
489 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
490 R: Send + 'static,
491 {
492 super::run_pooled_store_read(
493 Arc::clone(&self.pool),
494 StorageCapability::Notes,
495 op,
496 move |conn| f(conn).map_err(|error| map_err(error, op)),
497 )
498 .await
499 }
500}
501
502fn read_note(row: &rusqlite::Row<'_>) -> Result<Note, rusqlite::Error> {
507 let id_str: String = row.get(0)?;
508 let namespace: String = row.get(1)?;
509 let kind: String = row.get(2)?;
510 let status: String = row.get(3)?;
511 let name: Option<String> = row.get(4)?;
512 let content: String = row.get(5)?;
513 let salience: Option<f64> = row.get(6)?;
514 let decay_factor: Option<f64> = row.get(7)?;
515 let expires_at: Option<i64> = row.get(8)?;
516 let properties_str: Option<String> = row.get(9)?;
517 let created_at: i64 = row.get(10)?;
518 let updated_at: i64 = row.get(11)?;
519 let deleted_at: Option<i64> = row.get(12)?;
520 let key: Option<String> = row.get(13)?;
521 let version: i64 = row.get(14)?;
522
523 let id = parse_uuid(&id_str)?;
524
525 let properties = properties_str
526 .map(|s| {
527 serde_json::from_str(&s).map_err(|e| {
528 rusqlite::Error::FromSqlConversionFailure(
529 9,
530 rusqlite::types::Type::Text,
531 Box::new(e),
532 )
533 })
534 })
535 .transpose()?;
536
537 Ok(Note {
538 id,
539 namespace,
540 kind,
541 status,
542 name,
543 content,
544 salience,
545 decay_factor,
546 expires_at,
547 properties,
548 created_at,
549 updated_at,
550 deleted_at,
551 key,
552 version,
553 })
554}
555
556fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
557 Uuid::parse_str(s).map_err(|e| {
558 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
559 })
560}
561
562fn query_note_page_snapshot(
563 conn: &rusqlite::Connection,
564 operation: &'static str,
565 namespace: &str,
566 count_sql: &str,
567 count_params: &[Box<dyn rusqlite::types::ToSql>],
568 data_sql: &str,
569 data_params: &[Box<dyn rusqlite::types::ToSql>],
570) -> Result<Page<Note>, rusqlite::Error> {
571 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Deferred)?;
572
573 let total: i64 = {
574 let mut stmt = tx.prepare(count_sql)?;
575 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
576 count_params.iter().map(|param| param.as_ref()).collect();
577 stmt.query_row(param_refs.as_slice(), |row| row.get(0))?
578 };
579
580 #[cfg(test)]
581 tests::page_snapshot_seam::hook(operation, namespace);
582 #[cfg(not(test))]
583 let _ = (operation, namespace);
584
585 let items = {
586 let mut stmt = tx.prepare(data_sql)?;
587 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
588 data_params.iter().map(|param| param.as_ref()).collect();
589 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
590 rows.collect::<Result<Vec<_>, _>>()?
591 };
592
593 tx.commit()?;
594 Ok(Page {
595 items,
596 total: Some(total as u64),
597 })
598}
599
600fn batch_upsert_notes(
608 conn: &rusqlite::Connection,
609 notes: &[Note],
610 attempted: u64,
611) -> Result<BatchWriteSummary, rusqlite::Error> {
612 let mut summary = BatchWriteSummary {
613 attempted,
614 ..BatchWriteSummary::default()
615 };
616
617 let mut stmt = conn.prepare_cached(NOTE_UPSERT_SQL)?;
622
623 for (index, note) in notes.iter().enumerate() {
624 let id_str = note.id.to_string();
625 let kind_str = note.kind.to_string();
626 let status_str = note.status.clone();
627 let properties_str = note
628 .properties
629 .as_ref()
630 .map(|v| serde_json::to_string(v).unwrap_or_default());
631
632 match stmt.execute(rusqlite::params![
633 id_str,
634 ¬e.namespace,
635 kind_str,
636 status_str,
637 ¬e.name,
638 note.content,
639 note.salience,
640 note.decay_factor,
641 note.expires_at,
642 properties_str,
643 note.created_at,
644 note.updated_at,
645 note.deleted_at,
646 note.key,
647 ]) {
648 Ok(_) => {
649 assign_note_seq(conn, &id_str)?;
650 summary.affected = summary.affected.saturating_add(1);
651 }
652 Err(e) => {
653 let (class, retryability) = super::classify_batch_sqlite_error(&e);
654 summary.record_failure(index, Some(id_str), class, retryability, e.to_string());
655 }
656 }
657 }
658
659 Ok(summary)
660}
661
662fn assign_note_seq(conn: &rusqlite::Connection, note_id: &str) -> Result<(), rusqlite::Error> {
668 conn.execute(
669 "INSERT OR IGNORE INTO notes_seq (note_id) VALUES (?1)",
670 rusqlite::params![note_id],
671 )?;
672 Ok(())
673}
674
675fn build_note_where(
676 namespace: &str,
677 kind: Option<&str>,
678) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
679 let mut conditions: Vec<String> = vec![
680 "namespace = ?1".to_string(),
681 "deleted_at IS NULL".to_string(),
682 ];
683 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![Box::new(namespace.to_string())];
684
685 if let Some(k) = kind {
686 params.push(Box::new(k.to_string()));
687 conditions.push(format!("kind = ?{}", params.len()));
688 }
689
690 let clause = format!(" WHERE {}", conditions.join(" AND "));
691 (clause, params)
692}
693
694fn build_note_where_for_namespaces(
695 namespaces: &[String],
696 kind: Option<&str>,
697) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
698 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = namespaces
699 .iter()
700 .map(|namespace| -> Box<dyn rusqlite::types::ToSql> { Box::new(namespace.clone()) })
701 .collect();
702 let namespace_condition = match namespaces.len() {
703 0 => "0".to_string(),
704 1 => "namespace = ?1".to_string(),
705 _ => {
706 let placeholders: Vec<String> =
707 (1..=namespaces.len()).map(|i| format!("?{i}")).collect();
708 format!("namespace IN ({})", placeholders.join(", "))
709 }
710 };
711 let mut conditions = vec![namespace_condition, "deleted_at IS NULL".to_string()];
712
713 if let Some(kind) = kind {
714 params.push(Box::new(kind.to_string()));
715 conditions.push(format!("kind = ?{}", params.len()));
716 }
717
718 let clause = format!(" WHERE {}", conditions.join(" AND "));
719 (clause, params)
720}
721
722fn validate_json_path(path: &str) -> Result<(), StorageError> {
725 let valid = path.starts_with("$.")
726 && path[2..].split('.').all(|part| {
727 !part.is_empty() && part.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
728 });
729 if valid {
730 Ok(())
731 } else {
732 Err(StorageError::InvalidInput {
733 capability: StorageCapability::Notes,
734 operation: "query_notes_filtered".into(),
735 message: format!("invalid JSON path for note filter: {path:?}"),
736 })
737 }
738}
739
740fn json_extract_expr(path: &str) -> String {
741 format!("json_extract(properties, '{path}')")
742}
743
744fn json_type_expr(path: &str) -> String {
745 format!("json_type(properties, '{path}')")
746}
747
748fn note_filter_page_order_clause(filter: &NoteFilter) -> String {
752 if filter.unordered {
753 return String::new();
754 }
755 match &filter.order_by {
756 Some((path, dir)) => {
757 if filter.order_by_instant {
758 let expr = json_extract_expr(path);
759 return format!(" ORDER BY khive_rfc3339_key({expr}) ASC, {expr} ASC, id ASC");
760 }
761 let dir_str = match dir {
762 SortDir::Asc => "ASC",
763 SortDir::Desc => "DESC",
764 };
765 format!(
768 " ORDER BY {} {dir_str}, id {dir_str}",
769 json_extract_expr(path)
770 )
771 }
772 None => " ORDER BY created_at DESC, id ASC".to_string(),
775 }
776}
777
778fn json_type_literal(value: &SqlValue) -> Result<&str, rusqlite::Error> {
784 const JSON_TYPES: [&str; 8] = [
785 "true", "false", "integer", "real", "text", "array", "object", "null",
786 ];
787 match value {
788 SqlValue::Text(s) if JSON_TYPES.contains(&s.as_str()) => Ok(s.as_str()),
789 other => Err(rusqlite::Error::InvalidParameterName(format!(
790 "json_type comparison value must be one of SQLite's json_type strings \
791 ({JSON_TYPES:?}), got {other:?}"
792 ))),
793 }
794}
795
796fn text_prefix_upper_bound(prefix: &str) -> Option<String> {
804 let mut chars: Vec<char> = prefix.chars().collect();
805 while let Some(last) = chars.pop() {
806 let mut next = u32::from(last) + 1;
807 if (0xD800..=0xDFFF).contains(&next) {
808 next = 0xE000;
809 }
810 if let Some(next) = char::from_u32(next) {
811 chars.push(next);
812 return Some(chars.into_iter().collect());
813 }
814 }
815 None
816}
817
818fn sql_value_param(value: &SqlValue) -> Result<Box<dyn rusqlite::types::ToSql>, rusqlite::Error> {
819 Ok(match value {
820 SqlValue::Null => Box::new(Option::<String>::None),
821 SqlValue::Bool(v) => Box::new(*v as i64),
822 SqlValue::Integer(v) => Box::new(*v),
823 SqlValue::Float(v) => Box::new(*v),
824 SqlValue::Text(v) => Box::new(v.clone()),
825 SqlValue::Blob(v) => Box::new(v.clone()),
826 SqlValue::Json(v) => Box::new(
827 serde_json::to_string(v)
828 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?,
829 ),
830 SqlValue::Uuid(v) => Box::new(v.to_string()),
831 SqlValue::Timestamp(v) => Box::new(v.timestamp_micros()),
832 })
833}
834
835fn build_note_filter_where(
836 namespace: &str,
837 filter: &NoteFilter,
838) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
839 let (ns_condition, ns_params): (String, Vec<Box<dyn rusqlite::types::ToSql>>) =
842 if !filter.namespaces.is_empty() {
843 let placeholders: Vec<String> = (1..=filter.namespaces.len())
844 .map(|i| format!("?{i}"))
845 .collect();
846 let params: Vec<Box<dyn rusqlite::types::ToSql>> = filter
847 .namespaces
848 .iter()
849 .map(|ns| -> Box<dyn rusqlite::types::ToSql> { Box::new(ns.clone()) })
850 .collect();
851 (
852 format!("namespace IN ({})", placeholders.join(", ")),
853 params,
854 )
855 } else {
856 (
857 "namespace = ?1".to_string(),
858 vec![Box::new(namespace.to_string())],
859 )
860 };
861
862 let mut conditions = vec![ns_condition, "deleted_at IS NULL".to_string()];
863 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = ns_params;
864
865 if let Some(kind) = &filter.kind {
866 params.push(Box::new(kind.clone()));
867 conditions.push(format!("kind = ?{}", params.len()));
868 }
869
870 if let Some(since) = filter.min_updated_at {
871 params.push(Box::new(since));
872 conditions.push(format!("updated_at >= ?{}", params.len()));
873 }
874 if !filter.tags.is_empty() {
875 let mut tag_predicates = Vec::new();
876 for tag in &filter.tags {
877 params.push(Box::new(tag.clone()));
878 tag_predicates.push(format!(
879 "EXISTS (SELECT 1 FROM json_each(CASE WHEN json_type(properties,'$.tags')='array' \
880 THEN json_extract(properties,'$.tags') ELSE '[]' END) AS tag \
881 WHERE tag.type='text' AND tag.value = ?{} COLLATE NOCASE)",
882 params.len()
883 ));
884 }
885 let join = match filter.tag_mode {
886 NoteTagMode::Any => " OR ",
887 NoteTagMode::All => " AND ",
888 };
889 conditions.push(format!("({})", tag_predicates.join(join)));
890 }
891
892 for pf in &filter.property_filters {
893 match &pf.op {
894 FilterOp::Rfc3339Valid => {
895 let expr = json_extract_expr(&pf.json_path);
896 conditions.push(format!("khive_rfc3339_key({expr}) IS NOT NULL"));
897 }
898 FilterOp::Rfc3339Gte | FilterOp::Rfc3339Lte => {
899 let instant = match &pf.value {
900 SqlValue::Timestamp(instant) => *instant,
901 SqlValue::Text(text) => {
902 text.parse::<chrono::DateTime<chrono::Utc>>()
903 .map_err(|error| {
904 rusqlite::Error::ToSqlConversionFailure(Box::new(error))
905 })?
906 }
907 _ => {
908 return Err(rusqlite::Error::ToSqlConversionFailure(
909 "RFC 3339 filters require a timestamp or text value".into(),
910 ));
911 }
912 };
913 let expr = json_extract_expr(&pf.json_path);
914 let op = if matches!(&pf.op, FilterOp::Rfc3339Gte) {
915 ">="
916 } else {
917 "<="
918 };
919 params.push(Box::new(crate::pool::rfc3339_instant_key(instant)));
920 conditions.push(format!("khive_rfc3339_key({expr}) {op} ?{}", params.len()));
921 }
922 FilterOp::EqOrMissing => {
923 let expr = json_extract_expr(&pf.json_path);
924 params.push(sql_value_param(&pf.value)?);
925 conditions.push(format!(
926 "({expr} = ?{n} OR {expr} IS NULL)",
927 n = params.len()
928 ));
929 }
930 FilterOp::EqOrMissingIndexed => {
931 let expr = json_extract_expr(&pf.json_path);
932 params.push(sql_value_param(&pf.value)?);
933 conditions.push(format!("ifnull({expr}, '') = ?{}", params.len()));
934 }
935 FilterOp::TextEqOrNonText => {
936 let expr = json_extract_expr(&pf.json_path);
937 let type_expr = json_type_expr(&pf.json_path);
938 params.push(sql_value_param(&pf.value)?);
939 let n = params.len();
940 conditions.push(format!(
941 "CASE WHEN {type_expr} = 'text' THEN {expr} ELSE ?{n} END = ?{n}"
942 ));
943 }
944 FilterOp::TextInOrNonText(values) => {
945 let expr = json_extract_expr(&pf.json_path);
946 let type_expr = json_type_expr(&pf.json_path);
947 let mut placeholders = Vec::with_capacity(values.len());
948 for value in values {
949 params.push(sql_value_param(value)?);
950 placeholders.push(format!("?{}", params.len()));
951 }
952 let text_match = if placeholders.is_empty() {
953 "0".to_string()
954 } else {
955 format!("{expr} IN ({})", placeholders.join(", "))
956 };
957 conditions.push(format!(
958 "CASE WHEN {type_expr} = 'text' THEN {text_match} ELSE 1 END"
959 ));
960 }
961 FilterOp::JsonTypeEq => {
962 let type_expr = json_type_expr(&pf.json_path);
963 params.push(sql_value_param(&pf.value)?);
964 conditions.push(format!("{type_expr} = ?{}", params.len()));
965 }
966 FilterOp::JsonTypeMissing => {
967 let type_expr = json_type_expr(&pf.json_path);
968 conditions.push(format!("{type_expr} IS NULL"));
969 }
970 FilterOp::JsonTypeMissingOrNullIndexed => {
971 let expr = json_extract_expr(&pf.json_path);
972 let type_expr = json_type_expr(&pf.json_path);
973 conditions.push(format!(
974 "ifnull({expr}, '') = '' AND ({type_expr} IS NULL OR {type_expr} = 'null')"
975 ));
976 }
977 FilterOp::EqOrLegacyIndexed => {
978 let expr = json_extract_expr(&pf.json_path);
979 let type_expr = json_type_expr(&pf.json_path);
980 params.push(sql_value_param(&pf.value)?);
981 let n = params.len();
982 conditions.push(format!(
983 "ifnull({expr}, '') IN (?{n}, '') AND \
984 ({type_expr} IS NULL OR {type_expr} = 'null' OR ifnull({expr}, '') != '')"
985 ));
986 }
987 FilterOp::JsonTypeNeMissing => {
988 let type_expr = json_type_expr(&pf.json_path);
989 let literal = json_type_literal(&pf.value)?;
999 conditions.push(format!(
1000 "({type_expr} IS NULL OR {type_expr} != '{literal}')"
1001 ));
1002 }
1003 FilterOp::In(values) => {
1004 let expr = json_extract_expr(&pf.json_path);
1005 if values.is_empty() {
1006 conditions.push("0".to_string());
1008 continue;
1009 }
1010 let mut placeholders = Vec::with_capacity(values.len());
1011 for v in values {
1012 params.push(sql_value_param(v)?);
1013 placeholders.push(format!("?{}", params.len()));
1014 }
1015 conditions.push(format!("{expr} IN ({})", placeholders.join(", ")));
1016 }
1017 FilterOp::TextStartsWithIndexed => {
1018 let expr = json_extract_expr(&pf.json_path);
1019 let SqlValue::Text(prefix) = &pf.value else {
1020 return Err(rusqlite::Error::ToSqlConversionFailure(
1021 "TextStartsWithIndexed takes a text prefix in PropertyFilter.value".into(),
1022 ));
1023 };
1024 params.push(Box::new(prefix.clone()));
1025 let lower = params.len();
1026 match text_prefix_upper_bound(prefix) {
1027 Some(upper) => {
1028 params.push(Box::new(upper));
1029 conditions.push(format!(
1030 "({expr} >= ?{lower} AND {expr} < ?{})",
1031 params.len()
1032 ));
1033 }
1034 None => {
1035 let type_expr = json_type_expr(&pf.json_path);
1036 conditions.push(format!("({type_expr} = 'text' AND {expr} >= ?{lower})"));
1037 }
1038 }
1039 }
1040 FilterOp::NotInOrMissing(values) => {
1041 let expr = json_extract_expr(&pf.json_path);
1042 if values.is_empty() {
1043 continue;
1045 }
1046 let mut placeholders = Vec::with_capacity(values.len());
1047 for v in values {
1048 params.push(sql_value_param(v)?);
1049 placeholders.push(format!("?{}", params.len()));
1050 }
1051 conditions.push(format!(
1052 "({expr} IS NULL OR {expr} NOT IN ({}))",
1053 placeholders.join(", ")
1054 ));
1055 }
1056 _ => {
1057 let expr = json_extract_expr(&pf.json_path);
1058 let op = match pf.op {
1059 FilterOp::Eq => "=",
1060 FilterOp::Ne => "!=",
1061 FilterOp::Lt => "<",
1062 FilterOp::Lte => "<=",
1063 FilterOp::Gt => ">",
1064 FilterOp::Gte => ">=",
1065 FilterOp::EqOrMissing
1066 | FilterOp::EqOrMissingIndexed
1067 | FilterOp::TextEqOrNonText
1068 | FilterOp::TextInOrNonText(_)
1069 | FilterOp::JsonTypeEq
1070 | FilterOp::JsonTypeMissing
1071 | FilterOp::JsonTypeMissingOrNullIndexed
1072 | FilterOp::EqOrLegacyIndexed
1073 | FilterOp::JsonTypeNeMissing
1074 | FilterOp::In(_)
1075 | FilterOp::NotInOrMissing(_)
1076 | FilterOp::TextStartsWithIndexed => {
1077 unreachable!()
1078 }
1079 FilterOp::Rfc3339Valid | FilterOp::Rfc3339Gte | FilterOp::Rfc3339Lte => {
1080 unreachable!()
1081 }
1082 };
1083 params.push(sql_value_param(&pf.value)?);
1084 conditions.push(format!("{expr} {op} ?{}", params.len()));
1085 }
1086 }
1087 }
1088
1089 if let Some(min_ts) = filter.min_created_at {
1090 params.push(Box::new(min_ts));
1091 conditions.push(format!("created_at >= ?{}", params.len()));
1092 }
1093
1094 Ok((format!(" WHERE {}", conditions.join(" AND ")), params))
1095}
1096
1097fn comm_filter_index_clause(filter: &NoteFilter, where_sql: &str) -> &'static str {
1101 if filter.kind.as_deref() != Some("message") {
1102 return "";
1103 }
1104 let Some(predicate) = where_sql.strip_prefix(" WHERE ") else {
1105 return "";
1106 };
1107 let terms: Vec<_> = predicate.split(" AND ").collect();
1108 let numbered_param = |value: &str| {
1109 value.strip_prefix('?').is_some_and(|number| {
1110 !number.is_empty() && number.bytes().all(|byte| byte.is_ascii_digit())
1111 })
1112 };
1113 let equality = |prefix: &str| {
1114 terms
1115 .iter()
1116 .any(|term| term.strip_prefix(prefix).is_some_and(numbered_param))
1117 };
1118 let namespace = equality("namespace = ")
1119 || terms.iter().any(|term| {
1120 term.strip_prefix("namespace IN (")
1121 .and_then(|term| term.strip_suffix(')'))
1122 .is_some_and(|values| values.split(", ").all(numbered_param))
1123 });
1124 let recipient = equality("ifnull(json_extract(properties, '$.to_actor'), '') = ")
1125 || terms.contains(&"ifnull(json_extract(properties, '$.to_actor'), '') = ''")
1126 || terms.iter().any(|term| {
1127 term.strip_prefix("ifnull(json_extract(properties, '$.to_actor'), '') IN (")
1128 .and_then(|term| term.strip_suffix(", '')"))
1129 .is_some_and(numbered_param)
1130 });
1131 if !namespace
1132 || !terms.contains(&"deleted_at IS NULL")
1133 || !equality("kind = ")
1134 || !equality("json_extract(properties, '$.direction') = ")
1135 || !recipient
1136 {
1137 return "";
1138 }
1139 if terms.contains(
1140 &"(json_type(properties, '$.read') IS NULL OR json_type(properties, '$.read') != 'true')",
1141 ) {
1142 let typed_exact_recipient =
1146 equality("ifnull(json_extract(properties, '$.to_actor'), '') = ")
1147 && equality("json_type(properties, '$.to_actor') = ")
1148 && filter.property_filters.iter().any(|property| {
1149 property.json_path == "$.to_actor"
1150 && matches!(property.op, FilterOp::JsonTypeEq)
1151 && matches!(&property.value, SqlValue::Text(value) if value == "text")
1152 });
1153 if typed_exact_recipient {
1154 " INDEXED BY idx_notes_unread_probe_recipient_type_direction"
1155 } else {
1156 " INDEXED BY idx_notes_unread_probe_recipient_direction"
1157 }
1158 } else {
1159 " INDEXED BY idx_notes_message_recipient_direction"
1160 }
1161}
1162
1163fn build_note_filter_read_clause(
1164 namespace: &str,
1165 filter: &NoteFilter,
1166) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
1167 let (where_sql, params) = build_note_filter_where(namespace, filter)?;
1168 let index_clause = comm_filter_index_clause(filter, &where_sql);
1169 Ok((format!("{index_clause}{where_sql}"), params))
1170}
1171
1172const NOTE_COLUMNS: &str = "id, namespace, kind, status, name, content, salience, decay_factor, \
1179 expires_at, properties, created_at, updated_at, deleted_at, key, version";
1180
1181fn fetch_notes_after_instant(
1182 conn: &rusqlite::Connection,
1183 namespace: &str,
1184 base_filter: &NoteFilter,
1185 after: &NoteInstantSeekAfter,
1186 limit: i64,
1187) -> Result<Vec<Note>, rusqlite::Error> {
1188 let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
1189 params.push(Box::new(after.value.clone()));
1190 let value_idx = params.len();
1191 params.push(Box::new(after.id.to_string()));
1192 let id_idx = params.len();
1193 params.push(Box::new(limit));
1194 let limit_idx = params.len();
1195 let (path, _) = base_filter
1196 .order_by
1197 .as_ref()
1198 .expect("instant cursor requires order_by");
1199 let expr = json_extract_expr(path);
1200 let order_clause = note_filter_page_order_clause(base_filter);
1201 let sql = format!(
1202 "SELECT {NOTE_COLUMNS} FROM notes{where_sql} \
1203 AND (khive_rfc3339_key({expr}), {expr}, id) > \
1204 (khive_rfc3339_key(?{value_idx}), ?{value_idx}, ?{id_idx}) \
1205 {order_clause} LIMIT ?{limit_idx}"
1206 );
1207 let mut stmt = conn.prepare_cached(&sql)?;
1208 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1209 params.iter().map(|param| param.as_ref()).collect();
1210 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1211 rows.collect()
1212}
1213
1214fn fetch_notes_after(
1233 conn: &rusqlite::Connection,
1234 namespace: &str,
1235 base_filter: &NoteFilter,
1236 after: &NoteSeekAfter,
1237 limit: i64,
1238) -> Result<Vec<Note>, rusqlite::Error> {
1239 let mut items = Vec::new();
1240 if limit <= 0 {
1241 return Ok(items);
1242 }
1243
1244 {
1245 let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
1246 params.push(Box::new(after.created_at));
1247 let ts_idx = params.len();
1248 params.push(Box::new(after.id.to_string()));
1249 let id_idx = params.len();
1250 params.push(Box::new(limit));
1251 let limit_idx = params.len();
1252 let sql = format!(
1253 "SELECT {NOTE_COLUMNS} FROM notes{where_sql} AND created_at = ?{ts_idx} \
1254 AND id > ?{id_idx} ORDER BY id ASC LIMIT ?{limit_idx}"
1255 );
1256 let mut stmt = conn.prepare_cached(&sql)?;
1257 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1258 params.iter().map(|p| p.as_ref()).collect();
1259 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1260 for row in rows {
1261 items.push(row?);
1262 }
1263 }
1264
1265 let remaining = limit - items.len() as i64;
1266 if remaining > 0 {
1267 let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
1268 params.push(Box::new(after.created_at));
1269 let ts_idx = params.len();
1270 params.push(Box::new(remaining));
1271 let limit_idx = params.len();
1272 let sql = format!(
1273 "SELECT {NOTE_COLUMNS} FROM notes{where_sql} AND created_at < ?{ts_idx} \
1274 ORDER BY created_at DESC, id ASC LIMIT ?{limit_idx}"
1275 );
1276 let mut stmt = conn.prepare_cached(&sql)?;
1277 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1278 params.iter().map(|p| p.as_ref()).collect();
1279 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1280 for row in rows {
1281 items.push(row?);
1282 }
1283 }
1284
1285 Ok(items)
1286}
1287
1288fn execute_filtered_note_property_patch(
1289 conn: &rusqlite::Connection,
1290 id: Uuid,
1291 namespace: &str,
1292 filter: &NoteFilter,
1293 json_path: &str,
1294 value_json: &str,
1295 updated_at: i64,
1296) -> Result<usize, rusqlite::Error> {
1297 let (where_clause, mut params) = build_note_filter_where(namespace, filter)?;
1298
1299 let base = params.len();
1300 let sql = format!(
1301 "UPDATE notes SET properties = json_set(COALESCE(properties, '{{}}'), ?{p1}, json(?{p2})), \
1302 updated_at = ?{p3} {where_clause} \
1303 AND (properties IS NULL OR json_type(properties) = 'object') AND id = ?{p4}",
1304 p1 = base + 1,
1305 p2 = base + 2,
1306 p3 = base + 3,
1307 p4 = base + 4,
1308 );
1309 params.push(Box::new(json_path.to_string()));
1310 params.push(Box::new(value_json.to_string()));
1311 params.push(Box::new(updated_at));
1312 params.push(Box::new(id.to_string()));
1313
1314 let mut stmt = conn.prepare_cached(&sql)?;
1315 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1316 params.iter().map(|param| param.as_ref()).collect();
1317 stmt.execute(param_refs.as_slice())
1318}
1319
1320#[async_trait]
1325impl NoteStore for SqlNoteStore {
1326 async fn get_live_notes_by_key(
1327 &self,
1328 namespace: &str,
1329 key: &str,
1330 kind: Option<&str>,
1331 ) -> StorageResult<Vec<Note>> {
1332 let namespace = namespace.to_owned();
1333 let key = key.to_owned();
1334 let kind = kind.map(str::to_owned);
1335 self.with_reader("get_live_notes_by_key", move |conn| {
1336 let (mut clause, mut params) = build_note_where(&namespace, kind.as_deref());
1337 params.push(Box::new(key));
1338 clause.push_str(&format!(" AND key = ?{}", params.len()));
1339 let sql = format!("SELECT {NOTE_COLUMNS} FROM notes{clause} ORDER BY kind ASC");
1340 let mut stmt = conn.prepare(&sql)?;
1341 let params: Vec<&dyn rusqlite::types::ToSql> =
1342 params.iter().map(|p| p.as_ref()).collect();
1343 let rows = stmt.query_map(params.as_slice(), read_note)?.collect();
1344 rows
1345 })
1346 .await
1347 }
1348
1349 async fn query_keyed_notes(
1350 &self,
1351 namespace: &str,
1352 filter: &NoteFilter,
1353 prefix: &str,
1354 after: Option<&NoteKeyCursor>,
1355 page: PageRequest,
1356 ) -> StorageResult<(Vec<Note>, Option<NoteKeyCursor>)> {
1357 if !filter.namespaces.is_empty()
1358 || filter.order_by.is_some()
1359 || filter.after.is_some()
1360 || (after.is_some() && page.offset != 0)
1361 || page.limit == 0
1362 {
1363 return Err(StorageError::InvalidInput { capability: StorageCapability::Notes,
1364 operation: "query_keyed_notes".into(), message: "keyed paging requires primary namespace, keyed order and a positive limit; cursor excludes offset".into() });
1365 }
1366 for property in &filter.property_filters {
1367 validate_json_path(&property.json_path)?;
1368 }
1369 let offset = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1370 capability: StorageCapability::Notes,
1371 operation: "query_keyed_notes".into(),
1372 message: "offset exceeds the supported integer range".into(),
1373 })?;
1374 let namespace = namespace.to_owned();
1375 let filter = filter.clone();
1376 let prefix = prefix.to_owned();
1377 let after = after.cloned();
1378 self.with_reader("query_keyed_notes", move |conn| {
1379 let (mut clause, mut params) = build_note_filter_where(&namespace, &filter)?;
1380 params.push(Box::new(prefix.clone()));
1381 clause.push_str(&format!(" AND key IS NOT NULL AND key >= ?{}", params.len()));
1382 if let Some(upper) = note_key_prefix_successor(&prefix) {
1383 params.push(Box::new(upper));
1384 clause.push_str(&format!(" AND key < ?{}", params.len()));
1385 }
1386 if let Some(after) = after {
1387 params.push(Box::new(after.updated_at)); let u = params.len();
1388 params.push(Box::new(after.key)); let k = params.len();
1389 params.push(Box::new(after.id.to_string())); let id = params.len();
1390 clause.push_str(&format!(" AND (updated_at < ?{u} OR (updated_at = ?{u} AND key < ?{k}) \
1391 OR (updated_at = ?{u} AND key = ?{k} AND id > ?{id}))"));
1392 }
1393 params.push(Box::new(i64::from(page.limit) + 1)); let limit = params.len();
1394 params.push(Box::new(offset)); let offset = params.len();
1395 let sql = format!("SELECT {NOTE_COLUMNS} FROM notes{clause} ORDER BY updated_at DESC, key DESC, id ASC LIMIT ?{limit} OFFSET ?{offset}");
1396 let mut stmt = conn.prepare(&sql)?;
1397 let params: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
1398 let mut notes = stmt.query_map(params.as_slice(), read_note)?.collect::<Result<Vec<_>, _>>()?;
1399 let has_more = notes.len() > page.limit as usize;
1400 notes.truncate(page.limit as usize);
1401 let next = if has_more { notes.last().map(NoteKeyCursor::from) } else { None };
1402 Ok((notes, next))
1403 }).await
1404 }
1405
1406 async fn upsert_note(&self, note: Note) -> Result<(), StorageError> {
1407 let id_str = note.id.to_string();
1408 let statement = note_upsert_statement(¬e);
1409 self.with_writer_tx("upsert_note", move |conn| {
1410 let mut stmt = conn.prepare_cached(&statement.sql)?;
1411 bind_params(&mut stmt, &statement.params)?;
1412 stmt.raw_execute()?;
1413 assign_note_seq(conn, &id_str)?;
1414 Ok(())
1415 })
1416 .await
1417 }
1418
1419 async fn insert_note_if_absent(&self, note: Note) -> Result<bool, StorageError> {
1420 let id_str = note.id.to_string();
1421 let statement = note_insert_if_absent_statement(¬e);
1422 self.with_writer_tx("insert_note_if_absent", move |conn| {
1423 let mut stmt = conn.prepare_cached(&statement.sql)?;
1424 bind_params(&mut stmt, &statement.params)?;
1425 let inserted = stmt.raw_execute()? > 0;
1426 if inserted {
1430 assign_note_seq(conn, &id_str)?;
1431 }
1432 Ok(inserted)
1433 })
1434 .await
1435 }
1436
1437 async fn replace_note_if_unchanged(
1438 &self,
1439 note: Note,
1440 expected_updated_at: i64,
1441 expected_deleted_at: Option<i64>,
1442 ) -> Result<bool, StorageError> {
1443 let statement =
1444 note_replace_if_unchanged_statement(¬e, expected_updated_at, expected_deleted_at);
1445 self.with_writer("replace_note_if_unchanged", move |conn| {
1446 let mut stmt = conn.prepare(&statement.sql)?;
1447 bind_params(&mut stmt, &statement.params)?;
1448 Ok(stmt.raw_execute()? > 0)
1449 })
1450 .await
1451 }
1452
1453 async fn update_note_properties(
1454 &self,
1455 id: Uuid,
1456 properties: Option<serde_json::Value>,
1457 updated_at: i64,
1458 ) -> Result<bool, StorageError> {
1459 let statement = note_update_properties_statement(id, &properties, updated_at);
1460 self.with_writer("update_note_properties", move |conn| {
1461 let mut stmt = conn.prepare(&statement.sql)?;
1462 bind_params(&mut stmt, &statement.params)?;
1463 Ok(stmt.raw_execute()? > 0)
1464 })
1465 .await
1466 }
1467
1468 async fn set_note_property(
1469 &self,
1470 id: Uuid,
1471 key: &str,
1472 value: serde_json::Value,
1473 updated_at: i64,
1474 ) -> Result<bool, StorageError> {
1475 let statement = note_set_property_statement(id, key, &value, updated_at)?;
1476 self.with_writer("set_note_property", move |conn| {
1477 let mut stmt = conn.prepare(&statement.sql)?;
1478 bind_params(&mut stmt, &statement.params)?;
1479 Ok(stmt.raw_execute()? > 0)
1480 })
1481 .await
1482 }
1483
1484 async fn try_patch_note_property(
1485 &self,
1486 id: Uuid,
1487 namespace: &str,
1488 filter: &NoteFilter,
1489 json_path: &str,
1490 value: serde_json::Value,
1491 updated_at: i64,
1492 ) -> Result<bool, StorageError> {
1493 let namespace = namespace.to_string();
1494 let filter = filter.clone();
1495 let value_json = serde_json::to_string(&value).map_err(|e| {
1496 StorageError::driver(StorageCapability::Notes, "try_patch_note_property", e)
1497 })?;
1498 let json_path = json_path.to_string();
1499
1500 self.with_writer("try_patch_note_property", move |conn| {
1501 execute_filtered_note_property_patch(
1502 conn,
1503 id,
1504 &namespace,
1505 &filter,
1506 &json_path,
1507 &value_json,
1508 updated_at,
1509 )
1510 .map(|rows| rows > 0)
1511 })
1512 .await
1513 }
1514
1515 async fn patch_note_property_atomic(
1516 &self,
1517 mut ids: Vec<Uuid>,
1518 namespace: &str,
1519 filter: &NoteFilter,
1520 json_path: &str,
1521 value: serde_json::Value,
1522 updated_at: i64,
1523 ) -> Result<(), StorageError> {
1524 let mut seen = HashSet::with_capacity(ids.len());
1525 ids.retain(|id| seen.insert(*id));
1526 if ids.is_empty() {
1527 return Err(StorageError::InvalidInput {
1528 capability: StorageCapability::Notes,
1529 operation: "patch_note_property_atomic".into(),
1530 message: "at least one note id is required".to_string(),
1531 });
1532 }
1533
1534 let namespace = namespace.to_string();
1535 let filter = filter.clone();
1536 let value_json = serde_json::to_string(&value).map_err(|e| {
1537 StorageError::driver(StorageCapability::Notes, "patch_note_property_atomic", e)
1538 })?;
1539 let json_path = json_path.to_string();
1540
1541 self.with_writer_tx_storage("patch_note_property_atomic", move |conn| {
1542 for id in ids {
1543 let rows = execute_filtered_note_property_patch(
1544 conn,
1545 id,
1546 &namespace,
1547 &filter,
1548 &json_path,
1549 &value_json,
1550 updated_at,
1551 )
1552 .map_err(|error| map_err(error, "patch_note_property_atomic"))?;
1553 if rows != 1 {
1554 return Err(StorageError::Conflict {
1555 capability: StorageCapability::Notes,
1556 operation: "patch_note_property_atomic".into(),
1557 message: format!(
1558 "precondition failed for note {id}: guarded update changed {rows} rows; expected 1"
1559 ),
1560 });
1561 }
1562 }
1563 Ok(())
1564 })
1565 .await
1566 }
1567
1568 async fn try_insert_note(&self, note: Note) -> Result<bool, StorageError> {
1569 let namespace = note.namespace.clone();
1570 let id_str = note.id.to_string();
1571 let kind_str = note.kind.to_string();
1572 let status_str = note.status.clone();
1573 let properties_str = note
1574 .properties
1575 .as_ref()
1576 .map(|v| serde_json::to_string(v).unwrap_or_default());
1577
1578 let ext_id_opt: Option<String> = note
1580 .properties
1581 .as_ref()
1582 .and_then(|v| v.get("external_id"))
1583 .and_then(|v| v.as_str())
1584 .filter(|s| !s.is_empty())
1585 .map(|s| s.to_string());
1586
1587 self.with_writer_tx("try_insert_note", move |conn| {
1588 let rows = conn.execute(
1589 "INSERT OR IGNORE INTO notes \
1590 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1591 properties, created_at, updated_at, deleted_at, key) \
1592 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)",
1593 rusqlite::params![
1594 id_str,
1595 namespace,
1596 kind_str,
1597 status_str,
1598 note.name,
1599 note.content,
1600 note.salience,
1601 note.decay_factor,
1602 note.expires_at,
1603 properties_str,
1604 note.created_at,
1605 note.updated_at,
1606 note.deleted_at,
1607 note.key,
1608 ],
1609 )?;
1610
1611 if rows > 0 {
1612 assign_note_seq(conn, &id_str)?;
1613 return Ok(true);
1614 }
1615
1616 if let Some(ref ext_id) = ext_id_opt {
1622 let is_dedup: bool = conn.query_row(
1623 "SELECT COUNT(*) > 0 FROM notes \
1624 WHERE namespace = ?1 \
1625 AND kind = ?2 \
1626 AND json_extract(properties, '$.external_id') = ?3 \
1627 AND deleted_at IS NULL",
1628 rusqlite::params![namespace, kind_str, ext_id],
1629 |row| row.get(0),
1630 )?;
1631 if is_dedup {
1632 return Ok(false);
1633 }
1634 }
1635
1636 Err(rusqlite::Error::SqliteFailure(
1639 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
1640 Some(
1641 "try_insert_note: INSERT ignored for a constraint other than \
1642 external_id dedup; not masking as deduplication"
1643 .to_string(),
1644 ),
1645 ))
1646 })
1647 .await
1648 }
1649
1650 async fn upsert_notes(&self, notes: Vec<Note>) -> Result<BatchWriteSummary, StorageError> {
1651 let attempted = notes.len() as u64;
1652
1653 let origin = self.pool.origin();
1664 self.with_writer_tx("upsert_notes", move |conn| {
1665 let _tx_handle = khive_storage::tx_registry::register_scoped(
1666 Some("note_upsert_batch".to_string()),
1667 origin,
1668 );
1669 batch_upsert_notes(conn, ¬es, attempted)
1670 })
1671 .await
1672 }
1673
1674 async fn get_note(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
1675 let id_str = id.to_string();
1676
1677 self.with_reader("get_note", move |conn| {
1678 let mut stmt = conn.prepare(
1679 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1680 properties, created_at, updated_at, deleted_at, key, version \
1681 FROM notes WHERE id = ?1 AND deleted_at IS NULL",
1682 )?;
1683 let mut rows = stmt.query(rusqlite::params![id_str])?;
1684 match rows.next()? {
1685 Some(row) => Ok(Some(read_note(row)?)),
1686 None => Ok(None),
1687 }
1688 })
1689 .await
1690 }
1691
1692 async fn get_note_including_deleted(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
1693 let id_str = id.to_string();
1694
1695 self.with_reader("get_note_including_deleted", move |conn| {
1696 let mut stmt = conn.prepare(
1697 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1698 properties, created_at, updated_at, deleted_at, key, version \
1699 FROM notes WHERE id = ?1",
1700 )?;
1701 let mut rows = stmt.query(rusqlite::params![id_str])?;
1702 match rows.next()? {
1703 Some(row) => Ok(Some(read_note(row)?)),
1704 None => Ok(None),
1705 }
1706 })
1707 .await
1708 }
1709
1710 async fn note_sequence(&self, id: Uuid) -> Result<Option<i64>, StorageError> {
1711 let id = id.to_string();
1712 self.with_reader("note_sequence", move |conn| {
1713 conn.query_row(
1714 "SELECT seq FROM notes_seq WHERE note_id = ?1",
1715 rusqlite::params![id],
1716 |row| row.get(0),
1717 )
1718 .optional()
1719 })
1720 .await
1721 }
1722
1723 async fn get_notes_batch(&self, ids: &[Uuid]) -> Result<Vec<Note>, StorageError> {
1724 if ids.is_empty() {
1725 return Ok(vec![]);
1726 }
1727 const CHUNK: usize = 900;
1730 let id_strings: Vec<String> = ids.iter().map(|id| id.to_string()).collect();
1731
1732 let mut result = Vec::with_capacity(ids.len());
1733 for chunk in id_strings.chunks(CHUNK) {
1734 let chunk_owned = chunk.to_vec();
1735 let notes = self
1736 .with_reader("get_notes_batch", move |conn| {
1737 let placeholders: String = (1..=chunk_owned.len())
1738 .map(|i| format!("?{i}"))
1739 .collect::<Vec<_>>()
1740 .join(", ");
1741 let sql = format!(
1742 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1743 properties, created_at, updated_at, deleted_at, key, version \
1744 FROM notes WHERE id IN ({placeholders}) AND deleted_at IS NULL"
1745 );
1746 let mut stmt = conn.prepare(&sql)?;
1747 let params: Vec<&dyn rusqlite::types::ToSql> = chunk_owned
1748 .iter()
1749 .map(|s| s as &dyn rusqlite::types::ToSql)
1750 .collect();
1751 let rows = stmt.query_map(params.as_slice(), read_note)?;
1752 let mut notes = Vec::new();
1753 for row in rows {
1754 notes.push(row?);
1755 }
1756 Ok(notes)
1757 })
1758 .await?;
1759 result.extend(notes);
1760 }
1761 Ok(result)
1762 }
1763
1764 async fn delete_note(&self, id: Uuid, mode: DeleteMode) -> Result<bool, StorageError> {
1765 match mode {
1766 DeleteMode::Soft => {
1767 let now = chrono::Utc::now().timestamp_micros();
1768 let statement = note_soft_delete_statement(id, now);
1769 self.with_writer("delete_note_soft", move |conn| {
1770 let mut stmt = conn.prepare(&statement.sql)?;
1771 bind_params(&mut stmt, &statement.params)?;
1772 Ok(stmt.raw_execute()? > 0)
1773 })
1774 .await
1775 }
1776 DeleteMode::Hard => {
1777 let note_statement = note_hard_delete_statement(id);
1778 let attachment_statement =
1779 delete_record_attachments_statement(id, AttachmentSubstrate::Note);
1780 self.with_writer_tx("delete_note_hard", move |conn| {
1781 let mut note_stmt = conn.prepare(¬e_statement.sql)?;
1782 bind_params(&mut note_stmt, ¬e_statement.params)?;
1783 let deleted = note_stmt.raw_execute()? > 0;
1784 drop(note_stmt);
1785 if deleted {
1786 let mut attachment_stmt = conn.prepare(&attachment_statement.sql)?;
1787 bind_params(&mut attachment_stmt, &attachment_statement.params)?;
1788 attachment_stmt.raw_execute()?;
1789 }
1790 Ok(deleted)
1791 })
1792 .await
1793 }
1794 }
1795 }
1796
1797 async fn query_notes(
1798 &self,
1799 namespace: &str,
1800 kind: Option<&str>,
1801 page: PageRequest,
1802 ) -> Result<Page<Note>, StorageError> {
1803 let namespace = namespace.to_string();
1804 let kind = kind.map(|k| k.to_string());
1805 let limit_i64 = i64::from(page.limit);
1806 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1807 capability: StorageCapability::Notes,
1808 operation: "query_notes".into(),
1809 message: format!(
1810 "PageRequest: offset must be <= i64::MAX, got {}",
1811 page.offset
1812 ),
1813 })?;
1814
1815 self.with_reader("query_notes", move |conn| {
1816 let (count_sql, count_params) = build_note_where(&namespace, kind.as_deref());
1817 let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
1818
1819 let (where_sql, mut data_params) = build_note_where(&namespace, kind.as_deref());
1820 data_params.push(Box::new(limit_i64));
1821 data_params.push(Box::new(offset_i64));
1822
1823 let limit_idx = data_params.len() - 1;
1824 let offset_idx = data_params.len();
1825
1826 let data_sql = format!(
1827 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1828 properties, created_at, updated_at, deleted_at, key, version \
1829 FROM notes{} ORDER BY created_at DESC, id ASC LIMIT ?{} OFFSET ?{}",
1830 where_sql, limit_idx, offset_idx,
1831 );
1832
1833 query_note_page_snapshot(
1834 conn,
1835 "query_notes",
1836 &namespace,
1837 &count_sql,
1838 &count_params,
1839 &data_sql,
1840 &data_params,
1841 )
1842 })
1843 .await
1844 }
1845
1846 async fn query_notes_count_free(
1847 &self,
1848 namespace: &str,
1849 kind: Option<&str>,
1850 page: PageRequest,
1851 ) -> Result<Page<Note>, StorageError> {
1852 let namespace = namespace.to_string();
1853 let kind = kind.map(str::to_string);
1854 let limit_i64 = i64::from(page.limit);
1855 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1856 capability: StorageCapability::Notes,
1857 operation: "query_notes_count_free".into(),
1858 message: format!(
1859 "PageRequest: offset must be <= i64::MAX, got {}",
1860 page.offset
1861 ),
1862 })?;
1863
1864 self.with_reader("query_notes_count_free", move |conn| {
1865 let (where_sql, mut params) = build_note_where(&namespace, kind.as_deref());
1866 params.push(Box::new(limit_i64));
1867 params.push(Box::new(offset_i64));
1868 let limit_idx = params.len() - 1;
1869 let offset_idx = params.len();
1870 let sql = format!(
1871 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
1872 expires_at, properties, created_at, updated_at, deleted_at, key, version \
1873 FROM notes{where_sql} ORDER BY created_at DESC, id ASC \
1874 LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
1875 );
1876
1877 let mut stmt = conn.prepare(&sql)?;
1878 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1879 params.iter().map(|param| param.as_ref()).collect();
1880 let mut rows = stmt.query(param_refs.as_slice())?;
1881 let mut items = Vec::new();
1882 while let Some(row) = rows.next()? {
1883 items.push(read_note(row)?);
1884 }
1885
1886 Ok(Page { items, total: None })
1887 })
1888 .await
1889 }
1890
1891 async fn query_notes_filtered(
1892 &self,
1893 namespace: &str,
1894 filter: &NoteFilter,
1895 page: PageRequest,
1896 ) -> Result<Page<Note>, StorageError> {
1897 if filter.unordered {
1898 return Err(StorageError::InvalidInput {
1899 capability: StorageCapability::Notes,
1900 operation: "query_notes_filtered".into(),
1901 message: "NoteFilter.unordered is supported only by \
1902 query_notes_filtered_count_free"
1903 .into(),
1904 });
1905 }
1906 for pf in &filter.property_filters {
1908 validate_json_path(&pf.json_path)?;
1909 }
1910 if let Some((path, _)) = &filter.order_by {
1911 validate_json_path(path)?;
1912 }
1913 if filter.order_by_instant && !matches!(filter.order_by.as_ref(), Some((_, SortDir::Asc))) {
1914 return Err(StorageError::InvalidInput {
1915 capability: StorageCapability::Notes,
1916 operation: "query_notes_filtered".into(),
1917 message: "order_by_instant requires an ascending property order".into(),
1918 });
1919 }
1920 if filter.after.is_some() || filter.after_instant.is_some() {
1921 return Err(StorageError::InvalidInput {
1922 capability: StorageCapability::Notes,
1923 operation: "query_notes_filtered".into(),
1924 message: "NoteFilter.after or after_instant (keyset pagination) is not supported by this \
1925 method: it computes an exact COUNT(*) total over the whole \
1926 matching set, which has no defined meaning paired with a seek \
1927 boundary; use query_notes_filtered_count_free instead, which \
1928 seeks and returns total: None"
1929 .into(),
1930 });
1931 }
1932
1933 let namespace = namespace.to_string();
1934 let filter = filter.clone();
1935 let limit_i64 = i64::from(page.limit);
1936 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1937 capability: StorageCapability::Notes,
1938 operation: "query_notes_filtered".into(),
1939 message: format!(
1940 "PageRequest: offset must be <= i64::MAX, got {}",
1941 page.offset
1942 ),
1943 })?;
1944
1945 self.with_reader("query_notes_filtered", move |conn| {
1946 let (count_sql, count_params) = build_note_filter_read_clause(&namespace, &filter)?;
1947 let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
1948
1949 let (where_sql, mut data_params) = build_note_filter_read_clause(&namespace, &filter)?;
1950 data_params.push(Box::new(limit_i64));
1951 data_params.push(Box::new(offset_i64));
1952
1953 let order_clause = note_filter_page_order_clause(&filter);
1957
1958 let limit_idx = data_params.len() - 1;
1959 let offset_idx = data_params.len();
1960 let data_sql = format!(
1961 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
1962 expires_at, properties, created_at, updated_at, deleted_at, key, version \
1963 FROM notes{}{order_clause} LIMIT ?{} OFFSET ?{}",
1964 where_sql, limit_idx, offset_idx,
1965 );
1966
1967 query_note_page_snapshot(
1968 conn,
1969 "query_notes_filtered",
1970 &namespace,
1971 &count_sql,
1972 &count_params,
1973 &data_sql,
1974 &data_params,
1975 )
1976 })
1977 .await
1978 }
1979
1980 async fn query_notes_filtered_count_free(
1981 &self,
1982 namespace: &str,
1983 filter: &NoteFilter,
1984 page: PageRequest,
1985 ) -> Result<Page<Note>, StorageError> {
1986 for property_filter in &filter.property_filters {
1987 validate_json_path(&property_filter.json_path)?;
1988 }
1989 if let Some((path, _)) = &filter.order_by {
1990 validate_json_path(path)?;
1991 }
1992 if filter.order_by_instant && !matches!(filter.order_by.as_ref(), Some((_, SortDir::Asc))) {
1993 return Err(StorageError::InvalidInput {
1994 capability: StorageCapability::Notes,
1995 operation: "query_notes_filtered_count_free".into(),
1996 message: "order_by_instant requires an ascending property order".into(),
1997 });
1998 }
1999 if filter.order_by_instant && filter.unordered {
2000 return Err(StorageError::InvalidInput {
2001 capability: StorageCapability::Notes,
2002 operation: "query_notes_filtered_count_free".into(),
2003 message: "order_by_instant is incompatible with unordered pages".into(),
2004 });
2005 }
2006 if filter.after_instant.is_some() && !filter.order_by_instant {
2007 return Err(StorageError::InvalidInput {
2008 capability: StorageCapability::Notes,
2009 operation: "query_notes_filtered_count_free".into(),
2010 message: "after_instant requires order_by_instant".into(),
2011 });
2012 }
2013 if filter.after.is_some() && filter.after_instant.is_some() {
2014 return Err(StorageError::InvalidInput {
2015 capability: StorageCapability::Notes,
2016 operation: "query_notes_filtered_count_free".into(),
2017 message: "after and after_instant are mutually exclusive".into(),
2018 });
2019 }
2020 if filter.after.is_some() && filter.order_by.is_some() {
2021 return Err(StorageError::InvalidInput {
2022 capability: StorageCapability::Notes,
2023 operation: "query_notes_filtered_count_free".into(),
2024 message: "NoteFilter.after is incompatible with a custom order_by; it is \
2025 defined only over the default created_at DESC, id ASC order"
2026 .into(),
2027 });
2028 }
2029 if (filter.after.is_some() || filter.after_instant.is_some()) && page.offset != 0 {
2030 return Err(StorageError::InvalidInput {
2031 capability: StorageCapability::Notes,
2032 operation: "query_notes_filtered_count_free".into(),
2033 message: "NoteFilter.after or after_instant and a non-zero PageRequest.offset \
2034 are mutually exclusive pagination strategies; pass offset: 0"
2035 .into(),
2036 });
2037 }
2038
2039 let namespace = namespace.to_string();
2040 let filter = filter.clone();
2041 let limit_i64 = i64::from(page.limit);
2042 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
2043 capability: StorageCapability::Notes,
2044 operation: "query_notes_filtered_count_free".into(),
2045 message: format!(
2046 "PageRequest: offset must be <= i64::MAX, got {}",
2047 page.offset
2048 ),
2049 })?;
2050
2051 self.with_reader("query_notes_filtered_count_free", move |conn| {
2052 if let Some(after) = &filter.after_instant {
2053 let mut base_filter = filter.clone();
2054 base_filter.after_instant = None;
2055 let items =
2056 fetch_notes_after_instant(conn, &namespace, &base_filter, after, limit_i64)?;
2057 return Ok(Page { items, total: None });
2058 }
2059 if let Some(after) = &filter.after {
2060 let mut base_filter = filter.clone();
2061 base_filter.after = None;
2062 let items = fetch_notes_after(conn, &namespace, &base_filter, after, limit_i64)?;
2063 return Ok(Page { items, total: None });
2064 }
2065
2066 let (where_sql, mut params) = build_note_filter_read_clause(&namespace, &filter)?;
2067 params.push(Box::new(limit_i64));
2068 params.push(Box::new(offset_i64));
2069 let limit_idx = params.len() - 1;
2070 let offset_idx = params.len();
2071 let order_clause = note_filter_page_order_clause(&filter);
2072 let sql = format!(
2073 "SELECT {NOTE_COLUMNS} FROM notes{where_sql}{order_clause} \
2074 LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
2075 );
2076
2077 let mut stmt = conn.prepare(&sql)?;
2078 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2079 params.iter().map(|param| param.as_ref()).collect();
2080 let mut rows = stmt.query(param_refs.as_slice())?;
2081 let mut items = Vec::new();
2082 while let Some(row) = rows.next()? {
2083 items.push(read_note(row)?);
2084 #[cfg(test)]
2088 if items.len() == 1 {
2089 tests::page_snapshot_seam::hook("query_notes_filtered_count_free", &namespace);
2090 }
2091 }
2092
2093 Ok(Page { items, total: None })
2094 })
2095 .await
2096 }
2097
2098 async fn count_notes_filtered_in_snapshot(
2099 &self,
2100 namespace: &str,
2101 filters: &[NoteFilter],
2102 ) -> Result<Vec<u64>, StorageError> {
2103 for filter in filters {
2104 for pf in &filter.property_filters {
2105 validate_json_path(&pf.json_path)?;
2106 }
2107 }
2108
2109 let namespace = namespace.to_string();
2110 let filters = filters.to_vec();
2111 self.with_reader("count_notes_filtered_in_snapshot", move |conn| {
2112 let tx = rusqlite::Transaction::new_unchecked(
2113 conn,
2114 rusqlite::TransactionBehavior::Deferred,
2115 )?;
2116 let mut counts = Vec::with_capacity(filters.len());
2117 for filter in filters.iter() {
2118 #[cfg(test)]
2119 if !counts.is_empty() {
2120 tests::page_snapshot_seam::hook("count_notes_filtered_in_snapshot", &namespace);
2121 }
2122
2123 let (where_sql, params) = build_note_filter_read_clause(&namespace, filter)?;
2124 let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
2125 let mut stmt = tx.prepare(&sql)?;
2126 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2127 params.iter().map(|param| param.as_ref()).collect();
2128 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2129 counts.push(count as u64);
2130 }
2131 tx.commit()?;
2132 Ok(counts)
2133 })
2134 .await
2135 }
2136
2137 async fn count_notes_filtered_bounded_in_snapshot(
2138 &self,
2139 namespace: &str,
2140 filters: &[NoteFilter],
2141 cap: u32,
2142 ) -> Result<Vec<BoundedCount>, StorageError> {
2143 for filter in filters {
2144 for property_filter in &filter.property_filters {
2145 validate_json_path(&property_filter.json_path)?;
2146 }
2147 }
2148
2149 let namespace = namespace.to_string();
2150 let filters = filters.to_vec();
2151 let cap_u64 = u64::from(cap);
2152 let probe_limit_i64 = i64::from(cap) + 1;
2153 self.with_reader("count_notes_filtered_bounded_in_snapshot", move |conn| {
2154 let tx = rusqlite::Transaction::new_unchecked(
2155 conn,
2156 rusqlite::TransactionBehavior::Deferred,
2157 )?;
2158 let mut counts = Vec::with_capacity(filters.len());
2159 for filter in &filters {
2160 #[cfg(test)]
2161 if !counts.is_empty() {
2162 tests::page_snapshot_seam::hook(
2163 "count_notes_filtered_bounded_in_snapshot",
2164 &namespace,
2165 );
2166 }
2167
2168 let (where_sql, mut params) = build_note_filter_read_clause(&namespace, filter)?;
2169 params.push(Box::new(probe_limit_i64));
2170 let limit_idx = params.len();
2171 let sql = format!(
2176 "SELECT COUNT(*) FROM (SELECT 1 FROM notes{where_sql} LIMIT ?{limit_idx})"
2177 );
2178 let mut stmt = tx.prepare(&sql)?;
2179 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2180 params.iter().map(|param| param.as_ref()).collect();
2181 let observed: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2182 let observed = observed as u64;
2183 counts.push(BoundedCount {
2184 count: observed.min(cap_u64),
2185 cap: cap_u64,
2186 saturated: observed > cap_u64,
2187 });
2188 }
2189 tx.commit()?;
2190 Ok(counts)
2191 })
2192 .await
2193 }
2194
2195 async fn query_notes_filtered_after(
2196 &self,
2197 namespace: &str,
2198 filter: &NoteFilter,
2199 after: Option<SeekCursor>,
2200 limit: u32,
2201 ) -> Result<SeekPage<Note>, StorageError> {
2202 if limit == 0 {
2203 return Ok(SeekPage::default());
2204 }
2205 if filter.order_by.is_some() {
2206 return Err(StorageError::InvalidInput {
2207 capability: StorageCapability::Notes,
2208 operation: "query_notes_filtered_after".into(),
2209 message: "custom order_by is not compatible with insertion-sequence pagination"
2210 .into(),
2211 });
2212 }
2213 for property_filter in &filter.property_filters {
2214 validate_json_path(&property_filter.json_path)?;
2215 }
2216
2217 let namespace = namespace.to_string();
2218 let filter = filter.clone();
2219 let limit_usize = limit as usize;
2220 let probe_limit_i64 = i64::from(limit) + 1;
2221 self.with_reader("query_notes_filtered_after", move |conn| {
2222 let (mut where_sql, mut params) = build_note_filter_where(&namespace, &filter)?;
2223 if let Some(cursor) = after {
2224 params.push(Box::new(cursor.sequence));
2225 where_sql.push_str(&format!(" AND notes_seq.seq > ?{}", params.len()));
2226 }
2227 params.push(Box::new(probe_limit_i64));
2228 let limit_idx = params.len();
2229 let sql = format!(
2232 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
2233 expires_at, properties, created_at, updated_at, deleted_at, key, version, notes_seq.seq \
2234 FROM notes_seq CROSS JOIN notes ON notes.id = notes_seq.note_id{where_sql} \
2235 ORDER BY notes_seq.seq ASC LIMIT ?{limit_idx}"
2236 );
2237 let mut stmt = conn.prepare(&sql)?;
2238 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2239 params.iter().map(|param| param.as_ref()).collect();
2240 let rows = stmt.query_map(param_refs.as_slice(), |row| {
2241 Ok((read_note(row)?, row.get::<_, i64>(15)?))
2242 })?;
2243 let mut entries = rows.collect::<Result<Vec<_>, _>>()?;
2244 let has_more = entries.len() > limit_usize;
2245 if has_more {
2246 entries.truncate(limit_usize);
2247 }
2248 let next_after = if has_more {
2249 entries.last().map(|(note, sequence)| SeekCursor {
2250 sequence: *sequence,
2251 id: note.id,
2252 })
2253 } else {
2254 None
2255 };
2256 let items = entries.into_iter().map(|(note, _)| note).collect();
2257 Ok(SeekPage { items, next_after })
2258 })
2259 .await
2260 }
2261
2262 async fn query_notes_filtered_bounded(
2263 &self,
2264 namespace: &str,
2265 filter: &NoteFilter,
2266 max_rows: u32,
2267 ) -> Result<Vec<Note>, StorageError> {
2268 if filter.unordered {
2269 return Err(StorageError::InvalidInput {
2270 capability: StorageCapability::Notes,
2271 operation: "query_notes_filtered_bounded".into(),
2272 message: "NoteFilter.unordered is supported only by \
2273 query_notes_filtered_count_free"
2274 .into(),
2275 });
2276 }
2277 for pf in &filter.property_filters {
2278 validate_json_path(&pf.json_path)?;
2279 }
2280 if let Some((path, _)) = &filter.order_by {
2281 validate_json_path(path)?;
2282 }
2283
2284 let namespace = namespace.to_string();
2285 let filter = filter.clone();
2286 let limit_i64 = i64::from(max_rows) + 1;
2287
2288 self.with_reader("query_notes_filtered_bounded", move |conn| {
2289 let (where_sql, mut data_params) = build_note_filter_read_clause(&namespace, &filter)?;
2290 data_params.push(Box::new(limit_i64));
2291 let limit_idx = data_params.len();
2292
2293 let order_clause = match &filter.order_by {
2297 Some((path, dir)) => {
2298 let dir_str = match dir {
2299 SortDir::Asc => "ASC",
2300 SortDir::Desc => "DESC",
2301 };
2302 format!(" ORDER BY {} {dir_str}, id ASC", json_extract_expr(path))
2303 }
2304 None => " ORDER BY created_at DESC, id ASC".to_string(),
2305 };
2306
2307 let data_sql = format!(
2308 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
2309 expires_at, properties, created_at, updated_at, deleted_at, key, version \
2310 FROM notes{where_sql}{order_clause} LIMIT ?{limit_idx}",
2311 );
2312
2313 let mut stmt = conn.prepare(&data_sql)?;
2314 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2315 data_params.iter().map(|p| p.as_ref()).collect();
2316 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
2317
2318 let mut items = Vec::new();
2319 for row in rows {
2320 items.push(row?);
2321 }
2322 Ok(items)
2323 })
2324 .await
2325 }
2326
2327 async fn count_notes(&self, namespace: &str, kind: Option<&str>) -> Result<u64, StorageError> {
2328 let namespace = namespace.to_string();
2329 let kind = kind.map(|k| k.to_string());
2330
2331 self.with_reader("count_notes", move |conn| {
2332 let (where_sql, params) = build_note_where(&namespace, kind.as_deref());
2333 let sql = format!("SELECT COUNT(*) FROM notes{}", where_sql);
2334 let mut stmt = conn.prepare(&sql)?;
2335 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2336 params.iter().map(|p| p.as_ref()).collect();
2337 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2338 Ok(count as u64)
2339 })
2340 .await
2341 }
2342
2343 async fn count_notes_in_namespaces(
2344 &self,
2345 namespaces: &[String],
2346 kind: Option<&str>,
2347 ) -> Result<u64, StorageError> {
2348 let namespaces: Vec<String> = namespaces
2349 .iter()
2350 .cloned()
2351 .collect::<HashSet<_>>()
2352 .into_iter()
2353 .collect();
2354 let kind = kind.map(str::to_string);
2355
2356 self.with_reader("count_notes_in_namespaces", move |conn| {
2357 let mut total = 0;
2358 for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
2359 let (where_sql, params) = build_note_where_for_namespaces(chunk, kind.as_deref());
2360 let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
2361 let mut stmt = conn.prepare(&sql)?;
2362 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2363 params.iter().map(|p| p.as_ref()).collect();
2364 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2365 total += count as u64;
2366 }
2367 Ok(total)
2368 })
2369 .await
2370 }
2371}
2372
2373const NOTES_DDL: &str = include_str!("../../sql/notes-ddl.sql");
2378
2379const NOTES_SEQ_REPAIR_DDL: &str = include_str!("../../sql/008-notes-seq-repair.sql");
2386
2387pub(crate) fn ensure_notes_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
2388 conn.execute_batch(NOTES_DDL)
2389}
2390
2391pub(crate) fn repair_notes_seq(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
2397 conn.execute_batch(NOTES_SEQ_REPAIR_DDL)
2398}
2399
2400#[cfg(test)]
2401#[path = "note_tests.rs"]
2402mod tests;
2403
2404#[cfg(test)]
2405#[path = "comm_filter_plan_tests.rs"]
2406mod comm_filter_plan_tests;
2407
2408#[cfg(test)]
2409#[path = "note_list_plan_tests.rs"]
2410mod note_list_plan_tests;