aion_server/observability/
instrumented_store.rs1use std::sync::Arc;
4use std::time::Instant;
5
6use aion_core::{Event, TimerId, WorkflowFilter, WorkflowId, WorkflowSummary};
7use aion_store::{
8 EventStore, PackageRecord, PackageRouteRecord, PackageStore, ReadableEventStore, RunSummary,
9 StoreError, TimerEntry, WritableEventStore, WriteToken,
10};
11use async_trait::async_trait;
12use chrono::{DateTime, Utc};
13
14use super::metrics::Metrics;
15
16pub struct InstrumentedEventStore {
18 inner: Arc<dyn EventStore>,
19 metrics: Metrics,
20 namespace: String,
21}
22
23impl InstrumentedEventStore {
24 #[must_use]
26 pub fn new(inner: Arc<dyn EventStore>, metrics: Metrics, namespace: impl Into<String>) -> Self {
27 Self {
28 inner,
29 metrics,
30 namespace: namespace.into(),
31 }
32 }
33
34 fn record_events(&self, events: &[Event]) {
35 for event in events {
36 match event {
37 Event::WorkflowStarted { workflow_type, .. } => {
38 self.metrics
39 .workflow_started(&self.namespace, workflow_type.as_str());
40 }
41 Event::WorkflowCompleted { .. } => {
42 self.metrics
43 .workflow_completed(&self.namespace, "completed");
44 }
45 Event::WorkflowFailed { .. } => {
46 self.metrics.workflow_completed(&self.namespace, "failed");
47 }
48 Event::WorkflowCancelled { .. } => {
49 self.metrics
50 .workflow_completed(&self.namespace, "cancelled");
51 }
52 Event::WorkflowTimedOut { .. } => {
53 self.metrics
54 .workflow_completed(&self.namespace, "timed_out");
55 }
56 Event::WorkflowContinuedAsNew { .. } => {
57 self.metrics
58 .workflow_completed(&self.namespace, "continued_as_new");
59 }
60 Event::SignalReceived { .. } => {
61 self.metrics.signal_delivered(&self.namespace, "resident");
62 }
63 Event::ScheduleTriggered { .. } => {
64 self.metrics.schedule_fired(&self.namespace);
65 }
66 _ => {}
67 }
68 }
69 }
70
71 fn observe_since(&self, operation: &str, started: Instant) {
72 self.metrics.store_operation(operation, started.elapsed());
73 }
74}
75
76impl std::fmt::Debug for InstrumentedEventStore {
77 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78 f.debug_struct("InstrumentedEventStore")
79 .field("namespace", &self.namespace)
80 .finish_non_exhaustive()
81 }
82}
83
84#[async_trait]
85impl WritableEventStore for InstrumentedEventStore {
86 async fn append(
87 &self,
88 token: WriteToken,
89 workflow_id: &WorkflowId,
90 events: &[Event],
91 expected_seq: u64,
92 ) -> Result<(), StoreError> {
93 let started = Instant::now();
94 let result = self
95 .inner
96 .append(token, workflow_id, events, expected_seq)
97 .await;
98 self.observe_since("append", started);
99 if result.is_ok() {
100 self.record_events(events);
101 }
102 result
103 }
104}
105
106#[async_trait]
107impl ReadableEventStore for InstrumentedEventStore {
108 async fn read_history(&self, workflow_id: &WorkflowId) -> Result<Vec<Event>, StoreError> {
109 let started = Instant::now();
110 let result = self.inner.read_history(workflow_id).await;
111 self.observe_since("read_history", started);
112 result
113 }
114
115 async fn read_history_from(
116 &self,
117 workflow_id: &WorkflowId,
118 from_seq: u64,
119 ) -> Result<Vec<Event>, StoreError> {
120 let started = Instant::now();
121 let result = self.inner.read_history_from(workflow_id, from_seq).await;
122 self.observe_since("read_history_from", started);
123 result
124 }
125
126 async fn read_run_chain(
127 &self,
128 workflow_id: &WorkflowId,
129 ) -> Result<Vec<RunSummary>, StoreError> {
130 self.inner.read_run_chain(workflow_id).await
131 }
132
133 async fn list_workflow_ids(&self) -> Result<Vec<WorkflowId>, StoreError> {
134 let started = Instant::now();
135 let result = self.inner.list_workflow_ids().await;
136 self.observe_since("list_workflow_ids", started);
137 result
138 }
139
140 async fn list_active(&self) -> Result<Vec<WorkflowId>, StoreError> {
141 let started = Instant::now();
142 let result = self.inner.list_active().await;
143 self.observe_since("list_active", started);
144 result
145 }
146
147 async fn query(&self, filter: &WorkflowFilter) -> Result<Vec<WorkflowSummary>, StoreError> {
148 self.inner.query(filter).await
149 }
150
151 async fn schedule_timer(
152 &self,
153 workflow_id: &WorkflowId,
154 timer_id: &TimerId,
155 fire_at: DateTime<Utc>,
156 ) -> Result<(), StoreError> {
157 self.inner
158 .schedule_timer(workflow_id, timer_id, fire_at)
159 .await
160 }
161
162 async fn expired_timers(&self, as_of: DateTime<Utc>) -> Result<Vec<TimerEntry>, StoreError> {
163 self.inner.expired_timers(as_of).await
164 }
165}
166
167#[async_trait]
168impl PackageStore for InstrumentedEventStore {
169 async fn put_package(&self, record: PackageRecord) -> Result<(), StoreError> {
170 let started = Instant::now();
171 let result = self.inner.put_package(record).await;
172 self.observe_since("put_package", started);
173 result
174 }
175
176 async fn list_packages(&self) -> Result<Vec<PackageRecord>, StoreError> {
177 let started = Instant::now();
178 let result = self.inner.list_packages().await;
179 self.observe_since("list_packages", started);
180 result
181 }
182
183 async fn delete_package(
184 &self,
185 workflow_type: &str,
186 content_hash: &str,
187 ) -> Result<(), StoreError> {
188 let started = Instant::now();
189 let result = self.inner.delete_package(workflow_type, content_hash).await;
190 self.observe_since("delete_package", started);
191 result
192 }
193
194 async fn put_package_route(
195 &self,
196 workflow_type: &str,
197 content_hash: &str,
198 ) -> Result<(), StoreError> {
199 let started = Instant::now();
200 let result = self
201 .inner
202 .put_package_route(workflow_type, content_hash)
203 .await;
204 self.observe_since("put_package_route", started);
205 result
206 }
207
208 async fn list_package_routes(&self) -> Result<Vec<PackageRouteRecord>, StoreError> {
209 let started = Instant::now();
210 let result = self.inner.list_package_routes().await;
211 self.observe_since("list_package_routes", started);
212 result
213 }
214}