Skip to main content

aion_server/observability/
instrumented_store.rs

1//! [`InstrumentedEventStore`]: event-store decorator recording server metrics.
2
3use 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
16/// Event-store wrapper that observes operation latency and lifecycle events without changing engine crates.
17pub struct InstrumentedEventStore {
18    inner: Arc<dyn EventStore>,
19    metrics: Metrics,
20    namespace: String,
21}
22
23impl InstrumentedEventStore {
24    /// Wrap an event store with server-side metrics.
25    #[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}