Skip to main content

khive_storage/
event.rs

1//! Event storage capability — append-only operation log.
2
3use async_trait::async_trait;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use uuid::Uuid;
7
8use khive_types::{EventKind, EventOutcome, RefResolution, SubstrateKind};
9
10use crate::capability::StorageCapability;
11use crate::error::StorageError;
12use crate::types::{BatchWriteSummary, Page, PageRequest, SqlValue, StorageResult};
13
14/// Storage-level event record. Every verb execution produces one.
15/// Immutable once appended; projection rows are written beside it at append time.
16#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
17pub struct Event {
18    pub id: Uuid,
19    pub namespace: String,
20    pub verb: String,
21    pub substrate: SubstrateKind,
22    pub actor: String,
23    pub kind: EventKind,
24    pub outcome: EventOutcome,
25    pub payload: Value,
26    pub payload_schema_version: u32,
27    pub profile_state_version: Option<u64>,
28    pub duration_us: i64,
29    pub target_id: Option<Uuid>,
30    pub session_id: Option<Uuid>,
31    pub aggregate_kind: Option<String>,
32    pub aggregate_id: Option<Uuid>,
33    pub created_at: i64,
34    /// Original request parser position; absent for legacy/non-request events.
35    #[serde(default)]
36    pub op_index: Option<u32>,
37    /// Reference provenance, present exactly when `op_index` is present.
38    #[serde(default)]
39    pub ref_resolution: Option<RefResolution>,
40}
41
42impl Event {
43    /// Create a new event with a generated UUID and current timestamp.
44    pub fn new(
45        namespace: impl Into<String>,
46        verb: impl Into<String>,
47        kind: EventKind,
48        substrate: SubstrateKind,
49        actor: impl Into<String>,
50    ) -> Self {
51        let operation = crate::operation_context::current_operation_attribution();
52        Self {
53            id: Uuid::new_v4(),
54            namespace: namespace.into(),
55            verb: verb.into(),
56            substrate,
57            actor: actor.into(),
58            kind,
59            outcome: EventOutcome::Success,
60            payload: Value::Object(Default::default()),
61            payload_schema_version: 1,
62            profile_state_version: None,
63            duration_us: 0,
64            target_id: None,
65            session_id: None,
66            aggregate_kind: None,
67            aggregate_id: None,
68            created_at: chrono::Utc::now().timestamp_micros(),
69            op_index: operation.map(|operation| operation.op_index),
70            ref_resolution: operation.map(|operation| operation.ref_resolution),
71        }
72    }
73
74    /// Set the event outcome (success/failure).
75    pub fn with_outcome(mut self, o: EventOutcome) -> Self {
76        self.outcome = o;
77        self
78    }
79
80    /// Set the event payload JSON.
81    pub fn with_payload(mut self, payload: Value) -> Self {
82        self.payload = payload;
83        self
84    }
85
86    /// Set the payload schema version for forward compatibility.
87    pub fn with_payload_schema_version(mut self, version: u32) -> Self {
88        self.payload_schema_version = version;
89        self
90    }
91
92    /// Set the brain profile state version at event time.
93    pub fn with_profile_state_version(mut self, version: u64) -> Self {
94        self.profile_state_version = Some(version);
95        self
96    }
97
98    /// Set the operation duration in microseconds.
99    pub fn with_duration_us(mut self, us: i64) -> Self {
100        self.duration_us = us;
101        self
102    }
103
104    /// Set the exact subject UUID for this event; it need not be a graph record.
105    pub fn with_target(mut self, id: Uuid) -> Self {
106        self.target_id = Some(id);
107        self
108    }
109
110    /// Set the session ID for correlating related events.
111    pub fn with_session_id(mut self, id: Uuid) -> Self {
112        self.session_id = Some(id);
113        self
114    }
115
116    /// Set the aggregate kind and ID for event-sourced projections.
117    pub fn with_aggregate(mut self, kind: impl Into<String>, id: Uuid) -> Self {
118        self.aggregate_kind = Some(kind.into());
119        self.aggregate_id = Some(id);
120        self
121    }
122}
123
124/// Which durable record family an event observation refers to. Edge remains
125/// relational storage rather than a fourth [`SubstrateKind`]; this narrower
126/// discriminant exists so edge lifecycle events can still be queried through
127/// `observed=[edge_id]`.
128#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
129#[serde(rename_all = "snake_case")]
130pub enum ReferentKind {
131    Entity,
132    Note,
133    Edge,
134}
135
136impl ReferentKind {
137    /// Return the lowercase string name for this referent kind.
138    pub const fn name(self) -> &'static str {
139        match self {
140            Self::Entity => "entity",
141            Self::Note => "note",
142            Self::Edge => "edge",
143        }
144    }
145}
146
147/// Role of a referent in a brain observation (candidate, selected, target, signal).
148#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
149#[serde(rename_all = "snake_case")]
150pub enum ObservationRole {
151    Candidate,
152    Selected,
153    Target,
154    Signal,
155}
156
157impl ObservationRole {
158    /// Return the lowercase string name for this observation role.
159    pub const fn name(self) -> &'static str {
160        match self {
161            Self::Candidate => "candidate",
162            Self::Selected => "selected",
163            Self::Target => "target",
164            Self::Signal => "signal",
165        }
166    }
167}
168
169/// A single entity observation recorded alongside an event.
170#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
171pub struct EventObservation {
172    pub event_id: Uuid,
173    pub entity_id: Uuid,
174    pub referent_kind: ReferentKind,
175    pub role: ObservationRole,
176    pub position: u32,
177}
178
179/// An event together with its associated observations.
180#[derive(Debug, Clone, Serialize, Deserialize)]
181pub struct EventView {
182    pub event: Event,
183    pub observations: Vec<EventObservation>,
184}
185
186/// Filter for querying events. Namespace is implicit in the scoped EventStore.
187///
188/// New fields deserialize to defaults for older serialized filters. Complete
189/// external Rust struct literals must specify them or use `..Default::default()`.
190#[derive(Clone, Debug, Default, Serialize, Deserialize)]
191pub struct EventFilter {
192    pub ids: Vec<Uuid>,
193    /// Exact event subject, independent of the graph observation projection.
194    #[serde(default)]
195    pub target_id: Option<Uuid>,
196    pub kinds: Vec<EventKind>,
197    pub verbs: Vec<String>,
198    pub substrates: Vec<SubstrateKind>,
199    pub actors: Vec<String>,
200    pub after: Option<i64>,
201    pub before: Option<i64>,
202    pub session_id: Option<Uuid>,
203    pub observed: Vec<Uuid>,
204    pub selected: Vec<Uuid>,
205    pub payload_proposal_id: Option<Uuid>,
206    /// Exact outcome; no restriction when absent.
207    #[serde(default)]
208    pub outcome: Option<EventOutcome>,
209    /// ANDed SQL JSON equalities; see [`EventFilter::payload_eq`].
210    #[serde(default)]
211    pub payload_equalities: Vec<(String, SqlValue)>,
212}
213
214impl EventFilter {
215    /// Match one outcome, replacing any previous outcome selection.
216    pub fn outcome(mut self, outcome: EventOutcome) -> Self {
217        self.outcome = Some(outcome);
218        self
219    }
220
221    /// Require SQL `json_extract(payload, path) = value` equality.
222    ///
223    /// Paths must be `$.field[.subfield]` with nonempty ASCII alphanumeric or
224    /// underscore segments; queries reject other paths. Calls are ANDed.
225    /// SQL NULL matches neither a missing field nor explicit JSON null.
226    /// JSON booleans compare as 0/1, including numeric coercion. JSON objects
227    /// and arrays compare their serialized SQL text, not structural JSON equality.
228    pub fn payload_eq(mut self, path: impl Into<String>, value: SqlValue) -> Self {
229        self.payload_equalities.push((path.into(), value));
230        self
231    }
232}
233
234/// Persisted ordering key; the physical UUID spelling is not normalized.
235#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
236pub struct EventOrderKey {
237    pub created_at_us: i64,
238    pub physical_id: String,
239}
240
241/// Count-free ascending window within one scoped event store.
242#[derive(Clone, Debug, Serialize, Deserialize)]
243pub struct EventPageQuery {
244    pub since_us: i64,
245    pub until_us: i64,
246    pub kinds: Vec<EventKind>,
247    pub actors: Vec<String>,
248    pub exclude_namespaces: Vec<String>,
249    pub after: Option<EventOrderKey>,
250    pub max_rows: u32,
251}
252
253/// A decoded event with its exact persisted seek key.
254#[derive(Clone, Debug, Serialize, Deserialize)]
255pub struct EventPageRow {
256    pub event: Event,
257    pub order_key: EventOrderKey,
258}
259
260/// A bounded prefix, without a count of the complete matching window.
261///
262/// `budget_stop` names the row the leaf did not return because reading it would
263/// pass the leaf's raw-text budget; `rows` then holds every row ordered before it,
264/// whole, and is shorter than the requested row bound.
265#[derive(Clone, Debug, Serialize, Deserialize)]
266pub struct EventPageWindow {
267    pub rows: Vec<EventPageRow>,
268    #[serde(default)]
269    pub budget_stop: Option<EventOrderKey>,
270}
271
272/// Per-row outcome of an [`EventStore::append_events_idempotent`] call, in
273/// input order. Distinguishes a fresh insert from a retry that reproduced an
274/// identical row from a retry whose identity now disagrees with what is
275/// already stored.
276#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
277#[serde(rename_all = "snake_case")]
278pub enum EventAppendDisposition {
279    /// No prior row with this id existed; it was inserted.
280    Inserted,
281    /// A prior row with this id existed and every persisted column plus the
282    /// ordered observation projection matched exactly. Not re-inserted.
283    AlreadyPresentIdentical,
284    /// A prior row with this id existed but disagreed with the submitted
285    /// row. Not inserted; unrelated rows in the same batch are unaffected.
286    IdentityConflict,
287}
288
289/// Result of [`EventStore::append_events_idempotent`]. `rows` preserves the
290/// input order and length of the submitted batch.
291#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
292pub struct IdempotentEventBatchResult {
293    pub rows: Vec<EventAppendDisposition>,
294}
295
296/// Append-only operation log for verb executions.
297#[async_trait]
298pub trait EventStore: Send + Sync + 'static {
299    /// Append a single event to the log.
300    async fn append_event(&self, event: Event) -> StorageResult<()>;
301    /// Append a batch of events to the log.
302    async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
303    /// Fetch an event by UUID, returning `None` if absent.
304    async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
305    /// Query events matching a filter with pagination.
306    async fn query_events(
307        &self,
308        filter: EventFilter,
309        page: PageRequest,
310    ) -> StorageResult<Page<Event>>;
311    /// Read at most 4096 rows in ascending physical time/ID order, without counting.
312    /// Filters and exclusions apply before the row bound. A backend must refuse
313    /// unsupported paging rather than substitute an offset query.
314    async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
315        let _ = query;
316        Err(StorageError::Unsupported {
317            capability: StorageCapability::Events,
318            operation: "query_event_page".into(),
319            message: "this EventStore backend does not implement query_event_page".into(),
320        })
321    }
322
323    /// Count events matching a filter.
324    async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
325
326    /// Validate `event` against the exact insert/observation shape the
327    /// backend would build at append time, performing no I/O. Rejects a
328    /// malformed row before it is ever enqueued for a write, so one bad
329    /// producer input cannot poison a batch shared with other callers.
330    ///
331    /// Backends that do not implement pre-enqueue validation return
332    /// [`StorageError::Unsupported`].
333    fn preflight_event(&self, event: &Event) -> StorageResult<()> {
334        let _ = event;
335        Err(StorageError::Unsupported {
336            capability: StorageCapability::Events,
337            operation: "preflight_event".into(),
338            message: "this EventStore backend does not implement preflight_event".into(),
339        })
340    }
341
342    /// Append a batch of events with idempotent retry semantics: a row
343    /// carrying an id that already exists is compared against every
344    /// persisted column and its ordered observation projection rather than
345    /// treated as a write conflict. Exact equality reports
346    /// [`EventAppendDisposition::AlreadyPresentIdentical`] instead of
347    /// re-inserting; any mismatch reports
348    /// [`EventAppendDisposition::IdentityConflict`] for that row alone,
349    /// while unrelated rows in the same batch still commit.
350    ///
351    /// Backends that do not implement idempotent batching return
352    /// [`StorageError::Unsupported`].
353    async fn append_events_idempotent(
354        &self,
355        events: Vec<Event>,
356    ) -> StorageResult<IdempotentEventBatchResult> {
357        let _ = events;
358        Err(StorageError::Unsupported {
359            capability: StorageCapability::Events,
360            operation: "append_events_idempotent".into(),
361            message: "this EventStore backend does not implement append_events_idempotent".into(),
362        })
363    }
364
365    /// Whether this backend implements `preflight_event` and
366    /// `append_events_idempotent` for real, rather than inheriting their
367    /// `Unsupported`-returning defaults above.
368    ///
369    /// A caller that builds an ADR-133 audit-batch seam over a backend that
370    /// answers `false` here would have every audited row rejected at
371    /// preflight while the dispatch it audits still reports success — the
372    /// exact silent-loss failure mode the batch exists to prevent. Defaults
373    /// to `false` so an unmodified legacy backend is caught at registry
374    /// build time instead of appearing healthy.
375    fn supports_idempotent_audit_batch(&self) -> bool {
376        false
377    }
378}
379
380#[cfg(test)]
381mod tests {
382    use super::*;
383
384    #[test]
385    fn event_filter_target_id_roundtrips_and_defaults_for_legacy_frames() {
386        let target = Uuid::new_v4();
387        let filter = EventFilter {
388            target_id: Some(target),
389            ..EventFilter::default()
390        };
391        let mut wire = serde_json::to_value(&filter).unwrap();
392        assert_eq!(wire["target_id"], target.to_string());
393        let decoded: EventFilter = serde_json::from_value(wire.clone()).unwrap();
394        assert_eq!(decoded.target_id, Some(target));
395        wire.as_object_mut().unwrap().remove("target_id");
396        let legacy: EventFilter = serde_json::from_value(wire).unwrap();
397        assert_eq!(legacy.target_id, None);
398    }
399
400    #[test]
401    fn event_page_window_budget_stop_roundtrips_and_defaults_when_absent() {
402        let stop = EventOrderKey {
403            created_at_us: 7,
404            physical_id: Uuid::from_u128(9).to_string(),
405        };
406        let window = EventPageWindow {
407            rows: Vec::new(),
408            budget_stop: Some(stop.clone()),
409        };
410        let mut wire = serde_json::to_value(&window).unwrap();
411        let decoded: EventPageWindow = serde_json::from_value(wire.clone()).unwrap();
412        assert_eq!(decoded.budget_stop, Some(stop));
413        wire.as_object_mut().unwrap().remove("budget_stop");
414        let legacy: EventPageWindow = serde_json::from_value(wire).unwrap();
415        assert_eq!(legacy.budget_stop, None);
416    }
417}