use chrono::{DateTime, Utc};
use serde_json::Value;
use sqlx::{PgPool, Postgres, QueryBuilder};
use uuid::Uuid;
use super::Event;
use crate::error::{CoreError, Result};
#[derive(Debug, Clone)]
pub struct PersistedEvent {
pub id: Uuid,
pub event_type: String,
pub aggregate_id: Option<Uuid>,
pub user_id: Option<Uuid>,
pub payload: Value,
pub occurred_at: DateTime<Utc>,
pub persisted_at: DateTime<Utc>,
}
pub struct EventStore {
pool: PgPool,
}
impl EventStore {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
pub async fn persist(&self, event: &Event) -> Result<Uuid> {
let event_id = Uuid::new_v4();
let event_type = event.event_type().to_string();
let occurred_at = event.occurred_at();
let payload =
serde_json::to_value(event).map_err(|e| CoreError::Serialization(e.to_string()))?;
let (aggregate_id, user_id) = self.extract_ids(event);
sqlx::query(
"INSERT INTO events (id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at)
VALUES ($1, $2, $3, $4, $5, $6, NOW())"
)
.bind(event_id)
.bind(event_type)
.bind(aggregate_id)
.bind(user_id)
.bind(payload)
.bind(occurred_at)
.execute(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
Ok(event_id)
}
pub async fn persist_batch(&self, events: &[Event]) -> Result<Vec<Uuid>> {
if events.is_empty() {
return Ok(Vec::new());
}
let mut query_builder: QueryBuilder<Postgres> = QueryBuilder::new(
"INSERT INTO events (id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at) ",
);
let mut event_ids = Vec::with_capacity(events.len());
query_builder.push_values(events, |mut b, event| {
let event_id = Uuid::new_v4();
event_ids.push(event_id);
let event_type = event.event_type().to_string();
let occurred_at = event.occurred_at();
let payload = serde_json::to_value(event).expect("Failed to serialize event");
let (aggregate_id, user_id) = self.extract_ids(event);
b.push_bind(event_id)
.push_bind(event_type)
.push_bind(aggregate_id)
.push_bind(user_id)
.push_bind(payload)
.push_bind(occurred_at)
.push("NOW()");
});
query_builder
.build()
.execute(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
Ok(event_ids)
}
pub async fn get_by_type(
&self,
event_type: &str,
limit: i64,
offset: i64,
) -> Result<Vec<PersistedEvent>> {
let rows = sqlx::query_as::<
_,
(
Uuid,
String,
Option<Uuid>,
Option<Uuid>,
Value,
DateTime<Utc>,
DateTime<Utc>,
),
>(
"SELECT id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at
FROM events
WHERE event_type = $1
ORDER BY occurred_at DESC
LIMIT $2 OFFSET $3",
)
.bind(event_type)
.bind(limit)
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
let events = rows
.into_iter()
.map(
|(id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at)| {
PersistedEvent {
id,
event_type,
aggregate_id,
user_id,
payload,
occurred_at,
persisted_at,
}
},
)
.collect();
Ok(events)
}
pub async fn get_by_aggregate(
&self,
aggregate_id: Uuid,
limit: i64,
offset: i64,
) -> Result<Vec<PersistedEvent>> {
let rows = sqlx::query_as::<
_,
(
Uuid,
String,
Option<Uuid>,
Option<Uuid>,
Value,
DateTime<Utc>,
DateTime<Utc>,
),
>(
"SELECT id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at
FROM events
WHERE aggregate_id = $1
ORDER BY occurred_at DESC
LIMIT $2 OFFSET $3",
)
.bind(aggregate_id)
.bind(limit)
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
let events = rows
.into_iter()
.map(
|(id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at)| {
PersistedEvent {
id,
event_type,
aggregate_id,
user_id,
payload,
occurred_at,
persisted_at,
}
},
)
.collect();
Ok(events)
}
pub async fn get_by_user(
&self,
user_id: Uuid,
limit: i64,
offset: i64,
) -> Result<Vec<PersistedEvent>> {
let rows = sqlx::query_as::<
_,
(
Uuid,
String,
Option<Uuid>,
Option<Uuid>,
Value,
DateTime<Utc>,
DateTime<Utc>,
),
>(
"SELECT id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at
FROM events
WHERE user_id = $1
ORDER BY occurred_at DESC
LIMIT $2 OFFSET $3",
)
.bind(user_id)
.bind(limit)
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
let events = rows
.into_iter()
.map(
|(id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at)| {
PersistedEvent {
id,
event_type,
aggregate_id,
user_id,
payload,
occurred_at,
persisted_at,
}
},
)
.collect();
Ok(events)
}
pub async fn get_by_time_range(
&self,
start: DateTime<Utc>,
end: DateTime<Utc>,
limit: i64,
offset: i64,
) -> Result<Vec<PersistedEvent>> {
let rows = sqlx::query_as::<
_,
(
Uuid,
String,
Option<Uuid>,
Option<Uuid>,
Value,
DateTime<Utc>,
DateTime<Utc>,
),
>(
"SELECT id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at
FROM events
WHERE occurred_at >= $1 AND occurred_at <= $2
ORDER BY occurred_at DESC
LIMIT $3 OFFSET $4",
)
.bind(start)
.bind(end)
.bind(limit)
.bind(offset)
.fetch_all(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
let events = rows
.into_iter()
.map(
|(id, event_type, aggregate_id, user_id, payload, occurred_at, persisted_at)| {
PersistedEvent {
id,
event_type,
aggregate_id,
user_id,
payload,
occurred_at,
persisted_at,
}
},
)
.collect();
Ok(events)
}
pub async fn count_events(&self) -> Result<i64> {
let (count,): (i64,) = sqlx::query_as("SELECT COUNT(*) FROM events")
.fetch_one(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
Ok(count)
}
pub async fn count_by_type(&self, event_type: &str) -> Result<i64> {
let (count,): (i64,) = sqlx::query_as("SELECT COUNT(*) FROM events WHERE event_type = $1")
.bind(event_type)
.fetch_one(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
Ok(count)
}
pub async fn delete_before(&self, timestamp: DateTime<Utc>) -> Result<u64> {
let result = sqlx::query("DELETE FROM events WHERE occurred_at < $1")
.bind(timestamp)
.execute(&self.pool)
.await
.map_err(|e| CoreError::Database(e.to_string()))?;
Ok(result.rows_affected())
}
fn extract_ids(&self, event: &Event) -> (Option<Uuid>, Option<Uuid>) {
match event {
Event::UserRegistered(e) => (Some(e.user_id), Some(e.user_id)),
Event::UserUpdated(e) => (Some(e.user_id), Some(e.user_id)),
Event::UserDeleted(e) => (Some(e.user_id), Some(e.user_id)),
Event::TokenCreated(e) => (Some(e.token_id), Some(e.issuer_id)),
Event::TokenUpdated(e) => (Some(e.token_id), Some(e.updated_by)),
Event::TokenPaused(e) => (Some(e.token_id), Some(e.paused_by)),
Event::TokenResumed(e) => (Some(e.token_id), Some(e.resumed_by)),
Event::OrderCreated(e) => (Some(e.order_id), Some(e.user_id)),
Event::OrderCancelled(e) => (Some(e.order_id), Some(e.user_id)),
Event::OrderFilled(e) => (Some(e.order_id), Some(e.user_id)),
Event::OrderExpired(e) => (Some(e.order_id), Some(e.user_id)),
Event::TradeExecuted(e) => (Some(e.trade_id), Some(e.buyer_id)),
Event::PaymentCreated(e) => (Some(e.payment_id), Some(e.user_id)),
Event::PaymentDetected(e) => (Some(e.payment_id), Some(e.user_id)),
Event::PaymentConfirmed(e) => (Some(e.payment_id), Some(e.user_id)),
Event::PaymentCompleted(e) => (Some(e.payment_id), Some(e.user_id)),
Event::PaymentExpired(e) => (Some(e.payment_id), Some(e.user_id)),
Event::CommitmentCreated(e) => (Some(e.commitment_id), Some(e.user_id)),
Event::CommitmentFulfilled(e) => (Some(e.commitment_id), Some(e.user_id)),
Event::CommitmentFailed(e) => (Some(e.commitment_id), Some(e.user_id)),
Event::ReputationChanged(e) => (Some(e.user_id), Some(e.user_id)),
Event::KycSubmitted(e) => (Some(e.submission_id), Some(e.user_id)),
Event::KycApproved(e) => (Some(e.submission_id), Some(e.user_id)),
Event::KycRejected(e) => (Some(e.submission_id), Some(e.user_id)),
Event::ProposalCreated(e) => (Some(e.proposal_id), Some(e.proposer_id)),
Event::VoteCast(e) => (Some(e.vote_id), Some(e.voter_id)),
Event::ProposalExecuted(e) => (Some(e.proposal_id), Some(e.executed_by)),
Event::SecurityAlert(e) => (Some(e.alert_id), e.user_id),
Event::CircuitBreakerTriggered(e) => (Some(e.token_id), e.triggered_by),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::events::{TokenCreatedEvent, UserRegisteredEvent};
use rust_decimal_macros::dec;
#[test]
fn test_extract_user_ids() {
let user_id = Uuid::new_v4();
let event = Event::UserRegistered(UserRegisteredEvent {
user_id,
username: "test".to_string(),
email: "test@example.com".to_string(),
occurred_at: Utc::now(),
});
assert_eq!(event.event_type(), "user_registered");
}
#[test]
fn test_extract_token_ids() {
let token_id = Uuid::new_v4();
let issuer_id = Uuid::new_v4();
let event = Event::TokenCreated(TokenCreatedEvent {
token_id,
issuer_id,
symbol: "$TEST".to_string(),
name: "Test Token".to_string(),
initial_supply: dec!(1000000),
occurred_at: Utc::now(),
});
assert_eq!(event.event_type(), "token_created");
}
}