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(¬e).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}