use sqlx::{PgPool, Row};
use tonic::Status;
use uuid::Uuid;
use crate::runtime::native_catalog::NativeModel;
use super::super::auth_service::events::{ComplianceEnvelope, build_native_compliance_envelope};
use super::config::{LOCK_EXPIRY_SWEEP_BATCH, STATUS_EXPIRED, STATUS_HELD, TOPIC_EXPIRED};
use super::errors::lock_internal_status;
use super::events::lock_event_payload;
use super::model::lease_name;
use super::store::lock_model;
pub(crate) fn expired_locks_claim_sql(m: &NativeModel) -> String {
format!(
"UPDATE {rel} SET {status} = $1 \
WHERE {lock_id} IN ( \
SELECT {lock_id} FROM {rel} \
WHERE {status} = $2 AND {expires_at} < NOW() \
ORDER BY {expires_at} \
LIMIT $3 \
FOR UPDATE SKIP LOCKED) \
RETURNING {lock_id}::text AS lock_id, {tenant_id}::text AS tenant_id, \
{lock_name} AS lock_name, {owner_id} AS owner_id, \
{fencing_token} AS fencing_token",
rel = m.relation,
lock_id = m.q("lock_id"),
tenant_id = m.q("tenant_id"),
lock_name = m.q("lock_name"),
owner_id = m.q("owner_id"),
fencing_token = m.q("fencing_token"),
status = m.q("status"),
expires_at = m.q("expires_at"),
)
}
async fn insert_lock_expired_outbox(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
relation: &str,
lock_id: &str,
tenant_id: &str,
lock_name: &str,
owner_id: &str,
fencing_token: i64,
) -> Result<(), Status> {
let env = ComplianceEnvelope {
actor: "udb:lock".to_string(),
operation: "expired".to_string(),
outcome: "success".to_string(),
auth_method: "system".to_string(),
target_resource: lock_name.to_string(),
..ComplianceEnvelope::default()
};
let event_id = Uuid::new_v4();
let partition_key = lease_name(tenant_id, lock_name);
let envelope = build_native_compliance_envelope(
&event_id.to_string(),
TOPIC_EXPIRED,
&partition_key,
tenant_id,
"", &env,
lock_id, "none",
1,
&[],
lock_event_payload(tenant_id, "", lock_name, owner_id, fencing_token),
);
crate::runtime::cdc::insert_outbox_row(
&mut **tx,
relation,
event_id,
TOPIC_EXPIRED,
&partition_key,
&envelope,
)
.await
.map_err(|err| {
lock_internal_status(
"lock_expiry_outbox_insert",
format!("lock expiry outbox insert failed: {err}"),
)
})
}
pub(crate) async fn run_lock_expiry_once(
pool: &PgPool,
outbox_relation: Option<&str>,
batch_size: i64,
) -> Result<i64, Status> {
let Some(outbox_rel) = outbox_relation else {
tracing::warn!("lock expiry: no outbox relation configured; cannot expire locks");
return Ok(0);
};
let m = lock_model();
let claim_sql = expired_locks_claim_sql(&m);
let batch = batch_size.clamp(1, LOCK_EXPIRY_SWEEP_BATCH);
let mut tx = pool.begin().await.map_err(|err| {
lock_internal_status(
"lock_expiry_begin",
format!("lock expiry begin failed: {err}"),
)
})?;
let rows = sqlx::query(&claim_sql)
.bind(STATUS_EXPIRED)
.bind(STATUS_HELD)
.bind(batch)
.fetch_all(&mut *tx)
.await
.map_err(|err| {
lock_internal_status(
"lock_expiry_claim",
format!("lock expiry claim failed: {err}"),
)
})?;
let mut expired = 0i64;
for row in &rows {
let get = |c: &str| -> Result<String, Status> {
row.try_get::<String, _>(c).map_err(|e| {
lock_internal_status(
"lock_expiry_decode",
format!("lock expiry decode {c} failed: {e}"),
)
})
};
let lock_id = get("lock_id")?;
let tenant_id = get("tenant_id")?;
let lock_name = get("lock_name")?;
let owner_id = get("owner_id")?;
let fencing_token: i64 = row.try_get("fencing_token").map_err(|e| {
lock_internal_status(
"lock_expiry_decode",
format!("lock expiry decode fencing_token: {e}"),
)
})?;
insert_lock_expired_outbox(
&mut tx,
outbox_rel,
&lock_id,
&tenant_id,
&lock_name,
&owner_id,
fencing_token,
)
.await?;
expired += 1;
}
tx.commit().await.map_err(|err| {
lock_internal_status(
"lock_expiry_commit",
format!("lock expiry commit failed: {err}"),
)
})?;
Ok(expired)
}