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/// Per-row outcome of an [`EventStore::append_events_idempotent`] call, in
206/// input order. Distinguishes a fresh insert from a retry that reproduced an
207/// identical row from a retry whose identity now disagrees with what is
208/// already stored.
209#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
210#[serde(rename_all = "snake_case")]
211pub enum EventAppendDisposition {
212    /// No prior row with this id existed; it was inserted.
213    Inserted,
214    /// A prior row with this id existed and every persisted column plus the
215    /// ordered observation projection matched exactly. Not re-inserted.
216    AlreadyPresentIdentical,
217    /// A prior row with this id existed but disagreed with the submitted
218    /// row. Not inserted; unrelated rows in the same batch are unaffected.
219    IdentityConflict,
220}
221
222/// Result of [`EventStore::append_events_idempotent`]. `rows` preserves the
223/// input order and length of the submitted batch.
224#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
225pub struct IdempotentEventBatchResult {
226    pub rows: Vec<EventAppendDisposition>,
227}
228
229/// Append-only operation log for verb executions.
230#[async_trait]
231pub trait EventStore: Send + Sync + 'static {
232    /// Append a single event to the log.
233    async fn append_event(&self, event: Event) -> StorageResult<()>;
234    /// Append a batch of events to the log.
235    async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
236    /// Fetch an event by UUID, returning `None` if absent.
237    async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
238    /// Query events matching a filter with pagination.
239    async fn query_events(
240        &self,
241        filter: EventFilter,
242        page: PageRequest,
243    ) -> StorageResult<Page<Event>>;
244    /// Count events matching a filter.
245    async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
246
247    /// Validate `event` against the exact insert/observation shape the
248    /// backend would build at append time, performing no I/O. Rejects a
249    /// malformed row before it is ever enqueued for a write, so one bad
250    /// producer input cannot poison a batch shared with other callers.
251    ///
252    /// Backends that do not implement pre-enqueue validation return
253    /// [`StorageError::Unsupported`].
254    fn preflight_event(&self, event: &Event) -> StorageResult<()> {
255        let _ = event;
256        Err(StorageError::Unsupported {
257            capability: StorageCapability::Events,
258            operation: "preflight_event".into(),
259            message: "this EventStore backend does not implement preflight_event".into(),
260        })
261    }
262
263    /// Append a batch of events with idempotent retry semantics: a row
264    /// carrying an id that already exists is compared against every
265    /// persisted column and its ordered observation projection rather than
266    /// treated as a write conflict. Exact equality reports
267    /// [`EventAppendDisposition::AlreadyPresentIdentical`] instead of
268    /// re-inserting; any mismatch reports
269    /// [`EventAppendDisposition::IdentityConflict`] for that row alone,
270    /// while unrelated rows in the same batch still commit.
271    ///
272    /// Backends that do not implement idempotent batching return
273    /// [`StorageError::Unsupported`].
274    async fn append_events_idempotent(
275        &self,
276        events: Vec<Event>,
277    ) -> StorageResult<IdempotentEventBatchResult> {
278        let _ = events;
279        Err(StorageError::Unsupported {
280            capability: StorageCapability::Events,
281            operation: "append_events_idempotent".into(),
282            message: "this EventStore backend does not implement append_events_idempotent".into(),
283        })
284    }
285
286    /// Whether this backend implements `preflight_event` and
287    /// `append_events_idempotent` for real, rather than inheriting their
288    /// `Unsupported`-returning defaults above.
289    ///
290    /// A caller that builds an ADR-133 audit-batch seam over a backend that
291    /// answers `false` here would have every audited row rejected at
292    /// preflight while the dispatch it audits still reports success — the
293    /// exact silent-loss failure mode the batch exists to prevent. Defaults
294    /// to `false` so an unmodified legacy backend is caught at registry
295    /// build time instead of appearing healthy.
296    fn supports_idempotent_audit_batch(&self) -> bool {
297        false
298    }
299}
300
301#[cfg(test)]
302mod tests {
303    use super::*;
304
305    #[test]
306    fn event_filter_target_id_roundtrips_and_defaults_for_legacy_frames() {
307        let target = Uuid::new_v4();
308        let filter = EventFilter {
309            target_id: Some(target),
310            ..EventFilter::default()
311        };
312        let mut wire = serde_json::to_value(&filter).unwrap();
313        assert_eq!(wire["target_id"], target.to_string());
314        let decoded: EventFilter = serde_json::from_value(wire.clone()).unwrap();
315        assert_eq!(decoded.target_id, Some(target));
316        wire.as_object_mut().unwrap().remove("target_id");
317        let legacy: EventFilter = serde_json::from_value(wire).unwrap();
318        assert_eq!(legacy.target_id, None);
319    }
320}