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, SqlValue, 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)]
191pub struct EventFilter {
192 pub ids: Vec<Uuid>,
193 #[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 #[serde(default)]
208 pub outcome: Option<EventOutcome>,
209 #[serde(default)]
211 pub payload_equalities: Vec<(String, SqlValue)>,
212}
213
214impl EventFilter {
215 pub fn outcome(mut self, outcome: EventOutcome) -> Self {
217 self.outcome = Some(outcome);
218 self
219 }
220
221 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#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
236pub struct EventOrderKey {
237 pub created_at_us: i64,
238 pub physical_id: String,
239}
240
241#[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#[derive(Clone, Debug, Serialize, Deserialize)]
255pub struct EventPageRow {
256 pub event: Event,
257 pub order_key: EventOrderKey,
258}
259
260#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
277#[serde(rename_all = "snake_case")]
278pub enum EventAppendDisposition {
279 Inserted,
281 AlreadyPresentIdentical,
284 IdentityConflict,
287}
288
289#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
292pub struct IdempotentEventBatchResult {
293 pub rows: Vec<EventAppendDisposition>,
294}
295
296#[async_trait]
298pub trait EventStore: Send + Sync + 'static {
299 async fn append_event(&self, event: Event) -> StorageResult<()>;
301 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
303 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
305 async fn query_events(
307 &self,
308 filter: EventFilter,
309 page: PageRequest,
310 ) -> StorageResult<Page<Event>>;
311 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 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
325
326 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 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 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}