use std::sync::atomic::Ordering;
use cratestack_core::CratestackError;
use crate::SqlxRuntime;
use crate::sqlx;
pub const AUDIT_TABLE_DDL: &str = r#"
CREATE TABLE IF NOT EXISTS cratestack_audit (
event_id UUID PRIMARY KEY,
schema_name TEXT NOT NULL,
model TEXT NOT NULL,
operation TEXT NOT NULL,
primary_key JSONB NOT NULL,
actor JSONB NOT NULL,
tenant TEXT,
before JSONB,
after JSONB,
request_id TEXT,
occurred_at TIMESTAMPTZ NOT NULL,
delivered_at TIMESTAMPTZ,
attempts BIGINT NOT NULL DEFAULT 0,
last_error TEXT
);
CREATE INDEX IF NOT EXISTS cratestack_audit_model_idx
ON cratestack_audit (schema_name, model, occurred_at DESC);
CREATE INDEX IF NOT EXISTS cratestack_audit_tenant_idx
ON cratestack_audit (tenant, occurred_at DESC)
WHERE tenant IS NOT NULL;
CREATE INDEX IF NOT EXISTS cratestack_audit_undelivered_idx
ON cratestack_audit (occurred_at)
WHERE delivered_at IS NULL;
"#;
const AUDIT_OBJECTS_EXIST: &str = "SELECT to_regclass('cratestack_audit') IS NOT NULL \
AND to_regclass('cratestack_audit_model_idx') IS NOT NULL \
AND to_regclass('cratestack_audit_tenant_idx') IS NOT NULL \
AND to_regclass('cratestack_audit_undelivered_idx') IS NOT NULL";
pub(crate) async fn ensure_audit_table<'e, E>(
runtime: &SqlxRuntime,
probe: E,
) -> Result<(), CratestackError>
where
E: sqlx::Executor<'e, Database = sqlx::Postgres>,
{
if runtime.audit_table_ensured().load(Ordering::Acquire) {
return Ok(());
}
if runtime.bound().is_some() {
let exists: bool = sqlx::query_scalar(AUDIT_OBJECTS_EXIST)
.fetch_one(probe)
.await
.map_err(|error| CratestackError::Database(error.to_string()))?;
if exists {
runtime.audit_table_ensured().store(true, Ordering::Release);
return Ok(());
}
}
sqlx::raw_sql(AUDIT_TABLE_DDL)
.execute(runtime.pool())
.await
.map_err(|error| CratestackError::Database(error.to_string()))?;
runtime.audit_table_ensured().store(true, Ordering::Release);
Ok(())
}