1pub mod recipient;
4pub mod transport;
5
6use std::collections::HashSet;
7use std::sync::Arc;
8
9use async_trait::async_trait;
10use rusqlite::OptionalExtension;
11use uuid::Uuid;
12
13use khive_storage::attachment::{Attachment, AttachmentSubstrate};
14use khive_storage::error::StorageError;
15use khive_storage::note::{
16 FilterOp, Note, NoteFilter, NoteInstantSeekAfter, NoteKeyCursor, NoteSeekAfter, NoteTagMode,
17 NoteVisibility, SortDir,
18};
19use khive_storage::types::{
20 BatchWriteSummary, BoundedCount, DeleteMode, Page, PageRequest, SeekCursor, SeekPage,
21 SqlStatement, SqlValue,
22};
23use khive_storage::NoteStore;
24use khive_storage::{StorageCapability, StorageResult};
25
26use crate::error::SqliteError;
27use crate::pool::ConnectionPool;
28use crate::sql_bridge::bind_params;
29use crate::stores::attachment::{attachment_upsert_statement, delete_record_attachments_statement};
30use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
31
32fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
33 StorageError::driver(StorageCapability::Notes, op, e)
34}
35
36fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
37 e.into_storage_error(StorageCapability::Notes, op)
38}
39
40const NAMESPACE_COUNT_CHUNK_SIZE: usize = 500;
41
42pub fn note_key_prefix_successor(prefix: &str) -> Option<String> {
43 let mut chars: Vec<char> = prefix.chars().collect();
44 while let Some(last) = chars.pop() {
45 if last == char::MAX {
46 continue;
47 }
48 let next = if last == '\u{d7ff}' {
49 '\u{e000}'
50 } else {
51 char::from_u32(u32::from(last) + 1).expect("incremented non-max scalar")
52 };
53 chars.push(next);
54 return Some(chars.into_iter().collect());
55 }
56 None
57}
58
59pub const NOTE_UPSERT_SQL: &str = "INSERT INTO notes \
78 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
79 properties, created_at, updated_at, deleted_at, key, strict_due_key, due_source) \
80 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16) \
81 ON CONFLICT(id) DO UPDATE SET \
82 namespace = excluded.namespace, \
83 kind = excluded.kind, \
84 status = excluded.status, \
85 name = excluded.name, \
86 content = excluded.content, \
87 salience = excluded.salience, \
88 decay_factor = excluded.decay_factor, \
89 expires_at = excluded.expires_at, \
90 properties = excluded.properties, \
91 strict_due_key = excluded.strict_due_key, \
92 due_source = excluded.due_source, \
93 updated_at = excluded.updated_at, \
94 deleted_at = excluded.deleted_at";
95
96pub const NOTE_INSERT_IF_ABSENT_SQL: &str = "INSERT INTO notes \
105 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
106 properties, created_at, updated_at, deleted_at, key, strict_due_key, due_source) \
107 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16) \
108 ON CONFLICT(id) DO NOTHING";
109
110pub fn note_insert_if_absent_statement(note: &Note) -> SqlStatement {
112 let mut statement = note_upsert_statement(note);
113 statement.sql = NOTE_INSERT_IF_ABSENT_SQL.to_string();
114 statement.label = Some("note-insert-if-absent".to_string());
115 statement
116}
117
118pub fn note_insert_keyed_statement(note: &Note) -> SqlStatement {
120 let mut statement = note_upsert_statement(note);
121 statement.sql = "INSERT INTO notes \
122 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
123 properties, created_at, updated_at, deleted_at, key, strict_due_key, due_source) \
124 VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16) \
125 ON CONFLICT(namespace,kind,key) WHERE key IS NOT NULL AND deleted_at IS NULL DO NOTHING"
126 .into();
127 statement.label = Some("note-keyed-create".into());
128 statement
129}
130
131pub fn note_due_key_values(
135 properties: &Option<serde_json::Value>,
136) -> (Option<Vec<u8>>, Option<String>) {
137 let source = properties
138 .as_ref()
139 .and_then(|value| value.get("next_attempt_at"))
140 .and_then(|value| value.as_str());
141 match source.and_then(|text| crate::pool::strict_rfc3339_key(text).map(|key| (key, text))) {
142 Some((key, source)) => (Some(key), Some(source.to_string())),
143 None => (None, None),
144 }
145}
146
147pub fn note_upsert_statement(note: &Note) -> SqlStatement {
149 let (due_key, due_source) = note_due_key_values(¬e.properties);
150 let properties_str = note
151 .properties
152 .as_ref()
153 .map(|v| serde_json::to_string(v).unwrap_or_default());
154 SqlStatement {
155 sql: NOTE_UPSERT_SQL.to_string(),
156 params: vec![
157 SqlValue::Text(note.id.to_string()),
158 SqlValue::Text(note.namespace.clone()),
159 SqlValue::Text(note.kind.to_string()),
160 SqlValue::Text(note.status.clone()),
161 match ¬e.name {
162 Some(n) => SqlValue::Text(n.clone()),
163 None => SqlValue::Null,
164 },
165 SqlValue::Text(note.content.clone()),
166 match note.salience {
167 Some(s) => SqlValue::Float(s),
168 None => SqlValue::Null,
169 },
170 match note.decay_factor {
171 Some(d) => SqlValue::Float(d),
172 None => SqlValue::Null,
173 },
174 match note.expires_at {
175 Some(e) => SqlValue::Integer(e),
176 None => SqlValue::Null,
177 },
178 match properties_str {
179 Some(p) => SqlValue::Text(p),
180 None => SqlValue::Null,
181 },
182 SqlValue::Integer(note.created_at),
183 SqlValue::Integer(note.updated_at),
184 match note.deleted_at {
185 Some(d) => SqlValue::Integer(d),
186 None => SqlValue::Null,
187 },
188 match ¬e.key {
189 Some(key) => SqlValue::Text(key.clone()),
190 None => SqlValue::Null,
191 },
192 due_key.map_or(SqlValue::Null, SqlValue::Blob),
193 due_source.map_or(SqlValue::Null, SqlValue::Text),
194 ],
195 label: Some("note-upsert".to_string()),
196 }
197}
198
199pub fn note_replace_if_unchanged_statement(
206 note: &Note,
207 expected_updated_at: i64,
208 expected_deleted_at: Option<i64>,
209) -> SqlStatement {
210 let (due_key, due_source) = note_due_key_values(¬e.properties);
211 let properties_str = note
212 .properties
213 .as_ref()
214 .map(|v| serde_json::to_string(v).unwrap_or_default());
215 SqlStatement {
216 sql: "UPDATE notes SET \
217 namespace = ?1, kind = ?2, status = ?3, name = ?4, content = ?5, \
218 salience = ?6, decay_factor = ?7, expires_at = ?8, properties = ?9, \
219 updated_at = ?10, deleted_at = ?11, strict_due_key = ?15, due_source = ?16 \
220 WHERE id = ?12 AND updated_at = ?13 AND deleted_at IS ?14 \
221 AND ?10 > updated_at"
222 .to_string(),
223 params: vec![
224 SqlValue::Text(note.namespace.clone()),
225 SqlValue::Text(note.kind.to_string()),
226 SqlValue::Text(note.status.clone()),
227 match ¬e.name {
228 Some(name) => SqlValue::Text(name.clone()),
229 None => SqlValue::Null,
230 },
231 SqlValue::Text(note.content.clone()),
232 match note.salience {
233 Some(value) => SqlValue::Float(value),
234 None => SqlValue::Null,
235 },
236 match note.decay_factor {
237 Some(value) => SqlValue::Float(value),
238 None => SqlValue::Null,
239 },
240 match note.expires_at {
241 Some(value) => SqlValue::Integer(value),
242 None => SqlValue::Null,
243 },
244 match properties_str {
245 Some(value) => SqlValue::Text(value),
246 None => SqlValue::Null,
247 },
248 SqlValue::Integer(note.updated_at),
249 match note.deleted_at {
250 Some(value) => SqlValue::Integer(value),
251 None => SqlValue::Null,
252 },
253 SqlValue::Text(note.id.to_string()),
254 SqlValue::Integer(expected_updated_at),
255 match expected_deleted_at {
256 Some(value) => SqlValue::Integer(value),
257 None => SqlValue::Null,
258 },
259 due_key.map_or(SqlValue::Null, SqlValue::Blob),
260 due_source.map_or(SqlValue::Null, SqlValue::Text),
261 ],
262 label: Some("note-replace-if-unchanged".to_string()),
263 }
264}
265
266pub fn note_metadata_replace_if_unchanged_statement(
270 note: &Note,
271 expected_updated_at: i64,
272 expected_deleted_at: Option<i64>,
273) -> SqlStatement {
274 let mut statement =
275 note_replace_if_unchanged_statement(note, expected_updated_at, expected_deleted_at);
276 statement.sql = "UPDATE notes SET status=?3, name=?4, salience=?6, decay_factor=?7, expires_at=?8, updated_at=?10 \
277 WHERE id=?12 AND updated_at=?13 AND deleted_at IS ?14 AND ?10 > updated_at \
278 AND namespace=?1 AND kind=?2 AND content=?5 AND properties IS ?9 AND deleted_at IS ?11".into();
279 statement.params.truncate(14);
280 statement.label = Some("stream-note-metadata-cas".into());
281 statement
282}
283
284pub fn note_update_properties_statement(
292 id: Uuid,
293 properties: &Option<serde_json::Value>,
294 updated_at: i64,
295) -> SqlStatement {
296 let (due_key, due_source) = note_due_key_values(properties);
297 let properties_str = properties
298 .as_ref()
299 .map(|v| serde_json::to_string(v).unwrap_or_default());
300 SqlStatement {
301 sql: "UPDATE notes SET properties = ?1, updated_at = ?2, \
302 strict_due_key = ?4, due_source = ?5 \
303 WHERE id = ?3 AND deleted_at IS NULL"
304 .to_string(),
305 params: vec![
306 match properties_str {
307 Some(p) => SqlValue::Text(p),
308 None => SqlValue::Null,
309 },
310 SqlValue::Integer(updated_at),
311 SqlValue::Text(id.to_string()),
312 due_key.map_or(SqlValue::Null, SqlValue::Blob),
313 due_source.map_or(SqlValue::Null, SqlValue::Text),
314 ],
315 label: Some("note-update-properties".to_string()),
316 }
317}
318
319pub fn note_set_property_statement(
328 id: Uuid,
329 key: &str,
330 value: &serde_json::Value,
331 updated_at: i64,
332) -> Result<SqlStatement, StorageError> {
333 if key.contains('\0') {
337 return Err(StorageError::InvalidInput {
338 capability: StorageCapability::Notes,
339 operation: "set_note_property".into(),
340 message: "property key must not contain U+0000".to_string(),
341 });
342 }
343 let path = format!("$.{}", serde_json::Value::String(key.to_string()));
344 let updating_due = key == "next_attempt_at";
345 let (due_key, due_source) = if updating_due {
346 note_due_key_values(&Some(serde_json::json!({"next_attempt_at": value})))
347 } else {
348 (None, None)
349 };
350 Ok(SqlStatement {
351 sql: "UPDATE notes \
352 SET properties = json_set(COALESCE(properties, '{}'), ?1, json(?2)), \
353 updated_at = ?3, \
354 strict_due_key = CASE WHEN ?5 = 1 THEN ?6 ELSE strict_due_key END, \
355 due_source = CASE WHEN ?5 = 1 THEN ?7 ELSE due_source END \
356 WHERE id = ?4 AND deleted_at IS NULL \
357 AND (properties IS NULL OR json_type(properties) = 'object')"
358 .to_string(),
359 params: vec![
360 SqlValue::Text(path),
361 SqlValue::Text(value.to_string()),
362 SqlValue::Integer(updated_at),
363 SqlValue::Text(id.to_string()),
364 SqlValue::Integer(i64::from(updating_due)),
365 due_key.map_or(SqlValue::Null, SqlValue::Blob),
366 due_source.map_or(SqlValue::Null, SqlValue::Text),
367 ],
368 label: Some("note-set-property".to_string()),
369 })
370}
371
372pub fn note_soft_delete_statement(id: Uuid, deleted_at: i64) -> SqlStatement {
374 SqlStatement {
375 sql: "UPDATE notes SET status = 'deleted', deleted_at = ?1 \
376 WHERE id = ?2 AND deleted_at IS NULL"
377 .to_string(),
378 params: vec![
379 SqlValue::Integer(deleted_at),
380 SqlValue::Text(id.to_string()),
381 ],
382 label: Some("note-delete-soft".to_string()),
383 }
384}
385
386pub fn note_hard_delete_statement(id: Uuid) -> SqlStatement {
388 SqlStatement {
389 sql: "DELETE FROM notes WHERE id = ?1".to_string(),
390 params: vec![SqlValue::Text(id.to_string())],
391 label: Some("note-delete-hard".to_string()),
392 }
393}
394
395pub struct SqlNoteStore {
401 pool: Arc<ConnectionPool>,
402 index_repair: Option<super::index_repair::IndexRepairContext>,
403 writer_task: Option<WriterTaskHandle>,
404}
405
406impl SqlNoteStore {
407 pub fn new(pool: Arc<ConnectionPool>, _is_file_backed: bool) -> Self {
409 let writer_task = pool.writer_task_handle().ok().flatten();
416
417 Self {
418 pool,
419 writer_task,
420 index_repair: None,
421 }
422 }
423
424 pub(crate) fn with_index_repair(
425 mut self,
426 repair: super::index_repair::IndexRepairContext,
427 ) -> Self {
428 self.index_repair = Some(repair);
429 self
430 }
431
432 async fn with_indexed_reader<F, R>(&self, op: &'static str, read: F) -> Result<R, StorageError>
433 where
434 F: FnMut(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
435 R: Send + 'static,
436 {
437 super::index_repair::run_indexed_read(
438 Arc::clone(&self.pool),
439 self.index_repair.clone(),
440 StorageCapability::Notes,
441 op,
442 read,
443 )
444 .await
445 }
446
447 fn current_writer_task(
448 &self,
449 operation: &'static str,
450 ) -> Result<Option<WriterTaskHandle>, StorageError> {
451 self.pool
452 .writer_task_for_write(self.writer_task.as_ref(), operation)
453 }
454
455 async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
474 where
475 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
476 R: Send + 'static,
477 {
478 if let Some(writer_task) = self.current_writer_task(op)? {
479 return writer_task
480 .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
481 .await;
482 }
483
484 self.pool
485 .record_direct_route(crate::timeout_sink::Site::DirectRouteNote);
486 let pool = Arc::clone(&self.pool);
487 tokio::task::spawn_blocking(move || {
488 let guard = pool
489 .autocommit_write_unit()
490 .map_err(|e| map_sqlite_err(e, op))?;
491 f(guard.conn())
492 .map_err(|e| map_err(e, op))
493 .inspect_err(|error| pool.record_direct_writer_error(error))
494 })
495 .await
496 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
497 }
498
499 async fn with_writer_tx<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
511 where
512 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
513 R: Send + 'static,
514 {
515 self.with_writer_tx_storage(op, move |conn| f(conn).map_err(|error| map_err(error, op)))
516 .await
517 }
518
519 async fn with_writer_tx_storage<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
520 where
521 F: FnOnce(&rusqlite::Connection) -> Result<R, StorageError> + Send + 'static,
522 R: Send + 'static,
523 {
524 if let Some(writer_task) = self.current_writer_task(op)? {
525 return writer_task.send_bounded(f).await;
526 }
527
528 self.pool
529 .record_direct_route(crate::timeout_sink::Site::DirectRouteNote);
530 let pool = Arc::clone(&self.pool);
531 tokio::task::spawn_blocking(move || {
532 let guard = pool
533 .transaction_write_unit()
534 .map_err(|e| map_sqlite_err(e, op))
535 .inspect_err(|error| pool.record_direct_writer_error(error))?;
536 let conn = guard.conn();
537 let (result, terminal_state) = execute_wrapped_transaction(conn, op, f);
538 if terminal_state.is_some() {
539 pool.retire_pooled_writer(conn);
540 }
541 result.inspect_err(|error| pool.record_direct_writer_error(error))
542 })
543 .await
544 .map_err(|e| StorageError::driver(StorageCapability::Notes, op, e))?
545 }
546
547 async fn get_notes_batch_inner(
548 &self,
549 ids: &[Uuid],
550 include_deleted: bool,
551 ) -> Result<Vec<Note>, StorageError> {
552 if ids.is_empty() {
553 return Ok(vec![]);
554 }
555 let operation = if include_deleted {
556 "get_notes_batch_including_deleted"
557 } else {
558 "get_notes_batch"
559 };
560 let deletion_filter = if include_deleted {
561 ""
562 } else {
563 " AND deleted_at IS NULL"
564 };
565 const CHUNK: usize = 900;
568 let id_strings: Vec<String> = ids.iter().map(|id| id.to_string()).collect();
569
570 let mut result = Vec::with_capacity(ids.len());
571 for chunk in id_strings.chunks(CHUNK) {
572 let chunk_owned = chunk.to_vec();
573 let notes = self
574 .with_reader(operation, move |conn| {
575 let placeholders: String = (1..=chunk_owned.len())
576 .map(|i| format!("?{i}"))
577 .collect::<Vec<_>>()
578 .join(", ");
579 let sql = format!(
580 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
581 properties, created_at, updated_at, deleted_at, key, version \
582 FROM notes WHERE id IN ({placeholders}){deletion_filter}"
583 );
584 let mut stmt = conn.prepare(&sql)?;
585 let params: Vec<&dyn rusqlite::types::ToSql> = chunk_owned
586 .iter()
587 .map(|s| s as &dyn rusqlite::types::ToSql)
588 .collect();
589 let rows = stmt.query_map(params.as_slice(), read_note)?;
590 let mut notes = Vec::new();
591 for row in rows {
592 notes.push(row?);
593 }
594 Ok(notes)
595 })
596 .await?;
597 result.extend(notes);
598 }
599 Ok(result)
600 }
601
602 async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
603 where
604 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
605 R: Send + 'static,
606 {
607 super::run_pooled_store_read(
608 Arc::clone(&self.pool),
609 StorageCapability::Notes,
610 op,
611 move |conn| f(conn).map_err(|error| map_err(error, op)),
612 )
613 .await
614 }
615}
616
617fn read_note(row: &rusqlite::Row<'_>) -> Result<Note, rusqlite::Error> {
622 let id_str: String = row.get(0)?;
623 let namespace: String = row.get(1)?;
624 let kind: String = row.get(2)?;
625 let status: String = row.get(3)?;
626 let name: Option<String> = row.get(4)?;
627 let content: String = row.get(5)?;
628 let salience: Option<f64> = row.get(6)?;
629 let decay_factor: Option<f64> = row.get(7)?;
630 let expires_at: Option<i64> = row.get(8)?;
631 let properties_str: Option<String> = row.get(9)?;
632 let created_at: i64 = row.get(10)?;
633 let updated_at: i64 = row.get(11)?;
634 let deleted_at: Option<i64> = row.get(12)?;
635 let key: Option<String> = row.get(13)?;
636 let version: i64 = row.get(14)?;
637
638 let id = parse_uuid(&id_str)?;
639
640 let properties = properties_str
641 .map(|s| {
642 serde_json::from_str(&s).map_err(|e| {
643 rusqlite::Error::FromSqlConversionFailure(
644 9,
645 rusqlite::types::Type::Text,
646 Box::new(e),
647 )
648 })
649 })
650 .transpose()?;
651
652 Ok(Note {
653 id,
654 namespace,
655 kind,
656 status,
657 name,
658 content,
659 salience,
660 decay_factor,
661 expires_at,
662 properties,
663 created_at,
664 updated_at,
665 deleted_at,
666 key,
667 version,
668 })
669}
670
671fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
672 Uuid::parse_str(s).map_err(|e| {
673 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
674 })
675}
676
677fn query_note_page_snapshot(
678 conn: &rusqlite::Connection,
679 operation: &'static str,
680 namespace: &str,
681 count_sql: &str,
682 count_params: &[Box<dyn rusqlite::types::ToSql>],
683 data_sql: &str,
684 data_params: &[Box<dyn rusqlite::types::ToSql>],
685) -> Result<Page<Note>, rusqlite::Error> {
686 let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Deferred)?;
687
688 let total: i64 = {
689 let mut stmt = tx.prepare(count_sql)?;
690 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
691 count_params.iter().map(|param| param.as_ref()).collect();
692 stmt.query_row(param_refs.as_slice(), |row| row.get(0))?
693 };
694
695 #[cfg(test)]
696 tests::page_snapshot_seam::hook(operation, namespace);
697 #[cfg(not(test))]
698 let _ = (operation, namespace);
699
700 let items = {
701 let mut stmt = tx.prepare(data_sql)?;
702 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
703 data_params.iter().map(|param| param.as_ref()).collect();
704 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
705 rows.collect::<Result<Vec<_>, _>>()?
706 };
707
708 tx.commit()?;
709 Ok(Page {
710 items,
711 total: Some(total as u64),
712 })
713}
714
715fn batch_upsert_notes(
723 conn: &rusqlite::Connection,
724 notes: &[Note],
725 attempted: u64,
726) -> Result<BatchWriteSummary, rusqlite::Error> {
727 let mut summary = BatchWriteSummary {
728 attempted,
729 ..BatchWriteSummary::default()
730 };
731
732 let mut stmt = conn.prepare_cached(NOTE_UPSERT_SQL)?;
737
738 for (index, note) in notes.iter().enumerate() {
739 let id_str = note.id.to_string();
740 let kind_str = note.kind.to_string();
741 let status_str = note.status.clone();
742 let properties_str = note
743 .properties
744 .as_ref()
745 .map(|v| serde_json::to_string(v).unwrap_or_default());
746 let (due_key, due_source) = note_due_key_values(¬e.properties);
747
748 match stmt.execute(rusqlite::params![
749 id_str,
750 ¬e.namespace,
751 kind_str,
752 status_str,
753 ¬e.name,
754 note.content,
755 note.salience,
756 note.decay_factor,
757 note.expires_at,
758 properties_str,
759 note.created_at,
760 note.updated_at,
761 note.deleted_at,
762 note.key,
763 due_key,
764 due_source,
765 ]) {
766 Ok(_) => {
767 assign_note_seq(conn, &id_str)?;
768 summary.affected = summary.affected.saturating_add(1);
769 }
770 Err(e) => {
771 let (class, retryability) = super::classify_batch_sqlite_error(&e);
772 summary.record_failure(index, Some(id_str), class, retryability, e.to_string());
773 }
774 }
775 }
776
777 Ok(summary)
778}
779
780fn assign_note_seq(conn: &rusqlite::Connection, note_id: &str) -> Result<(), rusqlite::Error> {
786 conn.execute(
787 "INSERT OR IGNORE INTO notes_seq (note_id) VALUES (?1)",
788 rusqlite::params![note_id],
789 )?;
790 Ok(())
791}
792
793fn build_note_where(
794 namespace: &str,
795 kind: Option<&str>,
796) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
797 let mut conditions: Vec<String> = vec![
798 "namespace = ?1".to_string(),
799 "deleted_at IS NULL".to_string(),
800 ];
801 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![Box::new(namespace.to_string())];
802
803 if let Some(k) = kind {
804 params.push(Box::new(k.to_string()));
805 conditions.push(format!("kind = ?{}", params.len()));
806 }
807
808 let clause = format!(" WHERE {}", conditions.join(" AND "));
809 (clause, params)
810}
811
812fn build_note_where_for_namespaces(
813 namespaces: &[String],
814 kind: Option<&str>,
815) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
816 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = namespaces
817 .iter()
818 .map(|namespace| -> Box<dyn rusqlite::types::ToSql> { Box::new(namespace.clone()) })
819 .collect();
820 let namespace_condition = match namespaces.len() {
821 0 => "0".to_string(),
822 1 => "namespace = ?1".to_string(),
823 _ => {
824 let placeholders: Vec<String> =
825 (1..=namespaces.len()).map(|i| format!("?{i}")).collect();
826 format!("namespace IN ({})", placeholders.join(", "))
827 }
828 };
829 let mut conditions = vec![namespace_condition, "deleted_at IS NULL".to_string()];
830
831 if let Some(kind) = kind {
832 params.push(Box::new(kind.to_string()));
833 conditions.push(format!("kind = ?{}", params.len()));
834 }
835
836 let clause = format!(" WHERE {}", conditions.join(" AND "));
837 (clause, params)
838}
839
840fn validate_json_path(path: &str) -> Result<(), StorageError> {
843 let valid = path.starts_with("$.")
844 && path[2..].split('.').all(|part| {
845 !part.is_empty() && part.chars().all(|c| c.is_ascii_alphanumeric() || c == '_')
846 });
847 if valid {
848 Ok(())
849 } else {
850 Err(StorageError::InvalidInput {
851 capability: StorageCapability::Notes,
852 operation: "query_notes_filtered".into(),
853 message: format!("invalid JSON path for note filter: {path:?}"),
854 })
855 }
856}
857
858fn json_extract_expr(path: &str) -> String {
859 format!("json_extract(properties, '{path}')")
860}
861
862fn json_type_expr(path: &str) -> String {
863 format!("json_type(properties, '{path}')")
864}
865
866fn note_filter_page_order_clause(filter: &NoteFilter) -> String {
870 if filter.unordered {
871 return String::new();
872 }
873 match &filter.order_by {
874 Some((path, dir)) => {
875 if filter.order_by_instant {
876 let expr = json_extract_expr(path);
877 return format!(" ORDER BY khive_rfc3339_key({expr}) ASC, {expr} ASC, id ASC");
878 }
879 let dir_str = match dir {
880 SortDir::Asc => "ASC",
881 SortDir::Desc => "DESC",
882 };
883 format!(
886 " ORDER BY {} {dir_str}, id {dir_str}",
887 json_extract_expr(path)
888 )
889 }
890 None if filter
895 .property_filters
896 .iter()
897 .any(|property| matches!(&property.op, FilterOp::TextColonPrefixBucketIndexed))
898 && filter
899 .property_filters
900 .iter()
901 .any(|property| matches!(&property.op, FilterOp::Rfc3339LteOrInvalid)) =>
902 {
903 " ORDER BY +created_at DESC, id ASC".to_string()
904 }
905 None => " ORDER BY created_at DESC, id ASC".to_string(),
908 }
909}
910
911fn json_type_literal(value: &SqlValue) -> Result<&str, rusqlite::Error> {
917 const JSON_TYPES: [&str; 8] = [
918 "true", "false", "integer", "real", "text", "array", "object", "null",
919 ];
920 match value {
921 SqlValue::Text(s) if JSON_TYPES.contains(&s.as_str()) => Ok(s.as_str()),
922 other => Err(rusqlite::Error::InvalidParameterName(format!(
923 "json_type comparison value must be one of SQLite's json_type strings \
924 ({JSON_TYPES:?}), got {other:?}"
925 ))),
926 }
927}
928
929fn text_prefix_upper_bound(prefix: &str) -> Option<String> {
937 let mut chars: Vec<char> = prefix.chars().collect();
938 while let Some(last) = chars.pop() {
939 let mut next = u32::from(last) + 1;
940 if (0xD800..=0xDFFF).contains(&next) {
941 next = 0xE000;
942 }
943 if let Some(next) = char::from_u32(next) {
944 chars.push(next);
945 return Some(chars.into_iter().collect());
946 }
947 }
948 None
949}
950
951fn sql_value_param(value: &SqlValue) -> Result<Box<dyn rusqlite::types::ToSql>, rusqlite::Error> {
952 Ok(match value {
953 SqlValue::Null => Box::new(Option::<String>::None),
954 SqlValue::Bool(v) => Box::new(*v as i64),
955 SqlValue::Integer(v) => Box::new(*v),
956 SqlValue::Float(v) => Box::new(*v),
957 SqlValue::Text(v) => Box::new(v.clone()),
958 SqlValue::Blob(v) => Box::new(v.clone()),
959 SqlValue::Json(v) => Box::new(
960 serde_json::to_string(v)
961 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?,
962 ),
963 SqlValue::Uuid(v) => Box::new(v.to_string()),
964 SqlValue::Timestamp(v) => Box::new(v.timestamp_micros()),
965 })
966}
967
968fn build_note_filter_where(
969 namespace: &str,
970 filter: &NoteFilter,
971) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
972 let (ns_condition, ns_params): (String, Vec<Box<dyn rusqlite::types::ToSql>>) =
975 if !filter.namespaces.is_empty() {
976 let placeholders: Vec<String> = (1..=filter.namespaces.len())
977 .map(|i| format!("?{i}"))
978 .collect();
979 let params: Vec<Box<dyn rusqlite::types::ToSql>> = filter
980 .namespaces
981 .iter()
982 .map(|ns| -> Box<dyn rusqlite::types::ToSql> { Box::new(ns.clone()) })
983 .collect();
984 (
985 format!("namespace IN ({})", placeholders.join(", ")),
986 params,
987 )
988 } else {
989 (
990 "namespace = ?1".to_string(),
991 vec![Box::new(namespace.to_string())],
992 )
993 };
994
995 let mut conditions = vec![ns_condition, "deleted_at IS NULL".to_string()];
996 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = ns_params;
997
998 if let Some(kind) = &filter.kind {
999 params.push(Box::new(kind.clone()));
1000 conditions.push(format!("kind = ?{}", params.len()));
1001 }
1002
1003 if let Some(since) = filter.min_updated_at {
1004 params.push(Box::new(since));
1005 conditions.push(format!("updated_at >= ?{}", params.len()));
1006 }
1007 if !filter.tags.is_empty() {
1008 let mut tag_predicates = Vec::new();
1009 for tag in &filter.tags {
1010 params.push(Box::new(tag.clone()));
1011 tag_predicates.push(format!(
1012 "EXISTS (SELECT 1 FROM json_each(CASE WHEN json_type(properties,'$.tags')='array' \
1013 THEN json_extract(properties,'$.tags') ELSE '[]' END) AS tag \
1014 WHERE tag.type='text' AND tag.value = ?{} COLLATE NOCASE)",
1015 params.len()
1016 ));
1017 }
1018 let join = match filter.tag_mode {
1019 NoteTagMode::Any => " OR ",
1020 NoteTagMode::All => " AND ",
1021 };
1022 conditions.push(format!("({})", tag_predicates.join(join)));
1023 }
1024
1025 for pf in &filter.property_filters {
1026 match &pf.op {
1027 FilterOp::Rfc3339Valid => {
1028 let expr = json_extract_expr(&pf.json_path);
1029 conditions.push(format!("khive_rfc3339_key({expr}) IS NOT NULL"));
1030 }
1031 FilterOp::Rfc3339Gte | FilterOp::Rfc3339Lte | FilterOp::Rfc3339LteOrInvalid => {
1032 let instant = match &pf.value {
1033 SqlValue::Timestamp(instant) => *instant,
1034 SqlValue::Text(text) => {
1035 text.parse::<chrono::DateTime<chrono::Utc>>()
1036 .map_err(|error| {
1037 rusqlite::Error::ToSqlConversionFailure(Box::new(error))
1038 })?
1039 }
1040 _ => {
1041 return Err(rusqlite::Error::ToSqlConversionFailure(
1042 "RFC 3339 filters require a timestamp or text value".into(),
1043 ));
1044 }
1045 };
1046 let expr = json_extract_expr(&pf.json_path);
1047 let op = if matches!(&pf.op, FilterOp::Rfc3339Gte) {
1048 ">="
1049 } else {
1050 "<="
1051 };
1052 params.push(Box::new(crate::pool::rfc3339_instant_key(instant)));
1053 if matches!(&pf.op, FilterOp::Rfc3339LteOrInvalid) {
1054 conditions.push(format!(
1059 "CASE WHEN due_source = {expr} \
1060 THEN ifnull(strict_due_key, x'') ELSE x'' END <= ?{n} \
1061 AND ifnull(khive_rfc3339_strict_key({expr}), x'') <= ?{n}",
1062 n = params.len()
1063 ));
1064 } else {
1065 conditions.push(format!("khive_rfc3339_key({expr}) {op} ?{}", params.len()));
1066 }
1067 }
1068 FilterOp::EqOrMissing => {
1069 let expr = json_extract_expr(&pf.json_path);
1070 params.push(sql_value_param(&pf.value)?);
1071 conditions.push(format!(
1072 "({expr} = ?{n} OR {expr} IS NULL)",
1073 n = params.len()
1074 ));
1075 }
1076 FilterOp::EqOrMissingIndexed => {
1077 let expr = json_extract_expr(&pf.json_path);
1078 params.push(sql_value_param(&pf.value)?);
1079 conditions.push(format!("ifnull({expr}, '') = ?{}", params.len()));
1080 }
1081 FilterOp::TextEqOrNonText => {
1082 let expr = json_extract_expr(&pf.json_path);
1083 let type_expr = json_type_expr(&pf.json_path);
1084 params.push(sql_value_param(&pf.value)?);
1085 let n = params.len();
1086 conditions.push(format!(
1087 "CASE WHEN {type_expr} = 'text' THEN {expr} ELSE ?{n} END = ?{n}"
1088 ));
1089 }
1090 FilterOp::TextInOrNonText(values) => {
1091 let expr = json_extract_expr(&pf.json_path);
1092 let type_expr = json_type_expr(&pf.json_path);
1093 let mut placeholders = Vec::with_capacity(values.len());
1094 for value in values {
1095 params.push(sql_value_param(value)?);
1096 placeholders.push(format!("?{}", params.len()));
1097 }
1098 let text_match = if placeholders.is_empty() {
1099 "0".to_string()
1100 } else {
1101 format!("{expr} IN ({})", placeholders.join(", "))
1102 };
1103 conditions.push(format!(
1104 "CASE WHEN {type_expr} = 'text' THEN {text_match} ELSE 1 END"
1105 ));
1106 }
1107 FilterOp::JsonTypeEq => {
1108 let type_expr = json_type_expr(&pf.json_path);
1109 params.push(sql_value_param(&pf.value)?);
1110 conditions.push(format!("{type_expr} = ?{}", params.len()));
1111 }
1112 FilterOp::JsonTypeMissing => {
1113 let type_expr = json_type_expr(&pf.json_path);
1114 conditions.push(format!("{type_expr} IS NULL"));
1115 }
1116 FilterOp::JsonTypeMissingOrNullIndexed => {
1117 let expr = json_extract_expr(&pf.json_path);
1118 let type_expr = json_type_expr(&pf.json_path);
1119 conditions.push(format!(
1120 "ifnull({expr}, '') = '' AND ({type_expr} IS NULL OR {type_expr} = 'null')"
1121 ));
1122 }
1123 FilterOp::EqOrLegacyIndexed => {
1124 let expr = json_extract_expr(&pf.json_path);
1125 let type_expr = json_type_expr(&pf.json_path);
1126 params.push(sql_value_param(&pf.value)?);
1127 let n = params.len();
1128 conditions.push(format!(
1129 "ifnull({expr}, '') IN (?{n}, '') AND \
1130 ({type_expr} IS NULL OR {type_expr} = 'null' OR ifnull({expr}, '') != '')"
1131 ));
1132 }
1133 FilterOp::JsonTypeNeMissing => {
1134 let type_expr = json_type_expr(&pf.json_path);
1135 let literal = json_type_literal(&pf.value)?;
1145 conditions.push(format!(
1146 "({type_expr} IS NULL OR {type_expr} != '{literal}')"
1147 ));
1148 }
1149 FilterOp::In(values) => {
1150 let expr = json_extract_expr(&pf.json_path);
1151 if values.is_empty() {
1152 conditions.push("0".to_string());
1154 continue;
1155 }
1156 let mut placeholders = Vec::with_capacity(values.len());
1157 for v in values {
1158 params.push(sql_value_param(v)?);
1159 placeholders.push(format!("?{}", params.len()));
1160 }
1161 conditions.push(format!("{expr} IN ({})", placeholders.join(", ")));
1162 }
1163 FilterOp::TextStartsWithIndexed => {
1164 let expr = json_extract_expr(&pf.json_path);
1165 let SqlValue::Text(prefix) = &pf.value else {
1166 return Err(rusqlite::Error::ToSqlConversionFailure(
1167 "TextStartsWithIndexed takes a text prefix in PropertyFilter.value".into(),
1168 ));
1169 };
1170 params.push(Box::new(prefix.clone()));
1171 let lower = params.len();
1172 match text_prefix_upper_bound(prefix) {
1173 Some(upper) => {
1174 params.push(Box::new(upper));
1175 conditions.push(format!(
1176 "({expr} >= ?{lower} AND {expr} < ?{})",
1177 params.len()
1178 ));
1179 }
1180 None => {
1181 let type_expr = json_type_expr(&pf.json_path);
1182 conditions.push(format!("({type_expr} = 'text' AND {expr} >= ?{lower})"));
1183 }
1184 }
1185 }
1186 FilterOp::TextColonPrefixBucketIndexed => {
1187 let SqlValue::Text(prefix) = &pf.value else {
1188 return Err(rusqlite::Error::ToSqlConversionFailure(
1189 "TextColonPrefixBucketIndexed takes a text prefix".into(),
1190 ));
1191 };
1192 if prefix
1193 .strip_suffix(':')
1194 .is_none_or(|head| head.is_empty() || head.contains(':'))
1195 {
1196 return Err(rusqlite::Error::ToSqlConversionFailure(
1197 "TextColonPrefixBucketIndexed requires one trailing colon".into(),
1198 ));
1199 }
1200 let expr = json_extract_expr(&pf.json_path);
1201 params.push(Box::new(prefix.clone()));
1202 conditions.push(format!(
1203 "substr({expr}, 1, instr({expr}, ':')) = ?{}",
1204 params.len()
1205 ));
1206 }
1207 FilterOp::NotInOrMissing(values) => {
1208 let expr = json_extract_expr(&pf.json_path);
1209 if values.is_empty() {
1210 continue;
1212 }
1213 let mut placeholders = Vec::with_capacity(values.len());
1214 for v in values {
1215 params.push(sql_value_param(v)?);
1216 placeholders.push(format!("?{}", params.len()));
1217 }
1218 conditions.push(format!(
1219 "({expr} IS NULL OR {expr} NOT IN ({}))",
1220 placeholders.join(", ")
1221 ));
1222 }
1223 _ => {
1224 let expr = json_extract_expr(&pf.json_path);
1225 let op = match pf.op {
1226 FilterOp::Eq => "=",
1227 FilterOp::Ne => "!=",
1228 FilterOp::Lt => "<",
1229 FilterOp::Lte => "<=",
1230 FilterOp::Gt => ">",
1231 FilterOp::Gte => ">=",
1232 FilterOp::EqOrMissing
1233 | FilterOp::EqOrMissingIndexed
1234 | FilterOp::TextEqOrNonText
1235 | FilterOp::TextInOrNonText(_)
1236 | FilterOp::JsonTypeEq
1237 | FilterOp::JsonTypeMissing
1238 | FilterOp::JsonTypeMissingOrNullIndexed
1239 | FilterOp::EqOrLegacyIndexed
1240 | FilterOp::JsonTypeNeMissing
1241 | FilterOp::In(_)
1242 | FilterOp::NotInOrMissing(_)
1243 | FilterOp::TextStartsWithIndexed
1244 | FilterOp::TextColonPrefixBucketIndexed => {
1245 unreachable!()
1246 }
1247 FilterOp::Rfc3339Valid
1248 | FilterOp::Rfc3339Gte
1249 | FilterOp::Rfc3339Lte
1250 | FilterOp::Rfc3339LteOrInvalid => {
1251 unreachable!()
1252 }
1253 };
1254 params.push(sql_value_param(&pf.value)?);
1255 conditions.push(format!("{expr} {op} ?{}", params.len()));
1256 }
1257 }
1258 }
1259
1260 if let Some(min_ts) = filter.min_created_at {
1261 params.push(Box::new(min_ts));
1262 conditions.push(format!("created_at >= ?{}", params.len()));
1263 }
1264
1265 if let Some(scope) = &filter.mailbox {
1266 conditions.push(mailbox_condition(scope, &mut params));
1267 }
1268
1269 Ok((format!(" WHERE {}", conditions.join(" AND ")), params))
1270}
1271
1272fn mailbox_condition(
1278 scope: &khive_storage::note::NoteMailboxScope,
1279 params: &mut Vec<Box<dyn rusqlite::types::ToSql>>,
1280) -> String {
1281 params.push(Box::new(scope.actor_id.clone()));
1282 let actor = format!("?{}", params.len());
1283 let absent = |path: &str| {
1284 let ty = json_type_expr(path);
1285 format!("({ty} IS NULL OR {ty} = 'null')")
1286 };
1287 let not_text = |path: &str| format!("ifnull({}, '') != 'text'", json_type_expr(path));
1288 let direction = json_extract_expr("$.direction");
1289 let text_is_actor = |path: &str| {
1290 format!(
1291 "({} = 'text' AND {} = {actor})",
1292 json_type_expr(path),
1293 json_extract_expr(path)
1294 )
1295 };
1296 let to = text_is_actor("$.to_actor");
1297 let from = text_is_actor("$.from_actor");
1298 let (inbound_legacy, outbound_legacy, unrouted) = if scope.legacy_local {
1299 (
1300 format!(" OR {}", absent("$.to_actor")),
1301 format!(" OR {}", absent("$.from_actor")),
1302 format!(
1303 " OR ({} AND {} AND {})",
1304 absent("$.direction"),
1305 not_text("$.from_actor"),
1306 not_text("$.to_actor")
1307 ),
1308 )
1309 } else {
1310 Default::default()
1311 };
1312 format!(
1313 "(kind != 'message' OR ({direction} = 'inbound' AND ({to}{inbound_legacy})) \
1314 OR ({direction} = 'outbound' AND ({from}{outbound_legacy})){unrouted})"
1315 )
1316}
1317
1318fn comm_filter_index_clause(filter: &NoteFilter, where_sql: &str) -> &'static str {
1322 if filter.kind.as_deref() != Some("message") {
1323 return "";
1324 }
1325 let Some(predicate) = where_sql.strip_prefix(" WHERE ") else {
1326 return "";
1327 };
1328 let terms: Vec<_> = predicate.split(" AND ").collect();
1329 let numbered_param = |value: &str| {
1330 value.strip_prefix('?').is_some_and(|number| {
1331 !number.is_empty() && number.bytes().all(|byte| byte.is_ascii_digit())
1332 })
1333 };
1334 let equality = |prefix: &str| {
1335 terms
1336 .iter()
1337 .any(|term| term.strip_prefix(prefix).is_some_and(numbered_param))
1338 };
1339 let namespace = equality("namespace = ")
1340 || terms.iter().any(|term| {
1341 term.strip_prefix("namespace IN (")
1342 .and_then(|term| term.strip_suffix(')'))
1343 .is_some_and(|values| values.split(", ").all(numbered_param))
1344 });
1345 let recipient = equality("ifnull(json_extract(properties, '$.to_actor'), '') = ")
1346 || terms.contains(&"ifnull(json_extract(properties, '$.to_actor'), '') = ''")
1347 || terms.iter().any(|term| {
1348 term.strip_prefix("ifnull(json_extract(properties, '$.to_actor'), '') IN (")
1349 .and_then(|term| term.strip_suffix(", '')"))
1350 .is_some_and(numbered_param)
1351 });
1352 if !namespace
1353 || !terms.contains(&"deleted_at IS NULL")
1354 || !equality("kind = ")
1355 || !equality("json_extract(properties, '$.direction') = ")
1356 || !recipient
1357 {
1358 return "";
1359 }
1360 if terms.contains(
1361 &"(json_type(properties, '$.read') IS NULL OR json_type(properties, '$.read') != 'true')",
1362 ) {
1363 let typed_exact_recipient =
1367 equality("ifnull(json_extract(properties, '$.to_actor'), '') = ")
1368 && equality("json_type(properties, '$.to_actor') = ")
1369 && filter.property_filters.iter().any(|property| {
1370 property.json_path == "$.to_actor"
1371 && matches!(property.op, FilterOp::JsonTypeEq)
1372 && matches!(&property.value, SqlValue::Text(value) if value == "text")
1373 });
1374 if typed_exact_recipient {
1375 " INDEXED BY idx_notes_unread_probe_recipient_type_direction"
1376 } else {
1377 " INDEXED BY idx_notes_unread_probe_recipient_direction"
1378 }
1379 } else {
1380 " INDEXED BY idx_notes_message_recipient_direction"
1381 }
1382}
1383
1384fn build_note_filter_read_clause(
1385 namespace: &str,
1386 filter: &NoteFilter,
1387) -> Result<(String, Vec<Box<dyn rusqlite::types::ToSql>>), rusqlite::Error> {
1388 let (where_sql, params) = build_note_filter_where(namespace, filter)?;
1389 let index_clause = comm_filter_index_clause(filter, &where_sql);
1390 Ok((format!("{index_clause}{where_sql}"), params))
1391}
1392
1393const NOTE_COLUMNS: &str = "id, namespace, kind, status, name, content, salience, decay_factor, \
1400 expires_at, properties, created_at, updated_at, deleted_at, key, version";
1401
1402fn fetch_notes_after_instant(
1403 conn: &rusqlite::Connection,
1404 namespace: &str,
1405 base_filter: &NoteFilter,
1406 after: &NoteInstantSeekAfter,
1407 limit: i64,
1408) -> Result<Vec<Note>, rusqlite::Error> {
1409 let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
1410 params.push(Box::new(after.value.clone()));
1411 let value_idx = params.len();
1412 params.push(Box::new(after.id.to_string()));
1413 let id_idx = params.len();
1414 params.push(Box::new(limit));
1415 let limit_idx = params.len();
1416 let (path, _) = base_filter
1417 .order_by
1418 .as_ref()
1419 .expect("instant cursor requires order_by");
1420 let expr = json_extract_expr(path);
1421 let order_clause = note_filter_page_order_clause(base_filter);
1422 let sql = format!(
1423 "SELECT {NOTE_COLUMNS} FROM notes{where_sql} \
1424 AND (khive_rfc3339_key({expr}), {expr}, id) > \
1425 (khive_rfc3339_key(?{value_idx}), ?{value_idx}, ?{id_idx}) \
1426 {order_clause} LIMIT ?{limit_idx}"
1427 );
1428 let mut stmt = conn.prepare_cached(&sql)?;
1429 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1430 params.iter().map(|param| param.as_ref()).collect();
1431 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1432 rows.collect()
1433}
1434
1435fn fetch_notes_after(
1454 conn: &rusqlite::Connection,
1455 namespace: &str,
1456 base_filter: &NoteFilter,
1457 after: &NoteSeekAfter,
1458 limit: i64,
1459) -> Result<Vec<Note>, rusqlite::Error> {
1460 let mut items = Vec::new();
1461 if limit <= 0 {
1462 return Ok(items);
1463 }
1464
1465 {
1466 let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
1467 params.push(Box::new(after.created_at));
1468 let ts_idx = params.len();
1469 params.push(Box::new(after.id.to_string()));
1470 let id_idx = params.len();
1471 params.push(Box::new(limit));
1472 let limit_idx = params.len();
1473 let sql = format!(
1474 "SELECT {NOTE_COLUMNS} FROM notes{where_sql} AND created_at = ?{ts_idx} \
1475 AND id > ?{id_idx} ORDER BY id ASC LIMIT ?{limit_idx}"
1476 );
1477 let mut stmt = conn.prepare_cached(&sql)?;
1478 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1479 params.iter().map(|p| p.as_ref()).collect();
1480 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1481 for row in rows {
1482 items.push(row?);
1483 }
1484 }
1485
1486 let remaining = limit - items.len() as i64;
1487 if remaining > 0 {
1488 let (where_sql, mut params) = build_note_filter_read_clause(namespace, base_filter)?;
1489 params.push(Box::new(after.created_at));
1490 let ts_idx = params.len();
1491 params.push(Box::new(remaining));
1492 let limit_idx = params.len();
1493 let sql = format!(
1494 "SELECT {NOTE_COLUMNS} FROM notes{where_sql} AND created_at < ?{ts_idx} \
1495 ORDER BY created_at DESC, id ASC LIMIT ?{limit_idx}"
1496 );
1497 let mut stmt = conn.prepare_cached(&sql)?;
1498 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1499 params.iter().map(|p| p.as_ref()).collect();
1500 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
1501 for row in rows {
1502 items.push(row?);
1503 }
1504 }
1505
1506 Ok(items)
1507}
1508
1509fn execute_filtered_note_property_patch(
1510 conn: &rusqlite::Connection,
1511 id: Uuid,
1512 namespace: &str,
1513 filter: &NoteFilter,
1514 json_path: &str,
1515 value_json: &str,
1516 updated_at: i64,
1517) -> Result<usize, rusqlite::Error> {
1518 let (where_clause, mut params) = build_note_filter_where(namespace, filter)?;
1519 let updating_due = json_path == "$.next_attempt_at";
1520 let due_source = if updating_due {
1521 serde_json::from_str::<serde_json::Value>(value_json)
1522 .ok()
1523 .and_then(|value| value.as_str().map(str::to_owned))
1524 } else {
1525 None
1526 };
1527 let due_key = due_source
1528 .as_deref()
1529 .and_then(crate::pool::strict_rfc3339_key);
1530
1531 let base = params.len();
1532 let sql = format!(
1533 "UPDATE notes SET properties = json_set(COALESCE(properties, '{{}}'), ?{p1}, json(?{p2})), \
1534 updated_at = ?{p3}, \
1535 strict_due_key = CASE WHEN ?{p5} = 1 THEN ?{p6} ELSE strict_due_key END, \
1536 due_source = CASE WHEN ?{p5} = 1 THEN ?{p7} ELSE due_source END \
1537 {where_clause} \
1538 AND (properties IS NULL OR json_type(properties) = 'object') AND id = ?{p4}",
1539 p1 = base + 1,
1540 p2 = base + 2,
1541 p3 = base + 3,
1542 p4 = base + 4,
1543 p5 = base + 5,
1544 p6 = base + 6,
1545 p7 = base + 7,
1546 );
1547 params.push(Box::new(json_path.to_string()));
1548 params.push(Box::new(value_json.to_string()));
1549 params.push(Box::new(updated_at));
1550 params.push(Box::new(id.to_string()));
1551 params.push(Box::new(i64::from(updating_due)));
1552 params.push(Box::new(due_key));
1553 params.push(Box::new(due_source));
1554
1555 let mut stmt = conn.prepare_cached(&sql)?;
1556 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
1557 params.iter().map(|param| param.as_ref()).collect();
1558 stmt.execute(param_refs.as_slice())
1559}
1560
1561#[async_trait]
1566impl NoteStore for SqlNoteStore {
1567 async fn get_live_notes_by_key(
1568 &self,
1569 namespace: &str,
1570 key: &str,
1571 kind: Option<&str>,
1572 ) -> StorageResult<Vec<Note>> {
1573 let namespace = namespace.to_owned();
1574 let key = key.to_owned();
1575 let kind = kind.map(str::to_owned);
1576 self.with_reader("get_live_notes_by_key", move |conn| {
1577 let (mut clause, mut params) = build_note_where(&namespace, kind.as_deref());
1578 params.push(Box::new(key));
1579 clause.push_str(&format!(" AND key = ?{}", params.len()));
1580 let sql = format!("SELECT {NOTE_COLUMNS} FROM notes{clause} ORDER BY kind ASC");
1581 let mut stmt = conn.prepare(&sql)?;
1582 let params: Vec<&dyn rusqlite::types::ToSql> =
1583 params.iter().map(|p| p.as_ref()).collect();
1584 let rows = stmt.query_map(params.as_slice(), read_note)?.collect();
1585 rows
1586 })
1587 .await
1588 }
1589
1590 async fn query_keyed_notes(
1591 &self,
1592 namespace: &str,
1593 filter: &NoteFilter,
1594 prefix: &str,
1595 after: Option<&NoteKeyCursor>,
1596 page: PageRequest,
1597 ) -> StorageResult<(Vec<Note>, Option<NoteKeyCursor>)> {
1598 if !filter.namespaces.is_empty()
1599 || filter.order_by.is_some()
1600 || filter.after.is_some()
1601 || (after.is_some() && page.offset != 0)
1602 || page.limit == 0
1603 {
1604 return Err(StorageError::InvalidInput { capability: StorageCapability::Notes,
1605 operation: "query_keyed_notes".into(), message: "keyed paging requires primary namespace, keyed order and a positive limit; cursor excludes offset".into() });
1606 }
1607 for property in &filter.property_filters {
1608 validate_json_path(&property.json_path)?;
1609 }
1610 let offset = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
1611 capability: StorageCapability::Notes,
1612 operation: "query_keyed_notes".into(),
1613 message: "offset exceeds the supported integer range".into(),
1614 })?;
1615 let namespace = namespace.to_owned();
1616 let filter = filter.clone();
1617 let prefix = prefix.to_owned();
1618 let after = after.cloned();
1619 self.with_reader("query_keyed_notes", move |conn| {
1620 let (mut clause, mut params) = build_note_filter_where(&namespace, &filter)?;
1621 params.push(Box::new(prefix.clone()));
1622 clause.push_str(&format!(" AND key IS NOT NULL AND key >= ?{}", params.len()));
1623 if let Some(upper) = note_key_prefix_successor(&prefix) {
1624 params.push(Box::new(upper));
1625 clause.push_str(&format!(" AND key < ?{}", params.len()));
1626 }
1627 if let Some(after) = after {
1628 params.push(Box::new(after.updated_at)); let u = params.len();
1629 params.push(Box::new(after.key)); let k = params.len();
1630 params.push(Box::new(after.id.to_string())); let id = params.len();
1631 clause.push_str(&format!(" AND (updated_at < ?{u} OR (updated_at = ?{u} AND key < ?{k}) \
1632 OR (updated_at = ?{u} AND key = ?{k} AND id > ?{id}))"));
1633 }
1634 params.push(Box::new(i64::from(page.limit) + 1)); let limit = params.len();
1635 params.push(Box::new(offset)); let offset = params.len();
1636 let sql = format!("SELECT {NOTE_COLUMNS} FROM notes{clause} ORDER BY updated_at DESC, key DESC, id ASC LIMIT ?{limit} OFFSET ?{offset}");
1637 let mut stmt = conn.prepare(&sql)?;
1638 let params: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
1639 let mut notes = stmt.query_map(params.as_slice(), read_note)?.collect::<Result<Vec<_>, _>>()?;
1640 let has_more = notes.len() > page.limit as usize;
1641 notes.truncate(page.limit as usize);
1642 let next = if has_more { notes.last().map(NoteKeyCursor::from) } else { None };
1643 Ok((notes, next))
1644 }).await
1645 }
1646
1647 async fn upsert_note(&self, note: Note) -> Result<(), StorageError> {
1648 let id_str = note.id.to_string();
1649 let statement = note_upsert_statement(¬e);
1650 self.with_writer_tx("upsert_note", move |conn| {
1651 let mut stmt = conn.prepare_cached(&statement.sql)?;
1652 bind_params(&mut stmt, &statement.params)?;
1653 stmt.raw_execute()?;
1654 assign_note_seq(conn, &id_str)?;
1655 Ok(())
1656 })
1657 .await
1658 }
1659
1660 async fn insert_note_if_absent(&self, note: Note) -> Result<bool, StorageError> {
1661 let id_str = note.id.to_string();
1662 let statement = note_insert_if_absent_statement(¬e);
1663 self.with_writer_tx("insert_note_if_absent", move |conn| {
1664 let mut stmt = conn.prepare_cached(&statement.sql)?;
1665 bind_params(&mut stmt, &statement.params)?;
1666 let inserted = stmt.raw_execute()? > 0;
1667 if inserted {
1671 assign_note_seq(conn, &id_str)?;
1672 }
1673 Ok(inserted)
1674 })
1675 .await
1676 }
1677
1678 async fn replace_note_if_unchanged(
1679 &self,
1680 note: Note,
1681 expected_updated_at: i64,
1682 expected_deleted_at: Option<i64>,
1683 ) -> Result<bool, StorageError> {
1684 let statement =
1685 note_replace_if_unchanged_statement(¬e, expected_updated_at, expected_deleted_at);
1686 self.with_writer("replace_note_if_unchanged", move |conn| {
1687 let mut stmt = conn.prepare(&statement.sql)?;
1688 bind_params(&mut stmt, &statement.params)?;
1689 Ok(stmt.raw_execute()? > 0)
1690 })
1691 .await
1692 }
1693
1694 async fn update_note_properties(
1695 &self,
1696 id: Uuid,
1697 properties: Option<serde_json::Value>,
1698 updated_at: i64,
1699 ) -> Result<bool, StorageError> {
1700 let statement = note_update_properties_statement(id, &properties, updated_at);
1701 self.with_writer("update_note_properties", move |conn| {
1702 let mut stmt = conn.prepare(&statement.sql)?;
1703 bind_params(&mut stmt, &statement.params)?;
1704 Ok(stmt.raw_execute()? > 0)
1705 })
1706 .await
1707 }
1708
1709 async fn set_note_property(
1710 &self,
1711 id: Uuid,
1712 key: &str,
1713 value: serde_json::Value,
1714 updated_at: i64,
1715 ) -> Result<bool, StorageError> {
1716 let statement = note_set_property_statement(id, key, &value, updated_at)?;
1717 self.with_writer("set_note_property", move |conn| {
1718 let mut stmt = conn.prepare(&statement.sql)?;
1719 bind_params(&mut stmt, &statement.params)?;
1720 Ok(stmt.raw_execute()? > 0)
1721 })
1722 .await
1723 }
1724
1725 async fn try_patch_note_property(
1726 &self,
1727 id: Uuid,
1728 namespace: &str,
1729 filter: &NoteFilter,
1730 json_path: &str,
1731 value: serde_json::Value,
1732 updated_at: i64,
1733 ) -> Result<bool, StorageError> {
1734 let namespace = namespace.to_string();
1735 let filter = filter.clone();
1736 let value_json = serde_json::to_string(&value).map_err(|e| {
1737 StorageError::driver(StorageCapability::Notes, "try_patch_note_property", e)
1738 })?;
1739 let json_path = json_path.to_string();
1740
1741 self.with_writer("try_patch_note_property", move |conn| {
1742 execute_filtered_note_property_patch(
1743 conn,
1744 id,
1745 &namespace,
1746 &filter,
1747 &json_path,
1748 &value_json,
1749 updated_at,
1750 )
1751 .map(|rows| rows > 0)
1752 })
1753 .await
1754 }
1755
1756 async fn patch_note_property_atomic(
1757 &self,
1758 mut ids: Vec<Uuid>,
1759 namespace: &str,
1760 filter: &NoteFilter,
1761 json_path: &str,
1762 value: serde_json::Value,
1763 updated_at: i64,
1764 ) -> Result<(), StorageError> {
1765 let mut seen = HashSet::with_capacity(ids.len());
1766 ids.retain(|id| seen.insert(*id));
1767 if ids.is_empty() {
1768 return Err(StorageError::InvalidInput {
1769 capability: StorageCapability::Notes,
1770 operation: "patch_note_property_atomic".into(),
1771 message: "at least one note id is required".to_string(),
1772 });
1773 }
1774
1775 let namespace = namespace.to_string();
1776 let filter = filter.clone();
1777 let value_json = serde_json::to_string(&value).map_err(|e| {
1778 StorageError::driver(StorageCapability::Notes, "patch_note_property_atomic", e)
1779 })?;
1780 let json_path = json_path.to_string();
1781
1782 self.with_writer_tx_storage("patch_note_property_atomic", move |conn| {
1783 for id in ids {
1784 let rows = execute_filtered_note_property_patch(
1785 conn,
1786 id,
1787 &namespace,
1788 &filter,
1789 &json_path,
1790 &value_json,
1791 updated_at,
1792 )
1793 .map_err(|error| map_err(error, "patch_note_property_atomic"))?;
1794 if rows != 1 {
1795 return Err(StorageError::Conflict {
1796 capability: StorageCapability::Notes,
1797 operation: "patch_note_property_atomic".into(),
1798 message: format!(
1799 "precondition failed for note {id}: guarded update changed {rows} rows; expected 1"
1800 ),
1801 });
1802 }
1803 }
1804 Ok(())
1805 })
1806 .await
1807 }
1808
1809 async fn try_insert_note(&self, note: Note) -> Result<bool, StorageError> {
1810 self.try_insert_note_with_attachments(note, Vec::new())
1811 .await
1812 }
1813
1814 async fn try_insert_note_with_attachments(
1815 &self,
1816 note: Note,
1817 attachments: Vec<Attachment>,
1818 ) -> Result<bool, StorageError> {
1819 let mut attachment_statements = Vec::with_capacity(attachments.len());
1820 for attachment in attachments {
1821 attachment.validate()?;
1822 if attachment.record_uuid != note.id
1823 || attachment.substrate != AttachmentSubstrate::Note
1824 {
1825 return Err(StorageError::InvalidInput {
1826 capability: StorageCapability::Attachments,
1827 operation: "try_insert_note_with_attachments".into(),
1828 message: format!(
1829 "attachment {} must target note {}",
1830 attachment.role, note.id
1831 ),
1832 });
1833 }
1834 attachment_statements.push(attachment_upsert_statement(&attachment)?);
1835 }
1836 let namespace = note.namespace.clone();
1837 let id_str = note.id.to_string();
1838 let kind_str = note.kind.to_string();
1839 let status_str = note.status.clone();
1840 let properties_str = note
1841 .properties
1842 .as_ref()
1843 .map(|v| serde_json::to_string(v).unwrap_or_default());
1844 let (due_key, due_source) = note_due_key_values(¬e.properties);
1845
1846 let ext_id_opt: Option<String> = note
1848 .properties
1849 .as_ref()
1850 .and_then(|v| v.get("external_id"))
1851 .and_then(|v| v.as_str())
1852 .filter(|s| !s.is_empty())
1853 .map(|s| s.to_string());
1854 let channel_kind = note
1855 .properties
1856 .as_ref()
1857 .and_then(|v| v.get("channel_kind"))
1858 .and_then(|v| v.as_str())
1859 .unwrap_or_default()
1860 .to_string();
1861 let channel_slug = note
1862 .properties
1863 .as_ref()
1864 .and_then(|v| v.get("channel_slug"))
1865 .and_then(|v| v.as_str())
1866 .unwrap_or_default()
1867 .to_string();
1868
1869 self.with_writer_tx("try_insert_note", move |conn| {
1870 let rows = conn.execute(
1871 "INSERT OR IGNORE INTO notes \
1872 (id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1873 properties, created_at, updated_at, deleted_at, key, strict_due_key, due_source) \
1874 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)",
1875 rusqlite::params![
1876 id_str,
1877 namespace,
1878 kind_str,
1879 status_str,
1880 note.name,
1881 note.content,
1882 note.salience,
1883 note.decay_factor,
1884 note.expires_at,
1885 properties_str,
1886 note.created_at,
1887 note.updated_at,
1888 note.deleted_at,
1889 note.key,
1890 due_key,
1891 due_source,
1892 ],
1893 )?;
1894
1895 if rows > 0 {
1896 assign_note_seq(conn, &id_str)?;
1897 for statement in attachment_statements {
1898 let mut stmt = conn.prepare(&statement.sql)?;
1899 bind_params(&mut stmt, &statement.params)?;
1900 stmt.raw_execute()?;
1901 }
1902 return Ok(true);
1903 }
1904
1905 if let Some(ref ext_id) = ext_id_opt {
1912 let is_dedup: bool = conn.query_row(
1913 "SELECT COUNT(*) > 0 FROM notes \
1914 WHERE namespace = ?1 \
1915 AND kind = ?2 \
1916 AND json_extract(properties, '$.external_id') = ?3 \
1917 AND ifnull(json_extract(properties, '$.channel_kind'), '') = ?4 \
1918 AND ifnull(json_extract(properties, '$.channel_slug'), '') = ?5 \
1919 AND deleted_at IS NULL",
1920 rusqlite::params![namespace, kind_str, ext_id, channel_kind, channel_slug],
1921 |row| row.get(0),
1922 )?;
1923 if is_dedup {
1924 return Ok(false);
1925 }
1926 }
1927
1928 Err(rusqlite::Error::SqliteFailure(
1931 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
1932 Some(
1933 "try_insert_note: INSERT ignored for a constraint other than \
1934 external_id dedup; not masking as deduplication"
1935 .to_string(),
1936 ),
1937 ))
1938 })
1939 .await
1940 }
1941
1942 async fn upsert_notes(&self, notes: Vec<Note>) -> Result<BatchWriteSummary, StorageError> {
1943 let attempted = notes.len() as u64;
1944
1945 let origin = self.pool.origin();
1956 self.with_writer_tx("upsert_notes", move |conn| {
1957 let _tx_handle = khive_storage::tx_registry::register_scoped(
1958 Some("note_upsert_batch".to_string()),
1959 origin,
1960 );
1961 batch_upsert_notes(conn, ¬es, attempted)
1962 })
1963 .await
1964 }
1965
1966 async fn get_note(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
1967 let id_str = id.to_string();
1968
1969 self.with_reader("get_note", move |conn| {
1970 let mut stmt = conn.prepare(
1971 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1972 properties, created_at, updated_at, deleted_at, key, version \
1973 FROM notes WHERE id = ?1 AND deleted_at IS NULL",
1974 )?;
1975 let mut rows = stmt.query(rusqlite::params![id_str])?;
1976 match rows.next()? {
1977 Some(row) => Ok(Some(read_note(row)?)),
1978 None => Ok(None),
1979 }
1980 })
1981 .await
1982 }
1983
1984 async fn get_note_including_deleted(&self, id: Uuid) -> Result<Option<Note>, StorageError> {
1985 let id_str = id.to_string();
1986
1987 self.with_reader("get_note_including_deleted", move |conn| {
1988 let mut stmt = conn.prepare(
1989 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
1990 properties, created_at, updated_at, deleted_at, key, version \
1991 FROM notes WHERE id = ?1",
1992 )?;
1993 let mut rows = stmt.query(rusqlite::params![id_str])?;
1994 match rows.next()? {
1995 Some(row) => Ok(Some(read_note(row)?)),
1996 None => Ok(None),
1997 }
1998 })
1999 .await
2000 }
2001
2002 async fn note_sequence(&self, id: Uuid) -> Result<Option<i64>, StorageError> {
2003 let id = id.to_string();
2004 self.with_reader("note_sequence", move |conn| {
2005 conn.query_row(
2006 "SELECT seq FROM notes_seq WHERE note_id = ?1",
2007 rusqlite::params![id],
2008 |row| row.get(0),
2009 )
2010 .optional()
2011 })
2012 .await
2013 }
2014
2015 async fn get_notes_batch(&self, ids: &[Uuid]) -> Result<Vec<Note>, StorageError> {
2016 self.get_notes_batch_inner(ids, false).await
2017 }
2018
2019 async fn get_notes_batch_including_deleted(
2020 &self,
2021 ids: &[Uuid],
2022 ) -> Result<Vec<Note>, StorageError> {
2023 self.get_notes_batch_inner(ids, true).await
2024 }
2025
2026 async fn get_note_visibility_batch(
2027 &self,
2028 ids: &[Uuid],
2029 ) -> Result<Vec<NoteVisibility>, StorageError> {
2030 if ids.is_empty() {
2031 return Ok(vec![]);
2032 }
2033 const CHUNK: usize = 900;
2035 let id_strings: Vec<String> = ids.iter().map(Uuid::to_string).collect();
2036 let mut result = Vec::with_capacity(ids.len());
2037 for chunk in id_strings.chunks(CHUNK) {
2038 let chunk_owned = chunk.to_vec();
2039 let rows = self
2040 .with_reader("get_note_visibility_batch", move |conn| {
2041 let placeholders = (1..=chunk_owned.len())
2042 .map(|i| format!("?{i}"))
2043 .collect::<Vec<_>>()
2044 .join(", ");
2045 let sql = format!(
2046 "SELECT id, namespace, deleted_at FROM notes WHERE id IN ({placeholders})"
2047 );
2048 let mut stmt = conn.prepare(&sql)?;
2049 let params: Vec<&dyn rusqlite::types::ToSql> = chunk_owned
2050 .iter()
2051 .map(|s| s as &dyn rusqlite::types::ToSql)
2052 .collect();
2053 let rows = stmt.query_map(params.as_slice(), |row| {
2054 let id: String = row.get(0)?;
2055 Ok(NoteVisibility {
2056 id: parse_uuid(&id)?,
2057 namespace: row.get(1)?,
2058 deleted_at: row.get(2)?,
2059 })
2060 })?;
2061 rows.collect::<rusqlite::Result<Vec<_>>>()
2062 })
2063 .await?;
2064 result.extend(rows);
2065 }
2066 Ok(result)
2067 }
2068
2069 async fn delete_note(&self, id: Uuid, mode: DeleteMode) -> Result<bool, StorageError> {
2070 match mode {
2071 DeleteMode::Soft => {
2072 let now = chrono::Utc::now().timestamp_micros();
2073 let statement = note_soft_delete_statement(id, now);
2074 self.with_writer("delete_note_soft", move |conn| {
2075 let mut stmt = conn.prepare(&statement.sql)?;
2076 bind_params(&mut stmt, &statement.params)?;
2077 Ok(stmt.raw_execute()? > 0)
2078 })
2079 .await
2080 }
2081 DeleteMode::Hard => {
2082 let note_statement = note_hard_delete_statement(id);
2083 let attachment_statement =
2084 delete_record_attachments_statement(id, AttachmentSubstrate::Note);
2085 self.with_writer_tx("delete_note_hard", move |conn| {
2086 let mut note_stmt = conn.prepare(¬e_statement.sql)?;
2087 bind_params(&mut note_stmt, ¬e_statement.params)?;
2088 let deleted = note_stmt.raw_execute()? > 0;
2089 drop(note_stmt);
2090 if deleted {
2091 let mut attachment_stmt = conn.prepare(&attachment_statement.sql)?;
2092 bind_params(&mut attachment_stmt, &attachment_statement.params)?;
2093 attachment_stmt.raw_execute()?;
2094 }
2095 Ok(deleted)
2096 })
2097 .await
2098 }
2099 }
2100 }
2101
2102 async fn query_notes(
2103 &self,
2104 namespace: &str,
2105 kind: Option<&str>,
2106 page: PageRequest,
2107 ) -> Result<Page<Note>, StorageError> {
2108 let namespace = namespace.to_string();
2109 let kind = kind.map(|k| k.to_string());
2110 let limit_i64 = i64::from(page.limit);
2111 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
2112 capability: StorageCapability::Notes,
2113 operation: "query_notes".into(),
2114 message: format!(
2115 "PageRequest: offset must be <= i64::MAX, got {}",
2116 page.offset
2117 ),
2118 })?;
2119
2120 self.with_reader("query_notes", move |conn| {
2121 let (count_sql, count_params) = build_note_where(&namespace, kind.as_deref());
2122 let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
2123
2124 let (where_sql, mut data_params) = build_note_where(&namespace, kind.as_deref());
2125 data_params.push(Box::new(limit_i64));
2126 data_params.push(Box::new(offset_i64));
2127
2128 let limit_idx = data_params.len() - 1;
2129 let offset_idx = data_params.len();
2130
2131 let data_sql = format!(
2132 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, expires_at, \
2133 properties, created_at, updated_at, deleted_at, key, version \
2134 FROM notes{} ORDER BY created_at DESC, id ASC LIMIT ?{} OFFSET ?{}",
2135 where_sql, limit_idx, offset_idx,
2136 );
2137
2138 query_note_page_snapshot(
2139 conn,
2140 "query_notes",
2141 &namespace,
2142 &count_sql,
2143 &count_params,
2144 &data_sql,
2145 &data_params,
2146 )
2147 })
2148 .await
2149 }
2150
2151 async fn query_notes_count_free(
2152 &self,
2153 namespace: &str,
2154 kind: Option<&str>,
2155 page: PageRequest,
2156 ) -> Result<Page<Note>, StorageError> {
2157 let namespace = namespace.to_string();
2158 let kind = kind.map(str::to_string);
2159 let limit_i64 = i64::from(page.limit);
2160 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
2161 capability: StorageCapability::Notes,
2162 operation: "query_notes_count_free".into(),
2163 message: format!(
2164 "PageRequest: offset must be <= i64::MAX, got {}",
2165 page.offset
2166 ),
2167 })?;
2168
2169 self.with_reader("query_notes_count_free", move |conn| {
2170 let (where_sql, mut params) = build_note_where(&namespace, kind.as_deref());
2171 params.push(Box::new(limit_i64));
2172 params.push(Box::new(offset_i64));
2173 let limit_idx = params.len() - 1;
2174 let offset_idx = params.len();
2175 let sql = format!(
2176 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
2177 expires_at, properties, created_at, updated_at, deleted_at, key, version \
2178 FROM notes{where_sql} ORDER BY created_at DESC, id ASC \
2179 LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
2180 );
2181
2182 let mut stmt = conn.prepare(&sql)?;
2183 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2184 params.iter().map(|param| param.as_ref()).collect();
2185 let mut rows = stmt.query(param_refs.as_slice())?;
2186 let mut items = Vec::new();
2187 while let Some(row) = rows.next()? {
2188 items.push(read_note(row)?);
2189 }
2190
2191 Ok(Page { items, total: None })
2192 })
2193 .await
2194 }
2195
2196 async fn query_notes_filtered(
2197 &self,
2198 namespace: &str,
2199 filter: &NoteFilter,
2200 page: PageRequest,
2201 ) -> Result<Page<Note>, StorageError> {
2202 if filter.unordered {
2203 return Err(StorageError::InvalidInput {
2204 capability: StorageCapability::Notes,
2205 operation: "query_notes_filtered".into(),
2206 message: "NoteFilter.unordered is supported only by \
2207 query_notes_filtered_count_free"
2208 .into(),
2209 });
2210 }
2211 for pf in &filter.property_filters {
2213 validate_json_path(&pf.json_path)?;
2214 }
2215 if let Some((path, _)) = &filter.order_by {
2216 validate_json_path(path)?;
2217 }
2218 if filter.order_by_instant && !matches!(filter.order_by.as_ref(), Some((_, SortDir::Asc))) {
2219 return Err(StorageError::InvalidInput {
2220 capability: StorageCapability::Notes,
2221 operation: "query_notes_filtered".into(),
2222 message: "order_by_instant requires an ascending property order".into(),
2223 });
2224 }
2225 if filter.after.is_some() || filter.after_instant.is_some() {
2226 return Err(StorageError::InvalidInput {
2227 capability: StorageCapability::Notes,
2228 operation: "query_notes_filtered".into(),
2229 message: "NoteFilter.after or after_instant (keyset pagination) is not supported by this \
2230 method: it computes an exact COUNT(*) total over the whole \
2231 matching set, which has no defined meaning paired with a seek \
2232 boundary; use query_notes_filtered_count_free instead, which \
2233 seeks and returns total: None"
2234 .into(),
2235 });
2236 }
2237
2238 let namespace = namespace.to_string();
2239 let filter = filter.clone();
2240 let limit_i64 = i64::from(page.limit);
2241 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
2242 capability: StorageCapability::Notes,
2243 operation: "query_notes_filtered".into(),
2244 message: format!(
2245 "PageRequest: offset must be <= i64::MAX, got {}",
2246 page.offset
2247 ),
2248 })?;
2249
2250 self.with_indexed_reader("query_notes_filtered", move |conn| {
2251 let (count_sql, count_params) = build_note_filter_read_clause(&namespace, &filter)?;
2252 let count_sql = format!("SELECT COUNT(*) FROM notes{count_sql}");
2253
2254 let (where_sql, mut data_params) = build_note_filter_read_clause(&namespace, &filter)?;
2255 data_params.push(Box::new(limit_i64));
2256 data_params.push(Box::new(offset_i64));
2257
2258 let order_clause = note_filter_page_order_clause(&filter);
2262
2263 let limit_idx = data_params.len() - 1;
2264 let offset_idx = data_params.len();
2265 let data_sql = format!(
2266 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
2267 expires_at, properties, created_at, updated_at, deleted_at, key, version \
2268 FROM notes{}{order_clause} LIMIT ?{} OFFSET ?{}",
2269 where_sql, limit_idx, offset_idx,
2270 );
2271
2272 query_note_page_snapshot(
2273 conn,
2274 "query_notes_filtered",
2275 &namespace,
2276 &count_sql,
2277 &count_params,
2278 &data_sql,
2279 &data_params,
2280 )
2281 })
2282 .await
2283 }
2284
2285 async fn query_notes_filtered_count_free(
2286 &self,
2287 namespace: &str,
2288 filter: &NoteFilter,
2289 page: PageRequest,
2290 ) -> Result<Page<Note>, StorageError> {
2291 for property_filter in &filter.property_filters {
2292 validate_json_path(&property_filter.json_path)?;
2293 }
2294 if let Some((path, _)) = &filter.order_by {
2295 validate_json_path(path)?;
2296 }
2297 if filter.order_by_instant && !matches!(filter.order_by.as_ref(), Some((_, SortDir::Asc))) {
2298 return Err(StorageError::InvalidInput {
2299 capability: StorageCapability::Notes,
2300 operation: "query_notes_filtered_count_free".into(),
2301 message: "order_by_instant requires an ascending property order".into(),
2302 });
2303 }
2304 if filter.order_by_instant && filter.unordered {
2305 return Err(StorageError::InvalidInput {
2306 capability: StorageCapability::Notes,
2307 operation: "query_notes_filtered_count_free".into(),
2308 message: "order_by_instant is incompatible with unordered pages".into(),
2309 });
2310 }
2311 if filter.after_instant.is_some() && !filter.order_by_instant {
2312 return Err(StorageError::InvalidInput {
2313 capability: StorageCapability::Notes,
2314 operation: "query_notes_filtered_count_free".into(),
2315 message: "after_instant requires order_by_instant".into(),
2316 });
2317 }
2318 if filter.after.is_some() && filter.after_instant.is_some() {
2319 return Err(StorageError::InvalidInput {
2320 capability: StorageCapability::Notes,
2321 operation: "query_notes_filtered_count_free".into(),
2322 message: "after and after_instant are mutually exclusive".into(),
2323 });
2324 }
2325 if filter.after.is_some() && filter.order_by.is_some() {
2326 return Err(StorageError::InvalidInput {
2327 capability: StorageCapability::Notes,
2328 operation: "query_notes_filtered_count_free".into(),
2329 message: "NoteFilter.after is incompatible with a custom order_by; it is \
2330 defined only over the default created_at DESC, id ASC order"
2331 .into(),
2332 });
2333 }
2334 if (filter.after.is_some() || filter.after_instant.is_some()) && page.offset != 0 {
2335 return Err(StorageError::InvalidInput {
2336 capability: StorageCapability::Notes,
2337 operation: "query_notes_filtered_count_free".into(),
2338 message: "NoteFilter.after or after_instant and a non-zero PageRequest.offset \
2339 are mutually exclusive pagination strategies; pass offset: 0"
2340 .into(),
2341 });
2342 }
2343
2344 let namespace = namespace.to_string();
2345 let filter = filter.clone();
2346 let limit_i64 = i64::from(page.limit);
2347 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
2348 capability: StorageCapability::Notes,
2349 operation: "query_notes_filtered_count_free".into(),
2350 message: format!(
2351 "PageRequest: offset must be <= i64::MAX, got {}",
2352 page.offset
2353 ),
2354 })?;
2355
2356 self.with_indexed_reader("query_notes_filtered_count_free", move |conn| {
2357 if let Some(after) = &filter.after_instant {
2358 let mut base_filter = filter.clone();
2359 base_filter.after_instant = None;
2360 let items =
2361 fetch_notes_after_instant(conn, &namespace, &base_filter, after, limit_i64)?;
2362 return Ok(Page { items, total: None });
2363 }
2364 if let Some(after) = &filter.after {
2365 let mut base_filter = filter.clone();
2366 base_filter.after = None;
2367 let items = fetch_notes_after(conn, &namespace, &base_filter, after, limit_i64)?;
2368 return Ok(Page { items, total: None });
2369 }
2370
2371 let (where_sql, mut params) = build_note_filter_read_clause(&namespace, &filter)?;
2372 params.push(Box::new(limit_i64));
2373 params.push(Box::new(offset_i64));
2374 let limit_idx = params.len() - 1;
2375 let offset_idx = params.len();
2376 let order_clause = note_filter_page_order_clause(&filter);
2377 let sql = format!(
2378 "SELECT {NOTE_COLUMNS} FROM notes{where_sql}{order_clause} \
2379 LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
2380 );
2381
2382 let mut stmt = conn.prepare(&sql)?;
2383 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2384 params.iter().map(|param| param.as_ref()).collect();
2385 let mut rows = stmt.query(param_refs.as_slice())?;
2386 let mut items = Vec::new();
2387 while let Some(row) = rows.next()? {
2388 items.push(read_note(row)?);
2389 #[cfg(test)]
2393 if items.len() == 1 {
2394 tests::page_snapshot_seam::hook("query_notes_filtered_count_free", &namespace);
2395 }
2396 }
2397
2398 Ok(Page { items, total: None })
2399 })
2400 .await
2401 }
2402
2403 async fn count_notes_filtered_in_snapshot(
2404 &self,
2405 namespace: &str,
2406 filters: &[NoteFilter],
2407 ) -> Result<Vec<u64>, StorageError> {
2408 for filter in filters {
2409 for pf in &filter.property_filters {
2410 validate_json_path(&pf.json_path)?;
2411 }
2412 }
2413
2414 let namespace = namespace.to_string();
2415 let filters = filters.to_vec();
2416 self.with_indexed_reader("count_notes_filtered_in_snapshot", move |conn| {
2417 let tx = rusqlite::Transaction::new_unchecked(
2418 conn,
2419 rusqlite::TransactionBehavior::Deferred,
2420 )?;
2421 let mut counts = Vec::with_capacity(filters.len());
2422 for filter in filters.iter() {
2423 #[cfg(test)]
2424 if !counts.is_empty() {
2425 tests::page_snapshot_seam::hook("count_notes_filtered_in_snapshot", &namespace);
2426 }
2427
2428 let (where_sql, params) = build_note_filter_read_clause(&namespace, filter)?;
2429 let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
2430 let mut stmt = tx.prepare(&sql)?;
2431 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2432 params.iter().map(|param| param.as_ref()).collect();
2433 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2434 counts.push(count as u64);
2435 }
2436 tx.commit()?;
2437 Ok(counts)
2438 })
2439 .await
2440 }
2441
2442 async fn count_notes_filtered_bounded_in_snapshot(
2443 &self,
2444 namespace: &str,
2445 filters: &[NoteFilter],
2446 cap: u32,
2447 ) -> Result<Vec<BoundedCount>, StorageError> {
2448 for filter in filters {
2449 for property_filter in &filter.property_filters {
2450 validate_json_path(&property_filter.json_path)?;
2451 }
2452 }
2453
2454 let namespace = namespace.to_string();
2455 let filters = filters.to_vec();
2456 let cap_u64 = u64::from(cap);
2457 let probe_limit_i64 = i64::from(cap) + 1;
2458 self.with_indexed_reader("count_notes_filtered_bounded_in_snapshot", move |conn| {
2459 let tx = rusqlite::Transaction::new_unchecked(
2460 conn,
2461 rusqlite::TransactionBehavior::Deferred,
2462 )?;
2463 let mut counts = Vec::with_capacity(filters.len());
2464 for filter in &filters {
2465 #[cfg(test)]
2466 if !counts.is_empty() {
2467 tests::page_snapshot_seam::hook(
2468 "count_notes_filtered_bounded_in_snapshot",
2469 &namespace,
2470 );
2471 }
2472
2473 let (where_sql, mut params) = build_note_filter_read_clause(&namespace, filter)?;
2474 params.push(Box::new(probe_limit_i64));
2475 let limit_idx = params.len();
2476 let sql = format!(
2481 "SELECT COUNT(*) FROM (SELECT 1 FROM notes{where_sql} LIMIT ?{limit_idx})"
2482 );
2483 let mut stmt = tx.prepare(&sql)?;
2484 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2485 params.iter().map(|param| param.as_ref()).collect();
2486 let observed: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2487 let observed = observed as u64;
2488 counts.push(BoundedCount {
2489 count: observed.min(cap_u64),
2490 cap: cap_u64,
2491 saturated: observed > cap_u64,
2492 });
2493 }
2494 tx.commit()?;
2495 Ok(counts)
2496 })
2497 .await
2498 }
2499
2500 async fn query_notes_filtered_after(
2501 &self,
2502 namespace: &str,
2503 filter: &NoteFilter,
2504 after: Option<SeekCursor>,
2505 limit: u32,
2506 ) -> Result<SeekPage<Note>, StorageError> {
2507 if limit == 0 {
2508 return Ok(SeekPage::default());
2509 }
2510 if filter.order_by.is_some() {
2511 return Err(StorageError::InvalidInput {
2512 capability: StorageCapability::Notes,
2513 operation: "query_notes_filtered_after".into(),
2514 message: "custom order_by is not compatible with insertion-sequence pagination"
2515 .into(),
2516 });
2517 }
2518 for property_filter in &filter.property_filters {
2519 validate_json_path(&property_filter.json_path)?;
2520 }
2521
2522 let namespace = namespace.to_string();
2523 let filter = filter.clone();
2524 let limit_usize = limit as usize;
2525 let probe_limit_i64 = i64::from(limit) + 1;
2526 self.with_reader("query_notes_filtered_after", move |conn| {
2527 let (mut where_sql, mut params) = build_note_filter_where(&namespace, &filter)?;
2528 if let Some(cursor) = after {
2529 params.push(Box::new(cursor.sequence));
2530 where_sql.push_str(&format!(" AND notes_seq.seq > ?{}", params.len()));
2531 }
2532 params.push(Box::new(probe_limit_i64));
2533 let limit_idx = params.len();
2534 let sql = format!(
2537 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
2538 expires_at, properties, created_at, updated_at, deleted_at, key, version, notes_seq.seq \
2539 FROM notes_seq CROSS JOIN notes ON notes.id = notes_seq.note_id{where_sql} \
2540 ORDER BY notes_seq.seq ASC LIMIT ?{limit_idx}"
2541 );
2542 let mut stmt = conn.prepare(&sql)?;
2543 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2544 params.iter().map(|param| param.as_ref()).collect();
2545 let rows = stmt.query_map(param_refs.as_slice(), |row| {
2546 Ok((read_note(row)?, row.get::<_, i64>(15)?))
2547 })?;
2548 let mut entries = rows.collect::<Result<Vec<_>, _>>()?;
2549 let has_more = entries.len() > limit_usize;
2550 if has_more {
2551 entries.truncate(limit_usize);
2552 }
2553 let next_after = if has_more {
2554 entries.last().map(|(note, sequence)| SeekCursor {
2555 sequence: *sequence,
2556 id: note.id,
2557 })
2558 } else {
2559 None
2560 };
2561 let items = entries.into_iter().map(|(note, _)| note).collect();
2562 Ok(SeekPage { items, next_after })
2563 })
2564 .await
2565 }
2566
2567 async fn query_notes_filtered_bounded(
2568 &self,
2569 namespace: &str,
2570 filter: &NoteFilter,
2571 max_rows: u32,
2572 ) -> Result<Vec<Note>, StorageError> {
2573 if filter.unordered {
2574 return Err(StorageError::InvalidInput {
2575 capability: StorageCapability::Notes,
2576 operation: "query_notes_filtered_bounded".into(),
2577 message: "NoteFilter.unordered is supported only by \
2578 query_notes_filtered_count_free"
2579 .into(),
2580 });
2581 }
2582 for pf in &filter.property_filters {
2583 validate_json_path(&pf.json_path)?;
2584 }
2585 if let Some((path, _)) = &filter.order_by {
2586 validate_json_path(path)?;
2587 }
2588
2589 let namespace = namespace.to_string();
2590 let filter = filter.clone();
2591 let limit_i64 = i64::from(max_rows) + 1;
2592
2593 self.with_indexed_reader("query_notes_filtered_bounded", move |conn| {
2594 let (where_sql, mut data_params) = build_note_filter_read_clause(&namespace, &filter)?;
2595 data_params.push(Box::new(limit_i64));
2596 let limit_idx = data_params.len();
2597
2598 let order_clause = match &filter.order_by {
2602 Some((path, dir)) => {
2603 let dir_str = match dir {
2604 SortDir::Asc => "ASC",
2605 SortDir::Desc => "DESC",
2606 };
2607 format!(" ORDER BY {} {dir_str}, id ASC", json_extract_expr(path))
2608 }
2609 None => " ORDER BY created_at DESC, id ASC".to_string(),
2610 };
2611
2612 let data_sql = format!(
2613 "SELECT id, namespace, kind, status, name, content, salience, decay_factor, \
2614 expires_at, properties, created_at, updated_at, deleted_at, key, version \
2615 FROM notes{where_sql}{order_clause} LIMIT ?{limit_idx}",
2616 );
2617
2618 let mut stmt = conn.prepare(&data_sql)?;
2619 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2620 data_params.iter().map(|p| p.as_ref()).collect();
2621 let rows = stmt.query_map(param_refs.as_slice(), read_note)?;
2622
2623 let mut items = Vec::new();
2624 for row in rows {
2625 items.push(row?);
2626 }
2627 Ok(items)
2628 })
2629 .await
2630 }
2631
2632 async fn count_notes(&self, namespace: &str, kind: Option<&str>) -> Result<u64, StorageError> {
2633 let namespace = namespace.to_string();
2634 let kind = kind.map(|k| k.to_string());
2635
2636 self.with_reader("count_notes", move |conn| {
2637 let (where_sql, params) = build_note_where(&namespace, kind.as_deref());
2638 let sql = format!("SELECT COUNT(*) FROM notes{}", where_sql);
2639 let mut stmt = conn.prepare(&sql)?;
2640 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2641 params.iter().map(|p| p.as_ref()).collect();
2642 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2643 Ok(count as u64)
2644 })
2645 .await
2646 }
2647
2648 async fn count_notes_in_namespaces(
2649 &self,
2650 namespaces: &[String],
2651 kind: Option<&str>,
2652 ) -> Result<u64, StorageError> {
2653 let namespaces: Vec<String> = namespaces
2654 .iter()
2655 .cloned()
2656 .collect::<HashSet<_>>()
2657 .into_iter()
2658 .collect();
2659 let kind = kind.map(str::to_string);
2660
2661 self.with_reader("count_notes_in_namespaces", move |conn| {
2662 let mut total = 0;
2663 for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
2664 let (where_sql, params) = build_note_where_for_namespaces(chunk, kind.as_deref());
2665 let sql = format!("SELECT COUNT(*) FROM notes{where_sql}");
2666 let mut stmt = conn.prepare(&sql)?;
2667 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2668 params.iter().map(|p| p.as_ref()).collect();
2669 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
2670 total += count as u64;
2671 }
2672 Ok(total)
2673 })
2674 .await
2675 }
2676}
2677
2678const NOTES_DDL: &str = include_str!("../../sql/notes-ddl.sql");
2683
2684const NOTES_SEQ_REPAIR_DDL: &str = include_str!("../../sql/008-notes-seq-repair.sql");
2691
2692pub(crate) fn ensure_notes_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
2693 conn.execute_batch(NOTES_DDL)
2694}
2695
2696pub(crate) fn repair_notes_seq(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
2702 conn.execute_batch(NOTES_SEQ_REPAIR_DDL)
2703}
2704
2705#[cfg(test)]
2706#[path = "note_tests.rs"]
2707mod tests;
2708
2709#[cfg(test)]
2710#[path = "comm_filter_plan_tests.rs"]
2711mod comm_filter_plan_tests;
2712
2713#[cfg(test)]
2714#[path = "note_list_plan_tests.rs"]
2715mod note_list_plan_tests;
2716
2717#[cfg(test)]
2718#[path = "note_busy_tests.rs"]
2719mod direct_busy_tests;