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, 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#[derive(Clone, Debug, Default, Serialize, Deserialize)]
188pub struct EventFilter {
189    pub ids: Vec<Uuid>,
190    /// Exact event subject, independent of the graph observation projection.
191    #[serde(default)]
192    pub target_id: Option<Uuid>,
193    pub kinds: Vec<EventKind>,
194    pub verbs: Vec<String>,
195    pub substrates: Vec<SubstrateKind>,
196    pub actors: Vec<String>,
197    pub after: Option<i64>,
198    pub before: Option<i64>,
199    pub session_id: Option<Uuid>,
200    pub observed: Vec<Uuid>,
201    pub selected: Vec<Uuid>,
202    pub payload_proposal_id: Option<Uuid>,
203}
204
205/// Persisted ordering key; the physical UUID spelling is not normalized.
206#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
207pub struct EventOrderKey {
208    pub created_at_us: i64,
209    pub physical_id: String,
210}
211
212/// Count-free ascending window within one scoped event store.
213#[derive(Clone, Debug, Serialize, Deserialize)]
214pub struct EventPageQuery {
215    pub since_us: i64,
216    pub until_us: i64,
217    pub kinds: Vec<EventKind>,
218    pub actors: Vec<String>,
219    pub exclude_namespaces: Vec<String>,
220    pub after: Option<EventOrderKey>,
221    pub max_rows: u32,
222}
223
224/// A decoded event with its exact persisted seek key.
225#[derive(Clone, Debug, Serialize, Deserialize)]
226pub struct EventPageRow {
227    pub event: Event,
228    pub order_key: EventOrderKey,
229}
230
231/// A bounded prefix, without a count of the complete matching window.
232///
233/// `budget_stop` names the row the leaf did not return because reading it would
234/// pass the leaf's raw-text budget; `rows` then holds every row ordered before it,
235/// whole, and is shorter than the requested row bound.
236#[derive(Clone, Debug, Serialize, Deserialize)]
237pub struct EventPageWindow {
238    pub rows: Vec<EventPageRow>,
239    #[serde(default)]
240    pub budget_stop: Option<EventOrderKey>,
241}
242
243/// Per-row outcome of an [`EventStore::append_events_idempotent`] call, in
244/// input order. Distinguishes a fresh insert from a retry that reproduced an
245/// identical row from a retry whose identity now disagrees with what is
246/// already stored.
247#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
248#[serde(rename_all = "snake_case")]
249pub enum EventAppendDisposition {
250    /// No prior row with this id existed; it was inserted.
251    Inserted,
252    /// A prior row with this id existed and every persisted column plus the
253    /// ordered observation projection matched exactly. Not re-inserted.
254    AlreadyPresentIdentical,
255    /// A prior row with this id existed but disagreed with the submitted
256    /// row. Not inserted; unrelated rows in the same batch are unaffected.
257    IdentityConflict,
258}
259
260/// Result of [`EventStore::append_events_idempotent`]. `rows` preserves the
261/// input order and length of the submitted batch.
262#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
263pub struct IdempotentEventBatchResult {
264    pub rows: Vec<EventAppendDisposition>,
265}
266
267/// Append-only operation log for verb executions.
268#[async_trait]
269pub trait EventStore: Send + Sync + 'static {
270    /// Append a single event to the log.
271    async fn append_event(&self, event: Event) -> StorageResult<()>;
272    /// Append a batch of events to the log.
273    async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
274    /// Fetch an event by UUID, returning `None` if absent.
275    async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
276    /// Query events matching a filter with pagination.
277    async fn query_events(
278        &self,
279        filter: EventFilter,
280        page: PageRequest,
281    ) -> StorageResult<Page<Event>>;
282    /// Read at most 4096 rows in ascending physical time/ID order, without counting.
283    /// Filters and exclusions apply before the row bound. A backend must refuse
284    /// unsupported paging rather than substitute an offset query.
285    async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
286        let _ = query;
287        Err(StorageError::Unsupported {
288            capability: StorageCapability::Events,
289            operation: "query_event_page".into(),
290            message: "this EventStore backend does not implement query_event_page".into(),
291        })
292    }
293
294    /// Count events matching a filter.
295    async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
296
297    /// Validate `event` against the exact insert/observation shape the
298    /// backend would build at append time, performing no I/O. Rejects a
299    /// malformed row before it is ever enqueued for a write, so one bad
300    /// producer input cannot poison a batch shared with other callers.
301    ///
302    /// Backends that do not implement pre-enqueue validation return
303    /// [`StorageError::Unsupported`].
304    fn preflight_event(&self, event: &Event) -> StorageResult<()> {
305        let _ = event;
306        Err(StorageError::Unsupported {
307            capability: StorageCapability::Events,
308            operation: "preflight_event".into(),
309            message: "this EventStore backend does not implement preflight_event".into(),
310        })
311    }
312
313    /// Append a batch of events with idempotent retry semantics: a row
314    /// carrying an id that already exists is compared against every
315    /// persisted column and its ordered observation projection rather than
316    /// treated as a write conflict. Exact equality reports
317    /// [`EventAppendDisposition::AlreadyPresentIdentical`] instead of
318    /// re-inserting; any mismatch reports
319    /// [`EventAppendDisposition::IdentityConflict`] for that row alone,
320    /// while unrelated rows in the same batch still commit.
321    ///
322    /// Backends that do not implement idempotent batching return
323    /// [`StorageError::Unsupported`].
324    async fn append_events_idempotent(
325        &self,
326        events: Vec<Event>,
327    ) -> StorageResult<IdempotentEventBatchResult> {
328        let _ = events;
329        Err(StorageError::Unsupported {
330            capability: StorageCapability::Events,
331            operation: "append_events_idempotent".into(),
332            message: "this EventStore backend does not implement append_events_idempotent".into(),
333        })
334    }
335
336    /// Whether this backend implements `preflight_event` and
337    /// `append_events_idempotent` for real, rather than inheriting their
338    /// `Unsupported`-returning defaults above.
339    ///
340    /// A caller that builds an ADR-133 audit-batch seam over a backend that
341    /// answers `false` here would have every audited row rejected at
342    /// preflight while the dispatch it audits still reports success — the
343    /// exact silent-loss failure mode the batch exists to prevent. Defaults
344    /// to `false` so an unmodified legacy backend is caught at registry
345    /// build time instead of appearing healthy.
346    fn supports_idempotent_audit_batch(&self) -> bool {
347        false
348    }
349}
350
351#[cfg(test)]
352mod tests {
353    use super::*;
354
355    #[test]
356    fn event_filter_target_id_roundtrips_and_defaults_for_legacy_frames() {
357        let target = Uuid::new_v4();
358        let filter = EventFilter {
359            target_id: Some(target),
360            ..EventFilter::default()
361        };
362        let mut wire = serde_json::to_value(&filter).unwrap();
363        assert_eq!(wire["target_id"], target.to_string());
364        let decoded: EventFilter = serde_json::from_value(wire.clone()).unwrap();
365        assert_eq!(decoded.target_id, Some(target));
366        wire.as_object_mut().unwrap().remove("target_id");
367        let legacy: EventFilter = serde_json::from_value(wire).unwrap();
368        assert_eq!(legacy.target_id, None);
369    }
370
371    #[test]
372    fn event_page_window_budget_stop_roundtrips_and_defaults_when_absent() {
373        let stop = EventOrderKey {
374            created_at_us: 7,
375            physical_id: Uuid::from_u128(9).to_string(),
376        };
377        let window = EventPageWindow {
378            rows: Vec::new(),
379            budget_stop: Some(stop.clone()),
380        };
381        let mut wire = serde_json::to_value(&window).unwrap();
382        let decoded: EventPageWindow = serde_json::from_value(wire.clone()).unwrap();
383        assert_eq!(decoded.budget_stop, Some(stop));
384        wire.as_object_mut().unwrap().remove("budget_stop");
385        let legacy: EventPageWindow = serde_json::from_value(wire).unwrap();
386        assert_eq!(legacy.budget_stop, None);
387    }
388}