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, 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}
35
36impl Event {
37    /// Create a new event with a generated UUID and current timestamp.
38    pub fn new(
39        namespace: impl Into<String>,
40        verb: impl Into<String>,
41        kind: EventKind,
42        substrate: SubstrateKind,
43        actor: impl Into<String>,
44    ) -> Self {
45        Self {
46            id: Uuid::new_v4(),
47            namespace: namespace.into(),
48            verb: verb.into(),
49            substrate,
50            actor: actor.into(),
51            kind,
52            outcome: EventOutcome::Success,
53            payload: Value::Object(Default::default()),
54            payload_schema_version: 1,
55            profile_state_version: None,
56            duration_us: 0,
57            target_id: None,
58            session_id: None,
59            aggregate_kind: None,
60            aggregate_id: None,
61            created_at: chrono::Utc::now().timestamp_micros(),
62        }
63    }
64
65    /// Set the event outcome (success/failure).
66    pub fn with_outcome(mut self, o: EventOutcome) -> Self {
67        self.outcome = o;
68        self
69    }
70
71    /// Set the event payload JSON.
72    pub fn with_payload(mut self, payload: Value) -> Self {
73        self.payload = payload;
74        self
75    }
76
77    /// Set the payload schema version for forward compatibility.
78    pub fn with_payload_schema_version(mut self, version: u32) -> Self {
79        self.payload_schema_version = version;
80        self
81    }
82
83    /// Set the brain profile state version at event time.
84    pub fn with_profile_state_version(mut self, version: u64) -> Self {
85        self.profile_state_version = Some(version);
86        self
87    }
88
89    /// Set the operation duration in microseconds.
90    pub fn with_duration_us(mut self, us: i64) -> Self {
91        self.duration_us = us;
92        self
93    }
94
95    /// Set the target entity/note ID for this event.
96    pub fn with_target(mut self, id: Uuid) -> Self {
97        self.target_id = Some(id);
98        self
99    }
100
101    /// Set the session ID for correlating related events.
102    pub fn with_session_id(mut self, id: Uuid) -> Self {
103        self.session_id = Some(id);
104        self
105    }
106
107    /// Set the aggregate kind and ID for event-sourced projections.
108    pub fn with_aggregate(mut self, kind: impl Into<String>, id: Uuid) -> Self {
109        self.aggregate_kind = Some(kind.into());
110        self.aggregate_id = Some(id);
111        self
112    }
113}
114
115/// Which substrate (entity or note) the referent record lives in.
116#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
117#[serde(rename_all = "snake_case")]
118pub enum ReferentKind {
119    Entity,
120    Note,
121}
122
123impl ReferentKind {
124    /// Return the lowercase string name for this referent kind.
125    pub const fn name(self) -> &'static str {
126        match self {
127            Self::Entity => "entity",
128            Self::Note => "note",
129        }
130    }
131}
132
133/// Role of a referent in a brain observation (candidate, selected, target, signal).
134#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
135#[serde(rename_all = "snake_case")]
136pub enum ObservationRole {
137    Candidate,
138    Selected,
139    Target,
140    Signal,
141}
142
143impl ObservationRole {
144    /// Return the lowercase string name for this observation role.
145    pub const fn name(self) -> &'static str {
146        match self {
147            Self::Candidate => "candidate",
148            Self::Selected => "selected",
149            Self::Target => "target",
150            Self::Signal => "signal",
151        }
152    }
153}
154
155/// A single entity observation recorded alongside an event.
156#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
157pub struct EventObservation {
158    pub event_id: Uuid,
159    pub entity_id: Uuid,
160    pub referent_kind: ReferentKind,
161    pub role: ObservationRole,
162    pub position: u32,
163}
164
165/// An event together with its associated observations.
166#[derive(Debug, Clone, Serialize, Deserialize)]
167pub struct EventView {
168    pub event: Event,
169    pub observations: Vec<EventObservation>,
170}
171
172/// Filter for querying events. Namespace is implicit in the scoped EventStore.
173#[derive(Clone, Debug, Default, Serialize, Deserialize)]
174pub struct EventFilter {
175    pub ids: Vec<Uuid>,
176    pub kinds: Vec<EventKind>,
177    pub verbs: Vec<String>,
178    pub substrates: Vec<SubstrateKind>,
179    pub actors: Vec<String>,
180    pub after: Option<i64>,
181    pub before: Option<i64>,
182    pub session_id: Option<Uuid>,
183    pub observed: Vec<Uuid>,
184    pub selected: Vec<Uuid>,
185    pub payload_proposal_id: Option<Uuid>,
186}
187
188/// Per-row outcome of an [`EventStore::append_events_idempotent`] call, in
189/// input order. Distinguishes a fresh insert from a retry that reproduced an
190/// identical row from a retry whose identity now disagrees with what is
191/// already stored.
192#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
193#[serde(rename_all = "snake_case")]
194pub enum EventAppendDisposition {
195    /// No prior row with this id existed; it was inserted.
196    Inserted,
197    /// A prior row with this id existed and every persisted column plus the
198    /// ordered observation projection matched exactly. Not re-inserted.
199    AlreadyPresentIdentical,
200    /// A prior row with this id existed but disagreed with the submitted
201    /// row. Not inserted; unrelated rows in the same batch are unaffected.
202    IdentityConflict,
203}
204
205/// Result of [`EventStore::append_events_idempotent`]. `rows` preserves the
206/// input order and length of the submitted batch.
207#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
208pub struct IdempotentEventBatchResult {
209    pub rows: Vec<EventAppendDisposition>,
210}
211
212/// Append-only operation log for verb executions.
213#[async_trait]
214pub trait EventStore: Send + Sync + 'static {
215    /// Append a single event to the log.
216    async fn append_event(&self, event: Event) -> StorageResult<()>;
217    /// Append a batch of events to the log.
218    async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
219    /// Fetch an event by UUID, returning `None` if absent.
220    async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
221    /// Query events matching a filter with pagination.
222    async fn query_events(
223        &self,
224        filter: EventFilter,
225        page: PageRequest,
226    ) -> StorageResult<Page<Event>>;
227    /// Count events matching a filter.
228    async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
229
230    /// Validate `event` against the exact insert/observation shape the
231    /// backend would build at append time, performing no I/O. Rejects a
232    /// malformed row before it is ever enqueued for a write, so one bad
233    /// producer input cannot poison a batch shared with other callers.
234    ///
235    /// Backends that do not implement pre-enqueue validation return
236    /// [`StorageError::Unsupported`].
237    fn preflight_event(&self, event: &Event) -> StorageResult<()> {
238        let _ = event;
239        Err(StorageError::Unsupported {
240            capability: StorageCapability::Events,
241            operation: "preflight_event".into(),
242            message: "this EventStore backend does not implement preflight_event".into(),
243        })
244    }
245
246    /// Append a batch of events with idempotent retry semantics: a row
247    /// carrying an id that already exists is compared against every
248    /// persisted column and its ordered observation projection rather than
249    /// treated as a write conflict. Exact equality reports
250    /// [`EventAppendDisposition::AlreadyPresentIdentical`] instead of
251    /// re-inserting; any mismatch reports
252    /// [`EventAppendDisposition::IdentityConflict`] for that row alone,
253    /// while unrelated rows in the same batch still commit.
254    ///
255    /// Backends that do not implement idempotent batching return
256    /// [`StorageError::Unsupported`].
257    async fn append_events_idempotent(
258        &self,
259        events: Vec<Event>,
260    ) -> StorageResult<IdempotentEventBatchResult> {
261        let _ = events;
262        Err(StorageError::Unsupported {
263            capability: StorageCapability::Events,
264            operation: "append_events_idempotent".into(),
265            message: "this EventStore backend does not implement append_events_idempotent".into(),
266        })
267    }
268
269    /// Whether this backend implements `preflight_event` and
270    /// `append_events_idempotent` for real, rather than inheriting their
271    /// `Unsupported`-returning defaults above.
272    ///
273    /// A caller that builds an ADR-133 audit-batch seam over a backend that
274    /// answers `false` here would have every audited row rejected at
275    /// preflight while the dispatch it audits still reports success — the
276    /// exact silent-loss failure mode the batch exists to prevent. Defaults
277    /// to `false` so an unmodified legacy backend is caught at registry
278    /// build time instead of appearing healthy.
279    fn supports_idempotent_audit_batch(&self) -> bool {
280        false
281    }
282}