Skip to main content

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