use async_trait::async_trait;
use chrono::Utc;
use serde_json::json;
use std::sync::Arc;
use uuid::Uuid;
pub(crate) mod topics {
pub const USER_REGISTERED: &str = "udb.authn.user.registered.v1";
pub const USER_LOGGED_IN: &str = "udb.authn.user.login.v1";
pub const SESSION_REVOKED: &str = "udb.authn.session.revoked.v1";
pub const USER_LOCKED: &str = "udb.authn.user.locked.v1";
pub const PASSWORD_CHANGED: &str = "udb.authn.user.password.changed.v1";
pub const OTP_SENT: &str = "udb.authn.otp.sent.v1";
pub const USER_STATUS_CHANGED: &str = "udb.authn.user.status.changed.v1";
pub const EMAIL_VERIFIED: &str = "udb.authn.user.email.verified.v1";
pub const ROLE_CREATED: &str = "udb.authz.role.created.v1";
pub const ROLE_ASSIGNED: &str = "udb.authz.role.assigned.v1";
pub const ROLE_REVOKED: &str = "udb.authz.role.revoked.v1";
pub const ROLE_UPDATED: &str = "udb.authz.role.updated.v1";
pub const ACCESS_DENIED: &str = "udb.authz.access.denied.v1";
pub const API_KEY_CREATED: &str = "udb.apikey.created.v1";
pub const API_KEY_REVOKED: &str = "udb.apikey.revoked.v1";
pub const API_KEY_UPDATED: &str = "udb.apikey.updated.v1";
pub const AUTH_TOPIC_PATTERNS: &[&str] = &["udb.authn.*", "udb.authz.*", "udb.apikey.*"];
}
pub(crate) struct AuthEvent {
pub topic: &'static str,
pub document_id: String,
pub tenant_id: String,
pub correlation_id: String,
pub body: serde_json::Value,
}
impl AuthEvent {
pub(crate) fn new(
topic: &'static str,
document_id: impl Into<String>,
tenant_id: impl Into<String>,
body: serde_json::Value,
) -> Self {
Self {
topic,
document_id: document_id.into(),
tenant_id: tenant_id.into(),
correlation_id: String::new(),
body,
}
}
pub(crate) fn with_correlation(mut self, correlation_id: impl Into<String>) -> Self {
self.correlation_id = correlation_id.into();
self
}
}
#[async_trait]
pub(crate) trait AuthEventSink: Send + Sync {
async fn emit(&self, event: AuthEvent) -> Result<(), String>;
}
pub(crate) struct NoopAuthEventSink;
#[async_trait]
impl AuthEventSink for NoopAuthEventSink {
async fn emit(&self, _event: AuthEvent) -> Result<(), String> {
Ok(())
}
}
pub(crate) fn noop_sink() -> Arc<dyn AuthEventSink> {
Arc::new(NoopAuthEventSink)
}
pub(crate) struct OutboxAuthEventSink {
pool: sqlx::PgPool,
outbox_relation: String,
}
impl OutboxAuthEventSink {
pub(crate) fn new(pool: sqlx::PgPool, outbox_relation: impl Into<String>) -> Self {
Self {
pool,
outbox_relation: outbox_relation.into(),
}
}
fn envelope(event_id: &Uuid, event: &AuthEvent) -> serde_json::Value {
json!({
"event_id": event_id.to_string(),
"event_type": event.topic,
"timestamp": Utc::now().to_rfc3339(),
"correlation_id": event.correlation_id,
"document_id": event.document_id,
"payload": event.body,
})
}
}
#[async_trait]
impl AuthEventSink for OutboxAuthEventSink {
async fn emit(&self, event: AuthEvent) -> Result<(), String> {
let event_id = Uuid::new_v4();
let envelope = Self::envelope(&event_id, &event);
let sql = format!(
"INSERT INTO {rel} (event_id, topic, partition_key, payload, created_at) \
VALUES ($1::UUID, $2, $3, $4::JSONB, NOW())",
rel = self.outbox_relation
);
sqlx::query(&sql)
.bind(event_id)
.bind(event.topic)
.bind(&event.document_id)
.bind(envelope.to_string())
.execute(&self.pool)
.await
.map_err(|e| format!("auth outbox enqueue failed for topic {}: {e}", event.topic))?;
Ok(())
}
}