1use 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#[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 #[serde(default)]
36 pub op_index: Option<u32>,
37 #[serde(default)]
39 pub ref_resolution: Option<RefResolution>,
40}
41
42impl Event {
43 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 pub fn with_outcome(mut self, o: EventOutcome) -> Self {
76 self.outcome = o;
77 self
78 }
79
80 pub fn with_payload(mut self, payload: Value) -> Self {
82 self.payload = payload;
83 self
84 }
85
86 pub fn with_payload_schema_version(mut self, version: u32) -> Self {
88 self.payload_schema_version = version;
89 self
90 }
91
92 pub fn with_profile_state_version(mut self, version: u64) -> Self {
94 self.profile_state_version = Some(version);
95 self
96 }
97
98 pub fn with_duration_us(mut self, us: i64) -> Self {
100 self.duration_us = us;
101 self
102 }
103
104 pub fn with_target(mut self, id: Uuid) -> Self {
106 self.target_id = Some(id);
107 self
108 }
109
110 pub fn with_session_id(mut self, id: Uuid) -> Self {
112 self.session_id = Some(id);
113 self
114 }
115
116 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#[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 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#[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 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#[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#[derive(Debug, Clone, Serialize, Deserialize)]
181pub struct EventView {
182 pub event: Event,
183 pub observations: Vec<EventObservation>,
184}
185
186#[derive(Clone, Debug, Default, Serialize, Deserialize)]
188pub struct EventFilter {
189 pub ids: Vec<Uuid>,
190 #[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#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
207pub struct EventOrderKey {
208 pub created_at_us: i64,
209 pub physical_id: String,
210}
211
212#[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#[derive(Clone, Debug, Serialize, Deserialize)]
226pub struct EventPageRow {
227 pub event: Event,
228 pub order_key: EventOrderKey,
229}
230
231#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
248#[serde(rename_all = "snake_case")]
249pub enum EventAppendDisposition {
250 Inserted,
252 AlreadyPresentIdentical,
255 IdentityConflict,
258}
259
260#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
263pub struct IdempotentEventBatchResult {
264 pub rows: Vec<EventAppendDisposition>,
265}
266
267#[async_trait]
269pub trait EventStore: Send + Sync + 'static {
270 async fn append_event(&self, event: Event) -> StorageResult<()>;
272 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
274 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
276 async fn query_events(
278 &self,
279 filter: EventFilter,
280 page: PageRequest,
281 ) -> StorageResult<Page<Event>>;
282 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 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
296
297 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 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 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}