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}