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(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
210#[serde(rename_all = "snake_case")]
211pub enum EventAppendDisposition {
212 Inserted,
214 AlreadyPresentIdentical,
217 IdentityConflict,
220}
221
222#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
225pub struct IdempotentEventBatchResult {
226 pub rows: Vec<EventAppendDisposition>,
227}
228
229#[async_trait]
231pub trait EventStore: Send + Sync + 'static {
232 async fn append_event(&self, event: Event) -> StorageResult<()>;
234 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
236 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
238 async fn query_events(
240 &self,
241 filter: EventFilter,
242 page: PageRequest,
243 ) -> StorageResult<Page<Event>>;
244 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
246
247 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 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 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}