khive_runtime/
event_store_guard.rs1use std::sync::Arc;
10
11use async_trait::async_trait;
12use khive_storage::event::{EventPageQuery, EventPageWindow, IdempotentEventBatchResult};
13use khive_storage::{
14 BatchWriteSummary, Event, EventFilter, EventStore, Page, PageRequest, StorageResult,
15};
16use uuid::Uuid;
17
18use crate::NamespaceToken;
19
20#[derive(Clone, Debug, Eq, PartialEq)]
28pub struct EventAttribution {
29 namespace: String,
30 actor: String,
31 operation: Option<khive_types::OperationAttribution>,
32}
33
34impl EventAttribution {
35 pub fn from_token(token: &NamespaceToken) -> Self {
38 Self {
39 namespace: token.namespace().as_str().to_owned(),
40 actor: format!("{}:{}", token.actor().kind, token.actor().id),
41 operation: khive_storage::operation_context::current_operation_attribution(),
42 }
43 }
44
45 pub fn stamp(&self, mut event: Event) -> Event {
47 event.namespace.clone_from(&self.namespace);
48 event.actor.clone_from(&self.actor);
49 if let Some(operation) = self.operation {
52 event.op_index = Some(operation.op_index);
53 event.ref_resolution = Some(operation.ref_resolution);
54 }
55 event
56 }
57}
58
59pub(crate) struct AttributedEventStore {
60 inner: Arc<dyn EventStore>,
61 attribution: EventAttribution,
62}
63
64impl AttributedEventStore {
65 pub(crate) fn wrap(inner: Arc<dyn EventStore>, token: &NamespaceToken) -> Arc<dyn EventStore> {
66 let mut attribution = EventAttribution::from_token(token);
67 attribution.operation = None;
71 Arc::new(Self { inner, attribution })
72 }
73
74 fn attribute(&self, event: Event) -> Event {
75 self.attribution.stamp(event)
76 }
77
78 fn attribute_many(&self, events: Vec<Event>) -> Vec<Event> {
79 events
80 .into_iter()
81 .map(|event| self.attribute(event))
82 .collect()
83 }
84}
85
86#[async_trait]
87impl EventStore for AttributedEventStore {
88 async fn append_event(&self, event: Event) -> StorageResult<()> {
89 self.inner.append_event(self.attribute(event)).await
90 }
91
92 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
93 self.inner.append_events(self.attribute_many(events)).await
94 }
95
96 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
97 self.inner.get_event(id).await
98 }
99
100 async fn query_events(
101 &self,
102 filter: EventFilter,
103 page: PageRequest,
104 ) -> StorageResult<Page<Event>> {
105 self.inner.query_events(filter, page).await
106 }
107
108 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
109 self.inner.count_events(filter).await
110 }
111
112 async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
113 self.inner.query_event_page(query).await
114 }
115
116 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
117 self.inner.preflight_event(&self.attribute(event.clone()))
118 }
119
120 async fn append_events_idempotent(
121 &self,
122 events: Vec<Event>,
123 ) -> StorageResult<IdempotentEventBatchResult> {
124 self.inner
125 .append_events_idempotent(self.attribute_many(events))
126 .await
127 }
128
129 fn supports_idempotent_audit_batch(&self) -> bool {
130 self.inner.supports_idempotent_audit_batch()
131 }
132}
133
134#[cfg(test)]
135mod operation_tests {
136 use super::*;
137 use khive_storage::operation_context::scope_operation_attribution;
138 use khive_types::{EventKind, OperationAttribution, RefResolution, SubstrateKind};
139
140 fn event() -> Event {
141 Event::new(
142 "local",
143 "test.operation",
144 EventKind::Audit,
145 SubstrateKind::Event,
146 "fixture",
147 )
148 }
149
150 #[tokio::test]
151 async fn reusable_store_does_not_reuse_operation_but_explicit_capture_survives_defer() {
152 let runtime = crate::KhiveRuntime::memory().unwrap();
153 let token = runtime.authorize(crate::Namespace::local()).unwrap();
154 let operation = OperationAttribution {
155 op_index: 2,
156 ref_resolution: RefResolution::Resolved,
157 };
158 let (store, captured) = scope_operation_attribution(operation, async {
159 (
160 runtime.events(&token).unwrap(),
161 EventAttribution::from_token(&token),
162 )
163 })
164 .await;
165
166 let outside = event();
167 store.append_event(outside.clone()).await.unwrap();
168 let stored = store.get_event(outside.id).await.unwrap().unwrap();
169 assert_eq!((stored.op_index, stored.ref_resolution), (None, None));
170
171 let deferred = tokio::spawn(async move { captured.stamp(event()) })
172 .await
173 .unwrap();
174 assert_eq!(
175 (deferred.op_index, deferred.ref_resolution),
176 (Some(2), Some(RefResolution::Resolved))
177 );
178 store.append_event(deferred.clone()).await.unwrap();
179 let stored = store.get_event(deferred.id).await.unwrap().unwrap();
180 assert_eq!(
181 (stored.op_index, stored.ref_resolution),
182 (Some(2), Some(RefResolution::Resolved))
183 );
184
185 let next = scope_operation_attribution(
186 OperationAttribution {
187 op_index: 7,
188 ref_resolution: RefResolution::Literal,
189 },
190 async { event() },
191 )
192 .await;
193 store.append_event(next.clone()).await.unwrap();
194 let stored = store.get_event(next.id).await.unwrap().unwrap();
195 assert_eq!(
196 (stored.op_index, stored.ref_resolution),
197 (Some(7), Some(RefResolution::Literal))
198 );
199 }
200}