khive_storage/note.rs
1//! Note storage capability — temporal-referential record CRUD.
2
3use async_trait::async_trait;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use uuid::Uuid;
7
8use crate::types::{
9 BatchWriteSummary, BoundedCount, DeleteMode, Page, PageRequest, SeekCursor, SeekPage, SqlValue,
10 StorageResult,
11};
12
13/// A storage-level note record. Flat, SQL-friendly representation.
14#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
15pub struct Note {
16 pub id: Uuid,
17 pub namespace: String,
18 pub kind: String,
19 pub status: String,
20 pub name: Option<String>,
21 pub content: String,
22 pub salience: Option<f64>,
23 pub decay_factor: Option<f64>,
24 pub expires_at: Option<i64>,
25 pub properties: Option<Value>,
26 pub created_at: i64,
27 pub updated_at: i64,
28 pub deleted_at: Option<i64>,
29 /// Immutable caller-chosen identity, unique among live notes of the same namespace and kind.
30 #[serde(default, skip_serializing_if = "Option::is_none")]
31 pub key: Option<String>,
32 /// Persisted revision, assigned by storage and advanced on every matched update.
33 #[serde(default = "initial_note_version")]
34 pub version: i64,
35}
36
37const fn initial_note_version() -> i64 {
38 1
39}
40
41impl Note {
42 /// Create a new note with a generated UUID and current timestamp.
43 pub fn new(
44 namespace: impl Into<String>,
45 kind: impl Into<String>,
46 content: impl Into<String>,
47 ) -> Self {
48 let now = chrono::Utc::now().timestamp_micros();
49 Self {
50 id: Uuid::new_v4(),
51 namespace: namespace.into(),
52 kind: kind.into(),
53 status: "active".to_string(),
54 name: None,
55 content: content.into(),
56 salience: None,
57 decay_factor: None,
58 expires_at: None,
59 properties: None,
60 created_at: now,
61 updated_at: now,
62 deleted_at: None,
63 key: None,
64 version: 1,
65 }
66 }
67
68 /// Set the note display name.
69 pub fn with_name(mut self, n: impl Into<String>) -> Self {
70 self.name = Some(n.into());
71 self
72 }
73
74 /// Set salience (infallible). Rejects non-finite values by returning `self`
75 /// unchanged; clamps finite values to `[0.0, 1.0]`. Prefer
76 /// [`try_with_salience`](Self::try_with_salience) at public boundaries.
77 pub fn with_salience(mut self, s: f64) -> Self {
78 if !s.is_finite() {
79 return self;
80 }
81 self.salience = Some(s.clamp(0.0, 1.0));
82 self
83 }
84
85 /// Set decay factor (infallible). Rejects non-finite values by returning
86 /// `self` unchanged; floors finite values at `0.0`. Prefer
87 /// [`try_with_decay`](Self::try_with_decay) at public boundaries.
88 pub fn with_decay(mut self, d: f64) -> Self {
89 if !d.is_finite() {
90 return self;
91 }
92 self.decay_factor = Some(d.max(0.0));
93 self
94 }
95
96 /// Set salience with validation. Returns an error for non-finite or
97 /// out-of-range `[0.0, 1.0]` values.
98 pub fn try_with_salience(mut self, s: f64) -> Result<Self, String> {
99 if !s.is_finite() {
100 return Err(format!("salience must be finite, got {s}"));
101 }
102 if !(0.0..=1.0).contains(&s) {
103 return Err(format!("salience must be in [0.0, 1.0], got {s}"));
104 }
105 self.salience = Some(s);
106 Ok(self)
107 }
108
109 /// Set decay factor with validation. Returns an error for non-finite or
110 /// negative values.
111 pub fn try_with_decay(mut self, d: f64) -> Result<Self, String> {
112 if !d.is_finite() {
113 return Err(format!("decay_factor must be finite, got {d}"));
114 }
115 if d < 0.0 {
116 return Err(format!("decay_factor must be >= 0.0, got {d}"));
117 }
118 self.decay_factor = Some(d);
119 Ok(self)
120 }
121
122 /// Set the note properties JSON blob.
123 pub fn with_properties(mut self, p: Value) -> Self {
124 self.properties = Some(p);
125 self
126 }
127}
128
129#[cfg(test)]
130mod tests {
131 use super::*;
132
133 fn base_note() -> Note {
134 Note::new("ns:test", "memory", "hello world")
135 }
136
137 #[test]
138 fn note_key_is_optional_and_unkeyed_wire_shape_is_unchanged() {
139 let note = base_note();
140 assert_eq!(note.key, None);
141 let legacy = serde_json::to_value(¬e).unwrap();
142 assert!(legacy.get("key").is_none());
143 assert_eq!(serde_json::from_value::<Note>(legacy).unwrap(), note);
144
145 let mut keyed = note;
146 keyed.key = Some("operation-1".to_string());
147 let encoded = serde_json::to_value(&keyed).unwrap();
148 assert_eq!(encoded["key"], "operation-1");
149 assert_eq!(serde_json::from_value::<Note>(encoded).unwrap(), keyed);
150 }
151
152 // -- with_salience --
153
154 #[test]
155 fn with_salience_clamps_to_range() {
156 let n = base_note().with_salience(1.5);
157 assert_eq!(n.salience, Some(1.0));
158 let n = base_note().with_salience(-0.1);
159 assert_eq!(n.salience, Some(0.0));
160 let n = base_note().with_salience(0.7);
161 assert_eq!(n.salience, Some(0.7));
162 }
163
164 #[test]
165 fn with_salience_ignores_nan() {
166 let n = base_note().with_salience(f64::NAN);
167 assert_eq!(n.salience, None, "NaN must not set salience");
168 }
169
170 #[test]
171 fn with_salience_ignores_inf() {
172 let n = base_note().with_salience(f64::INFINITY);
173 assert_eq!(n.salience, None, "+Inf must not set salience");
174 let n = base_note().with_salience(f64::NEG_INFINITY);
175 assert_eq!(n.salience, None, "-Inf must not set salience");
176 }
177
178 // -- with_decay --
179
180 #[test]
181 fn with_decay_floors_at_zero() {
182 let n = base_note().with_decay(-1.0);
183 assert_eq!(n.decay_factor, Some(0.0));
184 let n = base_note().with_decay(0.5);
185 assert_eq!(n.decay_factor, Some(0.5));
186 }
187
188 #[test]
189 fn with_decay_ignores_nan() {
190 let n = base_note().with_decay(f64::NAN);
191 assert_eq!(n.decay_factor, None, "NaN must not set decay_factor");
192 }
193
194 #[test]
195 fn with_decay_ignores_inf() {
196 let n = base_note().with_decay(f64::INFINITY);
197 assert_eq!(n.decay_factor, None, "+Inf must not set decay_factor");
198 }
199
200 // -- try_with_salience --
201
202 #[test]
203 fn try_with_salience_accepts_valid_range() {
204 let n = base_note().try_with_salience(0.0).unwrap();
205 assert_eq!(n.salience, Some(0.0));
206 let n = base_note().try_with_salience(1.0).unwrap();
207 assert_eq!(n.salience, Some(1.0));
208 let n = base_note().try_with_salience(0.85).unwrap();
209 assert_eq!(n.salience, Some(0.85));
210 }
211
212 #[test]
213 fn try_with_salience_rejects_nan() {
214 let err = base_note().try_with_salience(f64::NAN).unwrap_err();
215 assert!(err.contains("finite"), "error must mention finite: {err}");
216 }
217
218 #[test]
219 fn try_with_salience_rejects_out_of_range() {
220 let err = base_note().try_with_salience(1.1).unwrap_err();
221 assert!(err.contains("1.0"), "error must mention bound: {err}");
222 let err = base_note().try_with_salience(-0.01).unwrap_err();
223 assert!(err.contains("0.0"), "error must mention bound: {err}");
224 }
225
226 // -- try_with_decay --
227
228 #[test]
229 fn try_with_decay_accepts_valid_values() {
230 let n = base_note().try_with_decay(0.0).unwrap();
231 assert_eq!(n.decay_factor, Some(0.0));
232 let n = base_note().try_with_decay(2.5).unwrap();
233 assert_eq!(n.decay_factor, Some(2.5));
234 }
235
236 #[test]
237 fn try_with_decay_rejects_nan() {
238 let err = base_note().try_with_decay(f64::NAN).unwrap_err();
239 assert!(err.contains("finite"), "error must mention finite: {err}");
240 }
241
242 #[test]
243 fn try_with_decay_rejects_negative() {
244 let err = base_note().try_with_decay(-0.1).unwrap_err();
245 assert!(err.contains("0.0"), "error must mention bound: {err}");
246 }
247
248 #[derive(Default)]
249 struct DefaultOnlyNoteStore {
250 query_calls: std::sync::atomic::AtomicUsize,
251 }
252
253 impl DefaultOnlyNoteStore {
254 fn query_call_count(&self) -> usize {
255 self.query_calls.load(std::sync::atomic::Ordering::SeqCst)
256 }
257 }
258
259 fn unsupported<T>(operation: &'static str) -> StorageResult<T> {
260 Err(crate::StorageError::Unsupported {
261 capability: crate::StorageCapability::Notes,
262 operation: operation.into(),
263 message: "test double does not implement this operation".into(),
264 })
265 }
266
267 #[async_trait]
268 impl NoteStore for DefaultOnlyNoteStore {
269 async fn upsert_note(&self, _note: Note) -> StorageResult<()> {
270 unsupported("upsert_note")
271 }
272
273 async fn upsert_notes(&self, _notes: Vec<Note>) -> StorageResult<BatchWriteSummary> {
274 unsupported("upsert_notes")
275 }
276
277 async fn get_note(&self, _id: Uuid) -> StorageResult<Option<Note>> {
278 unsupported("get_note")
279 }
280
281 async fn get_note_including_deleted(&self, _id: Uuid) -> StorageResult<Option<Note>> {
282 unsupported("get_note_including_deleted")
283 }
284
285 async fn delete_note(&self, _id: Uuid, _mode: DeleteMode) -> StorageResult<bool> {
286 unsupported("delete_note")
287 }
288
289 async fn update_note_properties(
290 &self,
291 _id: Uuid,
292 _properties: Option<Value>,
293 _updated_at: i64,
294 ) -> StorageResult<bool> {
295 unsupported("update_note_properties")
296 }
297
298 async fn set_note_property(
299 &self,
300 _id: Uuid,
301 _key: &str,
302 _value: Value,
303 _updated_at: i64,
304 ) -> StorageResult<bool> {
305 unsupported("set_note_property")
306 }
307
308 async fn try_patch_note_property(
309 &self,
310 _id: Uuid,
311 _namespace: &str,
312 _filter: &NoteFilter,
313 _json_path: &str,
314 _value: Value,
315 _updated_at: i64,
316 ) -> StorageResult<bool> {
317 unsupported("try_patch_note_property")
318 }
319
320 async fn query_notes(
321 &self,
322 _namespace: &str,
323 _kind: Option<&str>,
324 _page: PageRequest,
325 ) -> StorageResult<Page<Note>> {
326 self.query_calls
327 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
328 Ok(Page {
329 items: Vec::new(),
330 total: Some(0),
331 })
332 }
333
334 async fn query_notes_filtered(
335 &self,
336 _namespace: &str,
337 _filter: &NoteFilter,
338 _page: PageRequest,
339 ) -> StorageResult<Page<Note>> {
340 self.query_calls
341 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
342 Ok(Page {
343 items: Vec::new(),
344 total: Some(0),
345 })
346 }
347
348 async fn query_notes_filtered_bounded(
349 &self,
350 _namespace: &str,
351 _filter: &NoteFilter,
352 _max_rows: u32,
353 ) -> StorageResult<Vec<Note>> {
354 unsupported("query_notes_filtered_bounded")
355 }
356
357 async fn count_notes(&self, _namespace: &str, _kind: Option<&str>) -> StorageResult<u64> {
358 unsupported("count_notes")
359 }
360
361 async fn try_insert_note(&self, _note: Note) -> StorageResult<bool> {
362 unsupported("try_insert_note")
363 }
364 }
365
366 #[tokio::test]
367 async fn default_snapshot_count_is_fail_closed_without_query_fallback() {
368 let store = DefaultOnlyNoteStore::default();
369 let result = store
370 .count_notes_filtered_in_snapshot(
371 "ns:test",
372 &[NoteFilter::default(), NoteFilter::default()],
373 )
374 .await;
375
376 assert!(
377 matches!(
378 result,
379 Err(crate::StorageError::Unsupported { ref operation, .. })
380 if operation == "count_notes_filtered_in_snapshot"
381 ),
382 "inherited snapshot count must fail closed, got {result:?}"
383 );
384 assert_eq!(
385 store.query_call_count(),
386 0,
387 "inherited snapshot count must not fall back to independent queries"
388 );
389 }
390
391 #[tokio::test]
392 async fn count_free_page_default_is_fail_closed_without_exact_query_fallback() {
393 let store = DefaultOnlyNoteStore::default();
394 let result = store
395 .query_notes_filtered_count_free(
396 "ns:test",
397 &NoteFilter::default(),
398 PageRequest::default(),
399 )
400 .await;
401
402 assert!(
403 matches!(
404 result,
405 Err(crate::StorageError::Unsupported { ref operation, .. })
406 if operation == "query_notes_filtered_count_free"
407 ),
408 "inherited count-free query must fail closed, got {result:?}"
409 );
410 assert_eq!(
411 store.query_call_count(),
412 0,
413 "count-free default must not invoke the exact-count page"
414 );
415 }
416
417 #[tokio::test]
418 async fn unfiltered_count_free_page_default_is_fail_closed_without_exact_query_fallback() {
419 let store = DefaultOnlyNoteStore::default();
420 let result = store
421 .query_notes_count_free("ns:test", Some("message"), PageRequest::default())
422 .await;
423
424 assert!(
425 matches!(
426 result,
427 Err(crate::StorageError::Unsupported { ref operation, .. })
428 if operation == "query_notes_count_free"
429 ),
430 "inherited count-free query must fail closed, got {result:?}"
431 );
432 assert_eq!(
433 store.query_call_count(),
434 0,
435 "count-free default must not invoke the exact-count page"
436 );
437 }
438
439 #[tokio::test]
440 async fn bounded_snapshot_count_default_is_fail_closed_without_query_fallback() {
441 let store = DefaultOnlyNoteStore::default();
442 let result = store
443 .count_notes_filtered_bounded_in_snapshot(
444 "ns:test",
445 &[NoteFilter::default(), NoteFilter::default()],
446 1_000,
447 )
448 .await;
449
450 assert!(
451 matches!(
452 result,
453 Err(crate::StorageError::Unsupported { ref operation, .. })
454 if operation == "count_notes_filtered_bounded_in_snapshot"
455 ),
456 "inherited bounded snapshot count must fail closed, got {result:?}"
457 );
458 assert_eq!(
459 store.query_call_count(),
460 0,
461 "bounded snapshot count must not compose independent page queries"
462 );
463 }
464}
465
466/// Sort direction for filtered note queries.
467#[derive(Clone, Debug, Serialize, Deserialize)]
468#[serde(rename_all = "snake_case")]
469pub enum SortDir {
470 Asc,
471 Desc,
472}
473
474/// Comparison operator for a [`PropertyFilter`] on a JSON path.
475#[derive(Clone, Debug, Serialize, Deserialize)]
476#[serde(rename_all = "snake_case")]
477pub enum FilterOp {
478 Eq,
479 /// Matches rows where the JSON field equals the value OR the field is absent/NULL.
480 /// Used for properties that may be missing in legacy rows (e.g. `$.read`).
481 EqOrMissing,
482 /// Matches the supplied value using the indexable
483 /// `ifnull(json_extract(...), '')` expression. The value may be any
484 /// [`SqlValue`] accepted by the SQL adapter; in particular, `Text("")`
485 /// matches both missing values and present-but-empty values.
486 EqOrMissingIndexed,
487 /// Matches rows where the JSON field is absent or SQL-NULL.
488 JsonTypeMissing,
489 /// Matches rows where the JSON field is absent or explicitly JSON `null`,
490 /// while constraining its index key to the empty recipient key. This is
491 /// the index-friendly legacy-recipient partition used with
492 /// [`FilterOp::EqOrMissingIndexed`].
493 JsonTypeMissingOrNullIndexed,
494 /// Combines the exact-value and legacy-recipient partitions
495 /// (`EqOrMissingIndexed` + `JsonTypeMissingOrNullIndexed`) into one
496 /// predicate over the same indexable `ifnull(json_extract(...), '')`
497 /// expression, so a single index seek serves both partitions instead of
498 /// two separate bounded queries. Matches rows where the field equals the
499 /// value, OR the field is absent/JSON-`null`. A present-but-empty JSON
500 /// string value does NOT match through the legacy branch — the same
501 /// `json_type` guard `JsonTypeMissingOrNullIndexed` uses excludes it —
502 /// so this reproduces `EqOrMissing` exactly, given a non-empty compared
503 /// value.
504 EqOrLegacyIndexed,
505 /// Matches rows where a JSON text field equals the value, while treating
506 /// every missing or non-text value as that same value. The SQL adapter
507 /// emits `CASE WHEN json_type(...) = 'text' THEN json_extract(...) ELSE
508 /// value END = value`, mirroring callers whose read model assigns one
509 /// textual default to absent, JSON-null, and malformed legacy values.
510 TextEqOrNonText,
511 /// Matches a JSON text field against the supplied set, or any missing or
512 /// non-text value. The non-text branch still matches when the set is empty.
513 /// `PropertyFilter.value` is unused; the set lives in this variant.
514 TextInOrNonText(Vec<SqlValue>),
515 Ne,
516 Lt,
517 Lte,
518 Gt,
519 Gte,
520 /// Matches rows where `json_type(properties, path) = value`.
521 /// Value must be a SQLite json_type string literal: 'true', 'false', 'integer',
522 /// 'real', 'text', 'array', 'object', or 'null'.
523 JsonTypeEq,
524 /// Matches rows where the json_type is absent (NULL) OR differs from value.
525 /// Equivalent to `json_type IS NULL OR json_type != value`.
526 /// Value must be a SQLite json_type string literal: 'true', 'false',
527 /// 'integer', 'real', 'text', 'array', 'object', or 'null'. Used for
528 /// unread filter: matches any `$.read` that is NOT the JSON boolean true.
529 JsonTypeNeMissing,
530 /// Matches rows where `json_extract(properties, path)` equals any value in
531 /// the set. A row with a missing/NULL property does not match — use
532 /// `TextInOrNonText` when a textual read model includes missing/non-text
533 /// legacy values. `PropertyFilter.value` is unused for this op; the set
534 /// lives in the variant itself.
535 In(Vec<SqlValue>),
536 /// Matches rows where the property is missing/NULL OR its value is not in
537 /// the set. Used to exclude a closed set while treating an unset property
538 /// as included (e.g. comm inbox excludes `outbound`). `PropertyFilter.value` is
539 /// unused for this op; the set lives in the variant itself.
540 NotInOrMissing(Vec<SqlValue>),
541 /// Matches rows whose JSON text field starts with the supplied prefix.
542 /// Rendered as an index-seekable half-open range over the plain
543 /// `json_extract` expression (`expr >= prefix AND expr < next(prefix)`,
544 /// where `next` increments the prefix's last code point), so an index
545 /// keyed on that expression serves the seek instead of a scan. SQLite
546 /// sorts non-text values outside the text range, so a missing, null or
547 /// numeric field never matches. `PropertyFilter.value` is the prefix and
548 /// must be `SqlValue::Text`; an empty prefix matches every text value.
549 TextStartsWithIndexed,
550 /// Keep only RFC 3339 text values that parse as UTC instants.
551 Rfc3339Valid,
552 /// Compare parsed UTC instants, including subsecond precision and offsets.
553 /// The value is a `Timestamp` or RFC 3339 `Text`.
554 Rfc3339Gte,
555 Rfc3339Lte,
556}
557
558/// A single `json_extract(properties, '$.field') op value` predicate.
559///
560/// Callers import this as `khive_storage::note::PropertyFilter` to avoid
561/// collision with the vector-metadata `PropertyFilter` in `khive_storage::types`.
562#[derive(Clone, Debug, Serialize, Deserialize)]
563pub struct PropertyFilter {
564 pub json_path: String,
565 pub op: FilterOp,
566 pub value: SqlValue,
567}
568
569/// Keyset pagination boundary over the notes store's default total order
570/// (`created_at DESC, id ASC` — see `note_filter_page_order_clause`).
571///
572/// Only meaningful when [`NoteFilter::order_by`] is `None`: the boundary is
573/// expressed in terms of that default order, and combining it with a custom
574/// sort field has no defined meaning, so callers must not set both.
575#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
576pub struct NoteSeekAfter {
577 pub created_at: i64,
578 pub id: Uuid,
579}
580
581/// Boundary in an ascending RFC 3339 property order. Equal instants sort by
582/// their stored text and then note ID, so this cursor identifies one row even
583/// when distinct offsets encode the same instant.
584#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
585pub struct NoteInstantSeekAfter {
586 pub value: String,
587 pub id: Uuid,
588}
589
590#[derive(Clone, Copy, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
591#[serde(rename_all = "snake_case")]
592pub enum NoteTagMode {
593 #[default]
594 Any,
595 All,
596}
597
598#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
599#[serde(deny_unknown_fields)]
600pub struct NoteKeyCursor {
601 pub updated_at: i64,
602 pub key: String,
603 pub id: Uuid,
604}
605
606impl From<&Note> for NoteKeyCursor {
607 fn from(note: &Note) -> Self {
608 Self {
609 updated_at: note.updated_at,
610 key: note.key.clone().expect("keyed row"),
611 id: note.id,
612 }
613 }
614}
615
616/// Filter + sort options for [`NoteStore::query_notes_filtered`].
617///
618/// Designed for general property-based filtering on any JSON field, not
619/// schedule-specific, so D9 and future packs can reuse the same API.
620#[derive(Clone, Debug, Default, Serialize, Deserialize)]
621pub struct NoteFilter {
622 pub kind: Option<String>,
623 #[serde(default)]
624 pub property_filters: Vec<PropertyFilter>,
625 /// `(json_path, direction)` — `None` defaults to `created_at DESC`.
626 pub order_by: Option<(String, SortDir)>,
627 /// Interpret `order_by`'s text property as an RFC 3339 UTC instant before
628 /// sorting. Only ascending order is currently supported.
629 #[serde(default)]
630 pub order_by_instant: bool,
631 /// When true, omit the SQL ordering clause from count-free pages. The
632 /// caller is responsible for ordering the bounded result set.
633 #[serde(default)]
634 pub unordered: bool,
635 /// When non-empty, restricts results to any of these namespaces using
636 /// `namespace IN (...)`. Takes precedence over the `namespace` string
637 /// parameter passed to `query_notes_filtered`. When empty the
638 /// caller-supplied `namespace` parameter is used (backward-compatible).
639 #[serde(default)]
640 pub namespaces: Vec<String>,
641 /// Restrict to notes where `created_at >= min_created_at` (microseconds epoch).
642 /// `None` applies no lower-bound constraint.
643 pub min_created_at: Option<i64>,
644 #[serde(default)]
645 pub min_updated_at: Option<i64>,
646 #[serde(default)]
647 pub tags: Vec<String>,
648 #[serde(default)]
649 pub tag_mode: NoteTagMode,
650 /// Restrict results to rows strictly after this boundary in the default
651 /// `created_at DESC, id ASC` order, for keyset (seek) pagination that
652 /// avoids re-walking earlier pages the way `PageRequest.offset` does.
653 /// Requires `order_by` to be `None`.
654 ///
655 /// Honoured only by [`NoteStore::query_notes_filtered_count_free`], which
656 /// seeks directly to the boundary and returns `total: None`.
657 /// [`NoteStore::query_notes_filtered`] rejects a non-`None` value with
658 /// `StorageError::InvalidInput`: it computes an exact `COUNT(*)` total
659 /// over the whole matching set, which has no defined meaning paired with
660 /// a seek boundary.
661 #[serde(default)]
662 pub after: Option<NoteSeekAfter>,
663 /// Seek after a row in the `order_by_instant` total order. Honoured only
664 /// by `query_notes_filtered_count_free` with `offset: 0`.
665 #[serde(default)]
666 pub after_instant: Option<NoteInstantSeekAfter>,
667}
668
669/// Temporal-referential note CRUD over the notes substrate table.
670#[async_trait]
671pub trait NoteStore: Send + Sync + 'static {
672 async fn get_live_notes_by_key(
673 &self,
674 _namespace: &str,
675 _key: &str,
676 _kind: Option<&str>,
677 ) -> StorageResult<Vec<Note>> {
678 Err(crate::StorageError::Unsupported {
679 capability: crate::StorageCapability::Notes,
680 operation: "get_live_notes_by_key".into(),
681 message: "backend has no keyed note lookup".into(),
682 })
683 }
684
685 async fn query_keyed_notes(
686 &self,
687 _namespace: &str,
688 _filter: &NoteFilter,
689 _prefix: &str,
690 _after: Option<&NoteKeyCursor>,
691 _page: PageRequest,
692 ) -> StorageResult<(Vec<Note>, Option<NoteKeyCursor>)> {
693 Err(crate::StorageError::Unsupported {
694 capability: crate::StorageCapability::Notes,
695 operation: "query_keyed_notes".into(),
696 message: "backend has no keyed note pagination".into(),
697 })
698 }
699
700 /// Insert or update a single note. Updates preserve the stored immutable key.
701 async fn upsert_note(&self, note: Note) -> StorageResult<()>;
702 /// Replace a note only when the persisted row still matches the caller's
703 /// read snapshot.
704 ///
705 /// `expected_updated_at` is the snapshot revision and
706 /// `expected_deleted_at` closes the soft-delete race (legacy soft-delete
707 /// paths may change `deleted_at` without changing `updated_at`). The
708 /// replacement note's `updated_at` must be strictly greater than that
709 /// persisted revision. Returns `false` when the row disappeared, changed,
710 /// or was supplied a non-advancing replacement revision. This is the
711 /// full-note compare-and-swap seam used when a pack hook derives coupled
712 /// fields from that snapshot before persistence. The default returns
713 /// `Unsupported` rather than falling back to an unguarded upsert and
714 /// reintroducing the stale-snapshot race.
715 async fn replace_note_if_unchanged(
716 &self,
717 _note: Note,
718 _expected_updated_at: i64,
719 _expected_deleted_at: Option<i64>,
720 ) -> StorageResult<bool> {
721 Err(crate::StorageError::Unsupported {
722 capability: crate::StorageCapability::Notes,
723 operation: "replace_note_if_unchanged".into(),
724 message: "this backend does not implement guarded note replacement".into(),
725 })
726 }
727 /// Insert a note only if no row already holds its id, reporting whether
728 /// this call is the one that inserted it.
729 ///
730 /// This closes the other half of the read-modify-write race that
731 /// [`NoteStore::replace_note_if_unchanged`] closes. That one protects a
732 /// caller who read an existing row; this one protects a caller who read
733 /// *no* row. Where the id is derived rather than freshly generated — a
734 /// deterministic per-subject id — two callers can both read absence and
735 /// both write, and an upsert resolves that by overwriting, so the first
736 /// caller's write is lost with no error on either side. Returns `false`
737 /// when a row already existed, which the caller maps to a conflict rather
738 /// than to success.
739 ///
740 /// The pre-existing row is left exactly as it is: this must not be
741 /// implemented as an upsert, since overwriting is the behaviour being
742 /// avoided. The default returns `Unsupported` for that reason, rather
743 /// than falling back to `upsert_note` and silently reintroducing the
744 /// race under a name that promises otherwise.
745 async fn insert_note_if_absent(&self, _note: Note) -> StorageResult<bool> {
746 Err(crate::StorageError::Unsupported {
747 capability: crate::StorageCapability::Notes,
748 operation: "insert_note_if_absent".into(),
749 message: "this backend does not implement guarded note insertion".into(),
750 })
751 }
752 /// Insert or update a batch of notes. Updates preserve each stored immutable key.
753 async fn upsert_notes(&self, notes: Vec<Note>) -> StorageResult<BatchWriteSummary>;
754 /// Fetch a note by UUID, returning `None` if absent.
755 async fn get_note(&self, id: Uuid) -> StorageResult<Option<Note>>;
756 /// Fetch a note by UUID regardless of soft-deletion state.
757 ///
758 /// Returns the note row even when `deleted_at` is set. Callers use this
759 /// to distinguish "soft-deleted" from "never existed".
760 async fn get_note_including_deleted(&self, id: Uuid) -> StorageResult<Option<Note>>;
761 /// Delete a note by UUID using the specified delete mode.
762 async fn delete_note(&self, id: Uuid, mode: DeleteMode) -> StorageResult<bool>;
763 /// Patch `properties`/`updated_at` on an existing note in place via a real
764 /// `UPDATE`, leaving every other column (including the row's `rowid`)
765 /// untouched.
766 ///
767 /// Unlike `upsert_note`, which writes the complete note shape, this leaves
768 /// every non-property column untouched. It also never churns the row's
769 /// implicit `rowid`, which is required by callers relying on stable row
770 /// identity (#780).
771 /// Returns `true` when a live (non-soft-deleted) row with this `id` was
772 /// found and updated, `false` otherwise.
773 async fn update_note_properties(
774 &self,
775 id: Uuid,
776 properties: Option<Value>,
777 updated_at: i64,
778 ) -> StorageResult<bool>;
779 /// Atomically set one top-level key in a note's JSON `properties` object.
780 ///
781 /// The backend must perform the read/modify/write as one storage operation
782 /// so concurrent writes to different keys cannot overwrite each other.
783 /// `value` keeps its JSON type, including explicit JSON `null`. A SQL-NULL
784 /// property document is initialized as an empty object. A live row whose
785 /// stored document is a non-object is not modified and returns `false`, as
786 /// do missing and soft-deleted rows. Keys containing U+0000 must be
787 /// rejected: SQLite JSON-path labels cannot address them without risking
788 /// mutation of a shorter sibling key.
789 async fn set_note_property(
790 &self,
791 id: Uuid,
792 key: &str,
793 value: Value,
794 updated_at: i64,
795 ) -> StorageResult<bool>;
796 /// Atomically patch a single `properties` JSON key on a note, but only
797 /// when the row's *current* state (re-evaluated inside this same
798 /// statement, not a snapshot the caller fetched earlier) still satisfies
799 /// `filter`'s namespace/kind/property_filters.
800 ///
801 /// Unlike `set_note_property` (which patches unconditionally once the row
802 /// is live) or `update_note_properties` (which replaces the whole
803 /// `properties` column with a value the caller already computed — safe
804 /// only when nothing else can have written to the row since the caller's
805 /// read), this also rechecks `filter` against the row's live state before
806 /// writing, so a target that stopped matching an eligibility predicate
807 /// between validation and this call is not mutated. Any other property
808 /// written concurrently between the caller's read and this call survives
809 /// untouched either way. A live row whose stored `properties` document is
810 /// a non-object (scalar, array, or otherwise) is not modified and returns
811 /// `false`, mirroring `set_note_property`. Returns `Ok(false)` — not an
812 /// error — when no live row currently matches `filter` (id not found,
813 /// soft-deleted, an eligibility property changed since the caller last
814 /// validated it, or the stored document is not a JSON object); the
815 /// caller degrades that the same way as `update_note_properties`'s
816 /// `Ok(false)`.
817 async fn try_patch_note_property(
818 &self,
819 id: Uuid,
820 namespace: &str,
821 filter: &NoteFilter,
822 json_path: &str,
823 value: Value,
824 updated_at: i64,
825 ) -> StorageResult<bool>;
826 /// Atomically patch one JSON property on every supplied note.
827 ///
828 /// Each target is rechecked against `namespace` and `filter` inside the
829 /// same transaction. The operation commits only when every distinct id
830 /// matches exactly one live object-valued row; a missing, soft-deleted, or
831 /// no-longer-eligible target rolls the entire unit back. Other property
832 /// keys are preserved by the same storage-side `json_set` operation as
833 /// [`Self::try_patch_note_property`]. Backends without a transactional
834 /// multi-note implementation retain the default `Unsupported` result.
835 async fn patch_note_property_atomic(
836 &self,
837 _ids: Vec<Uuid>,
838 _namespace: &str,
839 _filter: &NoteFilter,
840 _json_path: &str,
841 _value: Value,
842 _updated_at: i64,
843 ) -> StorageResult<()> {
844 Err(crate::StorageError::Unsupported {
845 capability: crate::StorageCapability::Notes,
846 operation: "patch_note_property_atomic".into(),
847 message: "this backend does not implement atomic multi-note property patches".into(),
848 })
849 }
850 /// Query notes by namespace and optional kind with pagination.
851 /// The returned total and page items must come from one consistent
852 /// backend snapshot.
853 async fn query_notes(
854 &self,
855 namespace: &str,
856 kind: Option<&str>,
857 page: PageRequest,
858 ) -> StorageResult<Page<Note>>;
859 /// Query notes by namespace and optional kind without computing an exact total.
860 ///
861 /// Implementations must return `total: None`. Backends that cannot
862 /// guarantee a count-free projection must fail closed rather than invoke
863 /// [`Self::query_notes`], whose exact total remains available to callers
864 /// that expose a snapshot-consistent count.
865 async fn query_notes_count_free(
866 &self,
867 _namespace: &str,
868 _kind: Option<&str>,
869 _page: PageRequest,
870 ) -> StorageResult<Page<Note>> {
871 Err(crate::StorageError::Unsupported {
872 capability: crate::StorageCapability::Notes,
873 operation: "query_notes_count_free".into(),
874 message: "this backend does not implement count-free note paging".into(),
875 })
876 }
877 /// Query notes with property-based filtering and custom sort.
878 /// The returned total and page items must come from one consistent
879 /// backend snapshot.
880 async fn query_notes_filtered(
881 &self,
882 namespace: &str,
883 filter: &NoteFilter,
884 page: PageRequest,
885 ) -> StorageResult<Page<Note>>;
886 /// Query a filtered note page without computing an exact total.
887 ///
888 /// Implementations must return `total: None`; callers that need a
889 /// `has_more` bit should request one lookahead row. Backends that cannot
890 /// guarantee a count-free projection must fail closed rather than invoke
891 /// [`Self::query_notes_filtered`], whose exact total is intentionally
892 /// backlog-proportional.
893 async fn query_notes_filtered_count_free(
894 &self,
895 _namespace: &str,
896 _filter: &NoteFilter,
897 _page: PageRequest,
898 ) -> StorageResult<Page<Note>> {
899 Err(crate::StorageError::Unsupported {
900 capability: crate::StorageCapability::Notes,
901 operation: "query_notes_filtered_count_free".into(),
902 message: "this backend does not implement count-free filtered note paging".into(),
903 })
904 }
905 /// Count several filtered note populations in one consistent backend
906 /// snapshot. Backends that cannot provide that guarantee must return
907 /// [`crate::StorageError::Unsupported`] rather than composing independent
908 /// queries or treating an absent page total as zero. SQL backends should
909 /// override this operation so callers can retain separate index-friendly
910 /// predicates without racing between their counts.
911 async fn count_notes_filtered_in_snapshot(
912 &self,
913 _namespace: &str,
914 _filters: &[NoteFilter],
915 ) -> StorageResult<Vec<u64>> {
916 Err(crate::StorageError::Unsupported {
917 capability: crate::StorageCapability::Notes,
918 operation: "count_notes_filtered_in_snapshot".into(),
919 message: "this backend does not implement snapshot-consistent filtered note counts"
920 .into(),
921 })
922 }
923 /// Count several filtered note populations in one backend snapshot, with
924 /// each count bounded by `cap`.
925 ///
926 /// A saturated entry reports `count == cap` and `saturated == true`; an
927 /// unsaturated entry is exact. SQL implementations should count over a
928 /// limited `SELECT 1` subquery so work is proportional to `cap` and the
929 /// caller's own matching population, never an unrelated backlog. That
930 /// bound holds only when every `property_filter` in the given
931 /// `NoteFilter` is servable as an index key column (equality on a bound
932 /// parameter or a plan-time-provable partial-index predicate) all the way
933 /// down to the filter that actually narrows the population — a residual
934 /// filter evaluated after the index seek (a bound parameter the index
935 /// cannot use to skip rows) degrades the bound to work proportional to
936 /// the rows *scanned* before `cap` matches are found, not the rows that
937 /// match. Callers building filters against this bound must ensure the
938 /// full predicate set is indexed, not just the predicate that narrows the
939 /// population most. Backends that cannot preserve the shared snapshot or
940 /// work bound must fail closed.
941 async fn count_notes_filtered_bounded_in_snapshot(
942 &self,
943 _namespace: &str,
944 _filters: &[NoteFilter],
945 _cap: u32,
946 ) -> StorageResult<Vec<BoundedCount>> {
947 Err(crate::StorageError::Unsupported {
948 capability: crate::StorageCapability::Notes,
949 operation: "count_notes_filtered_bounded_in_snapshot".into(),
950 message: "this backend does not implement snapshot-consistent bounded note counts"
951 .into(),
952 })
953 }
954 /// Resolve a note id to its immutable insertion sequence.
955 async fn note_sequence(&self, _id: Uuid) -> StorageResult<Option<i64>> {
956 Err(crate::StorageError::Unsupported {
957 capability: crate::StorageCapability::Notes,
958 operation: "note_sequence".into(),
959 message: "this backend does not implement note insertion sequences".into(),
960 })
961 }
962 /// Query an immutable insertion-sequence keyset page with the same
963 /// predicates as [`Self::query_notes_filtered`].
964 async fn query_notes_filtered_after(
965 &self,
966 _namespace: &str,
967 _filter: &NoteFilter,
968 _after: Option<SeekCursor>,
969 _limit: u32,
970 ) -> StorageResult<SeekPage<Note>> {
971 Err(crate::StorageError::Unsupported {
972 capability: crate::StorageCapability::Notes,
973 operation: "query_notes_filtered_after".into(),
974 message: "this backend does not implement note seek pagination".into(),
975 })
976 }
977 /// Fetch up to `max_rows + 1` notes matching `filter` in a single
978 /// deterministically-ordered SQL statement, with no separate `COUNT(*)`
979 /// and no pagination loop.
980 ///
981 /// A single statement observes one consistent snapshot for its entire
982 /// execution, so the result cannot be split across a concurrent insert
983 /// the way a `COUNT(*)` followed by independent `LIMIT`/`OFFSET` pages
984 /// can. Callers detect the over-bound case by checking whether the
985 /// returned `Vec` has more than `max_rows` items — that means at least
986 /// `max_rows + 1` rows matched and the caller must reject the query
987 /// rather than silently return a truncated, possibly priority-incomplete
988 /// set.
989 async fn query_notes_filtered_bounded(
990 &self,
991 namespace: &str,
992 filter: &NoteFilter,
993 max_rows: u32,
994 ) -> StorageResult<Vec<Note>>;
995 /// Count notes in a namespace, optionally filtered by kind.
996 async fn count_notes(&self, namespace: &str, kind: Option<&str>) -> StorageResult<u64>;
997 /// Count notes across the given namespaces, optionally filtered by kind.
998 /// The default preserves compatibility by summing the existing
999 /// single-namespace operation; SQL backends should override this with one
1000 /// `IN` aggregate.
1001 async fn count_notes_in_namespaces(
1002 &self,
1003 namespaces: &[String],
1004 kind: Option<&str>,
1005 ) -> StorageResult<u64> {
1006 let mut total = 0;
1007 for namespace in namespaces {
1008 total += self.count_notes(namespace, kind).await?;
1009 }
1010 Ok(total)
1011 }
1012
1013 /// Attempt to insert a note without overwriting an existing row.
1014 ///
1015 /// Returns `true` when the row was newly written. Returns `false` only
1016 /// when a live note with the same non-empty `external_id` already exists in
1017 /// the same namespace and kind (confirmed dedup hit). Any other constraint
1018 /// violation (e.g. a primary key collision) is surfaced as a `StorageError`
1019 /// so that callers do not misinterpret unexpected failures as deduplication.
1020 async fn try_insert_note(&self, note: Note) -> StorageResult<bool>;
1021
1022 /// Fetch multiple notes by UUID in a single call.
1023 async fn get_notes_batch(&self, ids: &[Uuid]) -> StorageResult<Vec<Note>> {
1024 let mut out = Vec::with_capacity(ids.len());
1025 for &id in ids {
1026 if let Some(n) = self.get_note(id).await? {
1027 out.push(n);
1028 }
1029 }
1030 Ok(out)
1031 }
1032}