Skip to main content

khive_storage/
entity.rs

1//! Entity storage capability — graph node CRUD.
2
3use async_trait::async_trait;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use uuid::Uuid;
7
8use crate::attachment::Attachment;
9use crate::types::{
10    BatchWriteSummary, DeleteMode, Page, PageRequest, SeekCursor, SeekPage, SqlValue, StorageResult,
11};
12
13/// Storage-level entity record. Flat SQL-friendly representation.
14/// Maps to the `entities` substrate table.
15#[derive(Clone, Debug, Serialize, Deserialize)]
16pub struct Entity {
17    pub id: Uuid,
18    pub namespace: String,
19    pub kind: String,
20    /// Pack-governed subtype token. Maps to `entities.entity_type` column.
21    pub entity_type: Option<String>,
22    pub name: String,
23    pub description: Option<String>,
24    pub properties: Option<Value>,
25    pub tags: Vec<String>,
26    pub created_at: i64,
27    pub updated_at: i64,
28    /// Persisted row revision: inserts start at one; each committed update advances once.
29    /// Input to a snapshot replacement is the revision read before normalization.
30    #[serde(default = "initial_entity_version")]
31    pub version: i64,
32    pub deleted_at: Option<i64>,
33    /// When this entity was tombstoned by a merge, the `into` entity's ID.
34    pub merged_into: Option<Uuid>,
35    /// Opaque event ID for the merge that tombstoned this entity.
36    pub merge_event_id: Option<Uuid>,
37    /// Read-only compatibility projection of attachment role `"content"`.
38    ///
39    /// Entity writes ignore this field. Callers publish content through the
40    /// attachment substrate; reads populate it so existing response payloads
41    /// keep their `content_ref` field during the coordinated cutover.
42    pub content_ref: Option<String>,
43}
44
45fn initial_entity_version() -> i64 {
46    1
47}
48
49impl Entity {
50    /// Create a new entity with a generated UUID and current timestamp.
51    pub fn new(
52        namespace: impl Into<String>,
53        kind: impl Into<String>,
54        name: impl Into<String>,
55    ) -> Self {
56        let now = chrono::Utc::now().timestamp_micros();
57        Self::minimal(Uuid::new_v4(), namespace, kind, name, now, now)
58    }
59
60    /// Create an entity with supplied identity and microsecond timestamps.
61    /// Optional payload, revision, and tombstone fields use the same defaults as `new`.
62    pub fn minimal(
63        id: Uuid,
64        namespace: impl Into<String>,
65        kind: impl Into<String>,
66        name: impl Into<String>,
67        created_at: i64,
68        updated_at: i64,
69    ) -> Self {
70        Self {
71            id,
72            namespace: namespace.into(),
73            kind: kind.into(),
74            entity_type: None,
75            name: name.into(),
76            description: None,
77            properties: None,
78            tags: Vec::new(),
79            created_at,
80            updated_at,
81            version: 1,
82            deleted_at: None,
83            merged_into: None,
84            merge_event_id: None,
85            content_ref: None,
86        }
87    }
88
89    /// Set the pack-governed entity subtype token.
90    pub fn with_entity_type(mut self, t: Option<impl Into<String>>) -> Self {
91        self.entity_type = t.map(Into::into);
92        self
93    }
94
95    /// Set the entity description.
96    pub fn with_description(mut self, d: impl Into<String>) -> Self {
97        self.description = Some(d.into());
98        self
99    }
100
101    /// Set the entity properties JSON blob.
102    pub fn with_properties(mut self, p: Value) -> Self {
103        self.properties = Some(p);
104        self
105    }
106
107    /// Set the entity tags.
108    pub fn with_tags(mut self, t: Vec<String>) -> Self {
109        self.tags = t;
110        self
111    }
112}
113
114/// Entity liveness selected by an [`EntityFilter`].
115#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
116#[serde(rename_all = "snake_case")]
117pub enum EntityTombstones {
118    /// Only rows without a deletion timestamp.
119    #[default]
120    Live,
121    /// Both live and tombstoned rows.
122    All,
123    /// Only rows with a deletion timestamp.
124    Only,
125}
126
127/// Entity filter for query operations.
128///
129/// New fields deserialize to defaults for older serialized filters. Complete
130/// external Rust struct literals must specify them or use `..Default::default()`.
131#[derive(Clone, Debug, Default, Serialize, Deserialize)]
132pub struct EntityFilter {
133    pub ids: Vec<Uuid>,
134    pub kinds: Vec<String>,
135    /// ANDed SQL JSON equalities; see [`EntityFilter::property_eq`].
136    #[serde(default)]
137    pub property_equalities: Vec<(String, SqlValue)>,
138    /// Live-only by default. Liveness builders replace this selection.
139    #[serde(default)]
140    pub tombstones: EntityTombstones,
141    /// Filter by exact `entity_type` value. Multiple values are ORed.
142    pub entity_types: Vec<String>,
143    /// Kind-qualified accepted subtype spellings. Groups are ORed; each group
144    /// requires its kind and one of its values. An empty value group matches nothing.
145    #[serde(default)]
146    pub entity_types_by_kind: std::collections::BTreeMap<String, Vec<String>>,
147    /// For entity listing, fall back to a string `properties.type` only when
148    /// `entity_type` is null. Does not change the returned entity or apply when
149    /// both type filters are empty. Other query callers retain exact-column filtering.
150    #[serde(default)]
151    pub legacy_entity_type_fallback: bool,
152    pub name_prefix: Option<String>,
153    /// Deterministic, case-sensitive equality on `entities.name` (binary
154    /// comparison — SQLite's default collation for `=` on a `TEXT` column
155    /// without an explicit `COLLATE NOCASE`). Distinct from `name_prefix`:
156    /// that stage's `LIKE` is inherently prefix-shaped and, with SQLite's
157    /// default `NOCASE`-free `LIKE` on ASCII, still ranks a page by
158    /// `created_at DESC` — a match that is exact but not the newest can be
159    /// paged out. `name_exact` skips paging risk entirely by filtering to
160    /// only rows that equal `name` at the SQL layer.
161    pub name_exact: Option<String>,
162    pub tags_any: Vec<String>,
163    /// When non-empty, restricts results to any of these namespaces using
164    /// `namespace IN (...)`. Takes precedence over the `namespace` string
165    /// parameter passed to `query_entities` / `count_entities`. When empty the
166    /// caller-supplied `namespace` parameter is used (single-namespace path,
167    /// backward-compatible default).
168    #[serde(default)]
169    pub namespaces: Vec<String>,
170    /// ASCII-case-insensitive batched exact-name match (ADR-104 Stage C).
171    /// Compares a caller-bounded set of raw and ASCII-lowercased candidate
172    /// strings to `LOWER(name)`. Cased non-ASCII characters require exact form.
173    /// Distinct from single-value, case-sensitive `name_exact`. Results contain
174    /// at most one representative row per folded candidate before page limits
175    /// and offsets are applied.
176    /// Implementations may omit the page total to keep this lookup page-limited
177    /// instead of issuing a separate count.
178    #[serde(default)]
179    pub names_ci: Vec<String>,
180}
181
182impl EntityFilter {
183    /// Require SQL `json_extract(properties, path) = value` equality.
184    ///
185    /// Paths must be `$.field[.subfield]` with nonempty ASCII alphanumeric or
186    /// underscore segments; queries reject other paths. Calls are ANDed.
187    /// SQL NULL matches neither a missing field nor explicit JSON null.
188    /// JSON booleans compare as 0/1, including numeric coercion. JSON objects
189    /// and arrays compare their serialized SQL text, not structural JSON equality.
190    pub fn property_eq(mut self, path: impl Into<String>, value: SqlValue) -> Self {
191        self.property_equalities.push((path.into(), value));
192        self
193    }
194
195    /// Include live and tombstoned rows, replacing a previous liveness choice.
196    pub fn include_tombstones(mut self) -> Self {
197        self.tombstones = EntityTombstones::All;
198        self
199    }
200
201    /// Include only tombstoned rows, replacing a previous liveness choice.
202    pub fn tombstoned_only(mut self) -> Self {
203        self.tombstones = EntityTombstones::Only;
204        self
205    }
206}
207
208/// Exact nullable stored entity types and their live row counts.
209pub type EntityTypeCounts = Vec<(Option<String>, u64)>;
210
211/// Entity CRUD operations over the entities substrate table.
212#[async_trait]
213pub trait EntityStore: Send + Sync + 'static {
214    /// Insert at version one or update the current row and advance its version once.
215    /// Incoming version values do not override the stored counter. Raw SQL replacement
216    /// of an existing entity is forbidden: use a conflict UPDATE or the typed store.
217    async fn upsert_entity(&self, entity: Entity) -> StorageResult<()>;
218    /// Insert an entity only when no row with its id or another conflicting
219    /// key exists. Returns `true` when this call inserted the row and `false`
220    /// when an existing row won the race. The existing row is never updated.
221    ///
222    /// The default returns `Unsupported` rather than falling back to
223    /// [`EntityStore::upsert_entity`], because an upsert would overwrite the
224    /// winning row and violate this method's conditional-insert contract.
225    async fn insert_entity_if_absent(&self, _entity: Entity) -> StorageResult<bool> {
226        Err(crate::StorageError::Unsupported {
227            capability: crate::StorageCapability::Entities,
228            operation: "insert_entity_if_absent".into(),
229            message: "this backend does not implement conditional entity insert".into(),
230        })
231    }
232    /// Atomically insert/update an entity and all supplied attachment roles.
233    async fn upsert_entity_with_attachments(
234        &self,
235        _entity: Entity,
236        _attachments: Vec<Attachment>,
237    ) -> StorageResult<()> {
238        Err(crate::StorageError::Unsupported {
239            capability: crate::StorageCapability::Attachments,
240            operation: "upsert_entity_with_attachments".into(),
241            message: "this backend does not implement atomic entity attachment publication".into(),
242        })
243    }
244    /// Insert or update a batch of entities.
245    async fn upsert_entities(&self, entities: Vec<Entity>) -> StorageResult<BatchWriteSummary>;
246    /// Replace an entity only when the persisted row still matches the
247    /// caller's read snapshot.
248    ///
249    /// `entity.version` is the persisted revision read with the snapshot;
250    /// a successful replacement increments the stored version once.
251    /// `expected_updated_at` is an additional snapshot timestamp guard and
252    /// `expected_deleted_at` closes the soft-delete race. The replacement
253    /// entity's `updated_at` must be strictly greater than that persisted
254    /// revision. Returns `false` when the row disappeared, changed, or was
255    /// supplied a non-advancing replacement revision. This is the full-entity
256    /// compare-and-swap seam used when a caller derives coupled fields from
257    /// that snapshot before persistence — mirrors
258    /// [`crate::NoteStore::replace_note_if_unchanged`]. The default returns
259    /// `Unsupported` rather than falling back to an unguarded upsert and
260    /// reintroducing the stale-snapshot race.
261    async fn replace_entity_if_unchanged(
262        &self,
263        _entity: Entity,
264        _expected_updated_at: i64,
265        _expected_deleted_at: Option<i64>,
266    ) -> StorageResult<bool> {
267        Err(crate::StorageError::Unsupported {
268            capability: crate::StorageCapability::Entities,
269            operation: "replace_entity_if_unchanged".into(),
270            message: "this backend does not implement guarded entity replacement".into(),
271        })
272    }
273    /// Fetch an entity by UUID, returning `None` if absent.
274    async fn get_entity(&self, id: Uuid) -> StorageResult<Option<Entity>>;
275    /// Delete an entity by UUID using the specified delete mode.
276    async fn delete_entity(&self, id: Uuid, mode: DeleteMode) -> StorageResult<bool>;
277    /// Query entities by namespace with filter and pagination.
278    async fn query_entities(
279        &self,
280        namespace: &str,
281        filter: EntityFilter,
282        page: PageRequest,
283    ) -> StorageResult<Page<Entity>>;
284    /// Query an offset page without computing an exact total.
285    ///
286    /// The returned page has `total: None`. Callers needing a has-more signal
287    /// request one extra row. Backends must implement this operation directly;
288    /// falling back to `query_entities` could perform a discarded full count.
289    async fn query_entities_count_free(
290        &self,
291        _namespace: &str,
292        _filter: EntityFilter,
293        _page: PageRequest,
294    ) -> StorageResult<Page<Entity>> {
295        Err(crate::StorageError::Unsupported {
296            capability: crate::StorageCapability::Entities,
297            operation: "query_entities_count_free".into(),
298            message: "this backend does not implement count-free entity pages".into(),
299        })
300    }
301    /// Resolve an entity id to its immutable insertion sequence.
302    async fn entity_sequence(&self, _id: Uuid) -> StorageResult<Option<i64>> {
303        Err(crate::StorageError::Unsupported {
304            capability: crate::StorageCapability::Entities,
305            operation: "entity_sequence".into(),
306            message: "this backend does not implement entity insertion sequences".into(),
307        })
308    }
309    /// Query an immutable insertion-sequence keyset page.
310    ///
311    /// Backends that do not implement seek pagination may retain the default
312    /// unsupported result; callers must not silently fall back to offset
313    /// paging because that would weaken the no-gap/no-duplicate contract.
314    async fn query_entities_after(
315        &self,
316        _namespace: &str,
317        _filter: EntityFilter,
318        _after: Option<SeekCursor>,
319        _limit: u32,
320    ) -> StorageResult<SeekPage<Entity>> {
321        Err(crate::StorageError::Unsupported {
322            capability: crate::StorageCapability::Entities,
323            operation: "query_entities_after".into(),
324            message: "this backend does not implement entity seek pagination".into(),
325        })
326    }
327    /// Count entities in a namespace matching the given filter.
328    async fn count_entities(&self, namespace: &str, filter: EntityFilter) -> StorageResult<u64>;
329    /// Report complete, disjoint live entity counts by exact stored type across
330    /// the supplied namespace set, from one backend read snapshot. Duplicate
331    /// namespaces count once; an empty set produces no groups. SQL NULL is
332    /// `None`, distinct from every string label. Ordering is unspecified.
333    ///
334    /// `Ok(None)` means this reporting capability is unavailable. It is not an
335    /// empty report or a storage failure; callers may retain their scalar count
336    /// and omit the breakdown. Errors from an implemented report propagate.
337    async fn count_entities_by_type(
338        &self,
339        _namespaces: &[String],
340    ) -> StorageResult<Option<EntityTypeCounts>> {
341        Ok(None)
342    }
343    /// Fetch an entity by UUID regardless of soft-deletion state.
344    ///
345    /// Returns the entity row even when `deleted_at` is set. Callers use this
346    /// to distinguish "soft-deleted" from "never existed".
347    async fn get_entity_including_deleted(&self, id: Uuid) -> StorageResult<Option<Entity>>;
348}