use sqlx::PgPool;
use uuid::Uuid;
fn validate_outbox_relation(outbox_relation: &str) -> Result<(), String> {
if outbox_relation.is_empty()
|| outbox_relation
.chars()
.any(|c| !(c.is_ascii_alphanumeric() || c == '_' || c == '.' || c == '"'))
{
return Err(format!(
"reset_indoubt_publishing_rows: refusing unsafe outbox_relation '{outbox_relation}'"
));
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IndoubtRowOutcome {
Acked,
Pending,
Untouched,
}
pub async fn reset_indoubt_publishing_row(
pool: &PgPool,
outbox_relation: &str,
event_id: Uuid,
current_epoch: i64,
) -> Result<IndoubtRowOutcome, String> {
validate_outbox_relation(outbox_relation)?;
let ack_sql = format!(
"UPDATE {outbox_relation} SET \
delivery_state = 'acked', \
acked_at = NOW(), \
last_error = COALESCE(last_error, '') || \
' | udb-cdc-indoubt-recovery: commit confirmed by recorded offset; acking' \
WHERE event_id = $1 \
AND delivery_state = 'publishing' \
AND producer_epoch = $2 \
AND kafka_offset IS NOT NULL"
);
let acked = sqlx::query(&ack_sql)
.bind(event_id)
.bind(current_epoch)
.execute(pool)
.await
.map(|res| res.rows_affected())
.map_err(|e| format!("reset_indoubt_publishing_row (ack) failed: {e}"))?;
if acked > 0 {
return Ok(IndoubtRowOutcome::Acked);
}
let reset_sql = format!(
"UPDATE {outbox_relation} SET \
delivery_state = 'pending', \
last_error = COALESCE(last_error, '') || \
' | udb-cdc-indoubt-recovery: in-flight publish dropped on delivery timeout; re-publishing', \
publishing_started_at = NULL \
WHERE event_id = $1 \
AND delivery_state = 'publishing' \
AND producer_epoch = $2 \
AND kafka_offset IS NULL"
);
let reset = sqlx::query(&reset_sql)
.bind(event_id)
.bind(current_epoch)
.execute(pool)
.await
.map(|res| res.rows_affected())
.map_err(|e| format!("reset_indoubt_publishing_row failed: {e}"))?;
if reset > 0 {
Ok(IndoubtRowOutcome::Pending)
} else {
Ok(IndoubtRowOutcome::Untouched)
}
}
pub async fn reset_indoubt_publishing_rows(
pool: &PgPool,
outbox_relation: &str,
current_epoch: i64,
grace_secs: i64,
shard_band: Option<(i64, i64)>,
) -> Result<u64, String> {
validate_outbox_relation(outbox_relation)?;
if grace_secs < 0 {
return Err(format!(
"reset_indoubt_publishing_rows: grace_secs must be >= 0 (got {grace_secs})"
));
}
let ack_band = if shard_band.is_some() {
" AND producer_epoch >= $1 AND producer_epoch <= $2"
} else {
""
};
let ack_sql = format!(
"UPDATE {outbox_relation} SET \
delivery_state = 'acked', \
acked_at = NOW(), \
last_error = COALESCE(last_error, '') || \
' | udb-cdc-indoubt-recovery: commit confirmed by recorded offset; acking' \
WHERE delivery_state = 'publishing' \
AND kafka_offset IS NOT NULL{ack_band}"
);
let mut ack_query = sqlx::query(&ack_sql);
if let Some((lo, hi)) = shard_band {
ack_query = ack_query.bind(lo).bind(hi);
}
let acked = ack_query
.execute(pool)
.await
.map(|res| res.rows_affected())
.map_err(|e| format!("reset_indoubt_publishing_rows (ack) failed: {e}"))?;
let reset_band = if shard_band.is_some() {
" AND producer_epoch >= $3 AND producer_epoch <= $4"
} else {
""
};
let reset_sql = format!(
"UPDATE {outbox_relation} SET \
delivery_state = 'pending', \
last_error = COALESCE(last_error, '') || \
' | udb-cdc-indoubt-recovery: previous epoch fenced; re-publishing', \
publishing_started_at = NULL \
WHERE delivery_state = 'publishing' \
AND kafka_offset IS NULL \
AND ( \
producer_epoch < $1 \
OR publishing_started_at < NOW() - make_interval(secs => $2::double precision) \
){reset_band}"
);
let mut reset_query = sqlx::query(&reset_sql)
.bind(current_epoch)
.bind(grace_secs as f64);
if let Some((lo, hi)) = shard_band {
reset_query = reset_query.bind(lo).bind(hi);
}
let reset = reset_query
.execute(pool)
.await
.map(|res| res.rows_affected())
.map_err(|e| format!("reset_indoubt_publishing_rows failed: {e}"))?;
Ok(acked + reset)
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn rejects_unsafe_relation_name() {
let pool = PgPool::connect_lazy("postgres://invalid:0/none").unwrap();
let err = reset_indoubt_publishing_rows(&pool, "evil; DROP", 0, 300, None)
.await
.expect_err("should reject unsafe relation");
assert!(
err.contains("refusing unsafe outbox_relation"),
"got: {err}"
);
let err2 = reset_indoubt_publishing_rows(&pool, "", 0, 300, None)
.await
.expect_err("should reject empty relation");
assert!(err2.contains("refusing unsafe"));
}
#[tokio::test]
async fn per_event_reset_rejects_unsafe_relation_name() {
let pool = PgPool::connect_lazy("postgres://invalid:0/none").unwrap();
let err = reset_indoubt_publishing_row(&pool, "evil; DROP", Uuid::new_v4(), 0)
.await
.expect_err("should reject unsafe relation");
assert!(
err.contains("refusing unsafe outbox_relation"),
"got: {err}"
);
}
#[tokio::test]
async fn rejects_negative_grace() {
let pool = PgPool::connect_lazy("postgres://invalid:0/none").unwrap();
let err = reset_indoubt_publishing_rows(&pool, "udb_system.outbox_events", 0, -1, None)
.await
.expect_err("negative grace must be rejected");
assert!(err.contains("grace_secs must be >= 0"), "got: {err}");
}
#[tokio::test]
async fn accepts_schema_qualified_relation() {
let pool = PgPool::connect_lazy("postgres://invalid:0/none").unwrap();
let res =
reset_indoubt_publishing_rows(&pool, "udb_system.outbox_events", 0, 300, None).await;
match res {
Ok(_) => {} Err(e) => assert!(
!e.contains("refusing unsafe"),
"schema-qualified should not be rejected: {e}"
),
}
}
}