use std::sync::Arc;
use std::time::Instant;
use aion_core::{Event, TimerId, WorkflowFilter, WorkflowId, WorkflowSummary};
use aion_store::{
EventStore, OutboxRow, PackageRecord, PackageRouteRecord, PackageStore, ReadableEventStore,
RunSummary, StoreError, TimerEntry, WritableEventStore, WriteToken,
};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use super::metrics::Metrics;
pub struct InstrumentedEventStore {
inner: Arc<dyn EventStore>,
metrics: Metrics,
namespace: String,
outbox_wake: Arc<tokio::sync::Notify>,
}
impl InstrumentedEventStore {
#[must_use]
pub fn new(inner: Arc<dyn EventStore>, metrics: Metrics, namespace: impl Into<String>) -> Self {
Self {
inner,
metrics,
namespace: namespace.into(),
outbox_wake: Arc::new(tokio::sync::Notify::new()),
}
}
#[must_use]
pub fn with_outbox_wake(mut self, outbox_wake: Arc<tokio::sync::Notify>) -> Self {
self.outbox_wake = outbox_wake;
self
}
fn record_events(&self, events: &[Event]) {
for event in events {
match event {
Event::WorkflowStarted { workflow_type, .. } => {
self.metrics
.workflow_started(&self.namespace, workflow_type.as_str());
}
Event::WorkflowCompleted { .. } => {
self.metrics
.workflow_completed(&self.namespace, "completed");
}
Event::WorkflowFailed { .. } => {
self.metrics.workflow_completed(&self.namespace, "failed");
}
Event::WorkflowCancelled { .. } => {
self.metrics
.workflow_completed(&self.namespace, "cancelled");
}
Event::WorkflowTimedOut { .. } => {
self.metrics
.workflow_completed(&self.namespace, "timed_out");
}
Event::WorkflowContinuedAsNew { .. } => {
self.metrics
.workflow_completed(&self.namespace, "continued_as_new");
}
Event::WorkflowReopened { .. } => {
self.metrics.workflow_reopened(&self.namespace);
}
Event::SignalReceived { .. } => {
self.metrics.signal_delivered(&self.namespace, "resident");
}
Event::ScheduleTriggered { .. } => {
self.metrics.schedule_fired(&self.namespace);
}
_ => {}
}
}
}
fn observe_since(&self, operation: &str, started: Instant) {
self.metrics.store_operation(operation, started.elapsed());
}
}
impl std::fmt::Debug for InstrumentedEventStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("InstrumentedEventStore")
.field("namespace", &self.namespace)
.finish_non_exhaustive()
}
}
#[async_trait]
impl WritableEventStore for InstrumentedEventStore {
async fn append(
&self,
token: WriteToken,
workflow_id: &WorkflowId,
events: &[Event],
expected_seq: u64,
) -> Result<(), StoreError> {
let started = Instant::now();
let result = self
.inner
.append(token, workflow_id, events, expected_seq)
.await;
self.observe_since("append", started);
if result.is_ok() {
self.record_events(events);
}
result
}
async fn append_with_outbox(
&self,
token: WriteToken,
workflow_id: &WorkflowId,
events: &[Event],
expected_seq: u64,
outbox_rows: &[OutboxRow],
) -> Result<(), StoreError> {
let started = Instant::now();
let result = self
.inner
.append_with_outbox(token, workflow_id, events, expected_seq, outbox_rows)
.await;
self.observe_since("append", started);
if result.is_ok() {
self.record_events(events);
if !outbox_rows.is_empty() {
self.outbox_wake.notify_one();
}
}
result
}
async fn rearm_outbox_pending(&self, rows: &[OutboxRow]) -> Result<(), StoreError> {
let started = Instant::now();
let result = self.inner.rearm_outbox_pending(rows).await;
self.observe_since("append", started);
result
}
}
#[async_trait]
impl ReadableEventStore for InstrumentedEventStore {
fn set_owned_shards(&self, shards: Option<&[usize]>) {
self.inner.set_owned_shards(shards);
}
fn acquire_owned_shards(&self, shards: &[usize]) -> Result<(), StoreError> {
self.inner.acquire_owned_shards(shards)
}
fn extend_owned_shards(&self, shards: &[usize]) {
self.inner.extend_owned_shards(shards);
}
async fn read_history(&self, workflow_id: &WorkflowId) -> Result<Vec<Event>, StoreError> {
let started = Instant::now();
let result = self.inner.read_history(workflow_id).await;
self.observe_since("read_history", started);
result
}
async fn read_history_from(
&self,
workflow_id: &WorkflowId,
from_seq: u64,
) -> Result<Vec<Event>, StoreError> {
let started = Instant::now();
let result = self.inner.read_history_from(workflow_id, from_seq).await;
self.observe_since("read_history_from", started);
result
}
async fn read_run_chain(
&self,
workflow_id: &WorkflowId,
) -> Result<Vec<RunSummary>, StoreError> {
self.inner.read_run_chain(workflow_id).await
}
async fn list_workflow_ids(&self) -> Result<Vec<WorkflowId>, StoreError> {
let started = Instant::now();
let result = self.inner.list_workflow_ids().await;
self.observe_since("list_workflow_ids", started);
result
}
async fn list_active(&self) -> Result<Vec<WorkflowId>, StoreError> {
let started = Instant::now();
let result = self.inner.list_active().await;
self.observe_since("list_active", started);
result
}
async fn query(&self, filter: &WorkflowFilter) -> Result<Vec<WorkflowSummary>, StoreError> {
self.inner.query(filter).await
}
async fn schedule_timer(
&self,
workflow_id: &WorkflowId,
timer_id: &TimerId,
fire_at: DateTime<Utc>,
) -> Result<(), StoreError> {
self.inner
.schedule_timer(workflow_id, timer_id, fire_at)
.await
}
async fn expired_timers(&self, as_of: DateTime<Utc>) -> Result<Vec<TimerEntry>, StoreError> {
self.inner.expired_timers(as_of).await
}
}
#[async_trait]
impl PackageStore for InstrumentedEventStore {
async fn put_package(&self, record: PackageRecord) -> Result<(), StoreError> {
let started = Instant::now();
let result = self.inner.put_package(record).await;
self.observe_since("put_package", started);
result
}
async fn list_packages(&self) -> Result<Vec<PackageRecord>, StoreError> {
let started = Instant::now();
let result = self.inner.list_packages().await;
self.observe_since("list_packages", started);
result
}
async fn delete_package(
&self,
workflow_type: &str,
content_hash: &str,
) -> Result<(), StoreError> {
let started = Instant::now();
let result = self.inner.delete_package(workflow_type, content_hash).await;
self.observe_since("delete_package", started);
result
}
async fn put_package_route(
&self,
workflow_type: &str,
content_hash: &str,
) -> Result<(), StoreError> {
let started = Instant::now();
let result = self
.inner
.put_package_route(workflow_type, content_hash)
.await;
self.observe_since("put_package_route", started);
result
}
async fn list_package_routes(&self) -> Result<Vec<PackageRouteRecord>, StoreError> {
let started = Instant::now();
let result = self.inner.list_package_routes().await;
self.observe_since("list_package_routes", started);
result
}
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use aion_core::{
ContentType, Event, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId,
};
use aion_store::{OutboxRow, WritableEventStore, WriteToken};
use aion_store_libsql::LibSqlStore;
use chrono::Utc;
use super::InstrumentedEventStore;
use crate::observability::Metrics;
fn unique_temp_path(name: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |duration| duration.as_nanos());
std::env::temp_dir().join(format!(
"aion-server-instrumented-store-{name}-{}-{nanos}.db",
std::process::id()
))
}
fn workflow_started(workflow_id: &WorkflowId) -> Event {
Event::WorkflowStarted {
envelope: EventEnvelope {
seq: 1,
recorded_at: Utc::now(),
workflow_id: workflow_id.clone(),
},
workflow_type: String::from("checkout"),
input: Payload::new(ContentType::Json, b"{}".to_vec()),
run_id: RunId::new_v4(),
parent_run_id: None,
package_version: PackageVersion::new("a".repeat(64)),
}
}
async fn wake_fired(wake: &tokio::sync::Notify) -> bool {
tokio::time::timeout(Duration::from_millis(200), wake.notified())
.await
.is_ok()
}
#[tokio::test]
async fn append_with_outbox_fires_wake_on_successful_non_empty_stage()
-> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(LibSqlStore::open(unique_temp_path("fires")).await?);
let metrics = Metrics::new()?;
let wake = Arc::new(tokio::sync::Notify::new());
let instrumented = InstrumentedEventStore::new(store, metrics, "default")
.with_outbox_wake(Arc::clone(&wake));
let workflow_id = WorkflowId::new_v4();
let event = workflow_started(&workflow_id);
let row = OutboxRow::pending(
workflow_id.clone(),
0,
String::from("charge"),
Payload::new(ContentType::Json, b"{}".to_vec()),
Utc::now(),
);
instrumented
.append_with_outbox(
WriteToken::recorder(),
&workflow_id,
std::slice::from_ref(&event),
0,
std::slice::from_ref(&row),
)
.await?;
assert!(
wake_fired(&wake).await,
"a successful non-empty outbox stage must pulse the advisory wake"
);
Ok(())
}
#[tokio::test]
async fn append_with_outbox_does_not_fire_wake_on_empty_slice()
-> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(LibSqlStore::open(unique_temp_path("empty")).await?);
let metrics = Metrics::new()?;
let wake = Arc::new(tokio::sync::Notify::new());
let instrumented = InstrumentedEventStore::new(store, metrics, "default")
.with_outbox_wake(Arc::clone(&wake));
let workflow_id = WorkflowId::new_v4();
let event = workflow_started(&workflow_id);
instrumented
.append_with_outbox(
WriteToken::recorder(),
&workflow_id,
std::slice::from_ref(&event),
0,
&[],
)
.await?;
assert!(
!wake_fired(&wake).await,
"an empty outbox slice must not pulse the wake (nothing to dispatch)"
);
Ok(())
}
#[tokio::test]
async fn append_with_outbox_does_not_fire_wake_on_failed_append()
-> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(LibSqlStore::open(unique_temp_path("failed")).await?);
let metrics = Metrics::new()?;
let wake = Arc::new(tokio::sync::Notify::new());
let instrumented = InstrumentedEventStore::new(store, metrics, "default")
.with_outbox_wake(Arc::clone(&wake));
let workflow_id = WorkflowId::new_v4();
let event = workflow_started(&workflow_id);
let row = OutboxRow::pending(
workflow_id.clone(),
0,
String::from("charge"),
Payload::new(ContentType::Json, b"{}".to_vec()),
Utc::now(),
);
let result = instrumented
.append_with_outbox(
WriteToken::recorder(),
&workflow_id,
std::slice::from_ref(&event),
9,
std::slice::from_ref(&row),
)
.await;
assert!(result.is_err(), "the seq-conflict append must fail");
assert!(
!wake_fired(&wake).await,
"a failed append commits nothing, so it must not pulse the wake"
);
Ok(())
}
}