1use 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#[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 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 pub fn with_outcome(mut self, o: EventOutcome) -> Self {
67 self.outcome = o;
68 self
69 }
70
71 pub fn with_payload(mut self, payload: Value) -> Self {
73 self.payload = payload;
74 self
75 }
76
77 pub fn with_payload_schema_version(mut self, version: u32) -> Self {
79 self.payload_schema_version = version;
80 self
81 }
82
83 pub fn with_profile_state_version(mut self, version: u64) -> Self {
85 self.profile_state_version = Some(version);
86 self
87 }
88
89 pub fn with_duration_us(mut self, us: i64) -> Self {
91 self.duration_us = us;
92 self
93 }
94
95 pub fn with_target(mut self, id: Uuid) -> Self {
97 self.target_id = Some(id);
98 self
99 }
100
101 pub fn with_session_id(mut self, id: Uuid) -> Self {
103 self.session_id = Some(id);
104 self
105 }
106
107 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#[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 pub const fn name(self) -> &'static str {
126 match self {
127 Self::Entity => "entity",
128 Self::Note => "note",
129 }
130 }
131}
132
133#[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 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#[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#[derive(Debug, Clone, Serialize, Deserialize)]
167pub struct EventView {
168 pub event: Event,
169 pub observations: Vec<EventObservation>,
170}
171
172#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
193#[serde(rename_all = "snake_case")]
194pub enum EventAppendDisposition {
195 Inserted,
197 AlreadyPresentIdentical,
200 IdentityConflict,
203}
204
205#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
208pub struct IdempotentEventBatchResult {
209 pub rows: Vec<EventAppendDisposition>,
210}
211
212#[async_trait]
214pub trait EventStore: Send + Sync + 'static {
215 async fn append_event(&self, event: Event) -> StorageResult<()>;
217 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary>;
219 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>>;
221 async fn query_events(
223 &self,
224 filter: EventFilter,
225 page: PageRequest,
226 ) -> StorageResult<Page<Event>>;
227 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64>;
229
230 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 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 fn supports_idempotent_audit_batch(&self) -> bool {
280 false
281 }
282}