use sqlx::PgPool;
use uuid::Uuid;
use crate::application::service::event_error::EventError;
use super::seat_repository::{record_audit, RegisterCommand, SaleLink, SeatRepository};
#[derive(Debug, Clone, serde::Serialize)]
pub struct SeamOutcome {
pub already_applied: bool,
pub registrations_touched: usize,
pub booths_latched_paid: usize,
}
pub fn grand_total_is_zero(raw: &str) -> bool {
match raw.trim().parse::<f64>() {
Ok(v) => v == 0.0,
Err(_) => false,
}
}
pub struct SaleSeamRepository {
pool: PgPool,
}
impl SaleSeamRepository {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
pub fn pool(&self) -> &PgPool {
&self.pool
}
fn rpool(&self) -> PgPool {
crate::request_pool::current().unwrap_or_else(|| self.pool.clone())
}
pub async fn multi_slot_flags(
&self,
event_ids: &[uuid::Uuid],
) -> Result<std::collections::HashMap<uuid::Uuid, bool>, EventError> {
let rows: Vec<(uuid::Uuid, bool)> = sqlx::query_as(
"SELECT id, is_multi_slots FROM event.events WHERE id = ANY($1)",
)
.bind(event_ids)
.fetch_all(&self.rpool())
.await?;
Ok(rows.into_iter().collect())
}
async fn claim(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
consumer: &str,
external_id: &str,
) -> Result<bool, EventError> {
let claimed = sqlx::query_scalar::<_, i64>(
r#"INSERT INTO event.seam_inbox (consumer, external_id)
VALUES ($1, $2)
ON CONFLICT (consumer, external_id) DO NOTHING
RETURNING 1::int8"#,
)
.bind(consumer)
.bind(external_id)
.fetch_optional(&mut **tx)
.await?
.is_some();
Ok(claimed)
}
fn claim_key(verb: &str, delivery_id: Option<&str>, order_id: Uuid) -> String {
match delivery_id {
Some(d) => format!("{verb}:{d}"),
None => format!("{verb}:{order_id}"),
}
}
pub async fn on_order_confirmed(
&self,
consumer_key: &str,
delivery_id: Option<&str>,
order_id: Uuid,
grand_total: &str,
specs: Vec<RegisterCommand>,
actor: Option<Uuid>,
) -> Result<SeamOutcome, EventError> {
let free = grand_total_is_zero(grand_total);
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
if !Self::claim(
&mut tx,
consumer_key,
&Self::claim_key("confirmed", delivery_id, order_id),
)
.await?
{
tx.commit().await?;
return Ok(SeamOutcome {
already_applied: true,
registrations_touched: 0,
booths_latched_paid: 0,
});
}
let linked: i64 = sqlx::query_scalar::<_, i64>(
"SELECT count(*) FROM event.registrations WHERE sale_order_id = $1",
)
.bind(order_id)
.fetch_one(&mut *tx)
.await?;
let mut touched = 0usize;
if linked == 0 {
let link = SaleLink {
sale_order_id: order_id,
sale_order_state: "sale".to_string(),
sale_status: if free { "free" } else { "to_pay" }.to_string(),
initial_state: if free { "open" } else { "draft" }.to_string(),
};
for spec in &specs {
SeatRepository::register_core(&mut tx, spec, Some(&link)).await?;
touched += 1;
}
tx.commit().await?;
record_audit(
&self.rpool(),
"sale_seam_confirmed",
actor,
"sale_order",
order_id,
serde_json::json!({
"verb": "mint",
"free": free,
"minted": touched,
"sale_status": link.sale_status,
}),
)
.await;
} else {
let detail = if free {
let rows = sqlx::query(
r#"UPDATE event.registrations
SET state = 'open',
sale_order_state = 'sale',
sale_status = 'free'
WHERE sale_order_id = $1
AND state IN ('draft','cancel')
AND active
RETURNING id"#,
)
.bind(order_id)
.fetch_all(&mut *tx)
.await?;
touched = rows.len();
sqlx::query(
r#"UPDATE event.mails m SET scheduled_date = now(), mail_done = false
WHERE m.interval_kind = 'after_sub'
AND m.event_id IN (
SELECT DISTINCT event_id FROM event.registrations
WHERE sale_order_id = $1)"#,
)
.bind(order_id)
.execute(&mut *tx)
.await?;
serde_json::json!({
"verb": "heal",
"free": true,
"healed_to_open": touched,
})
} else {
let rows = sqlx::query(
r#"UPDATE event.registrations
SET sale_order_state = 'sale',
sale_status = 'to_pay'
WHERE sale_order_id = $1 AND active
RETURNING id"#,
)
.bind(order_id)
.fetch_all(&mut *tx)
.await?;
touched = rows.len();
serde_json::json!({
"verb": "heal",
"free": false,
"held": touched,
"note": "held at draft until the paid fact",
})
};
record_audit(
&self.rpool(),
"sale_seam_confirmed",
actor,
"sale_order",
order_id,
detail,
)
.await;
tx.commit().await?;
}
Ok(SeamOutcome {
already_applied: false,
registrations_touched: touched,
booths_latched_paid: 0,
})
}
pub async fn on_order_cancelled(
&self,
consumer_key: &str,
delivery_id: Option<&str>,
order_id: Uuid,
actor: Option<Uuid>,
) -> Result<SeamOutcome, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
if !Self::claim(
&mut tx,
consumer_key,
&Self::claim_key("cancelled", delivery_id, order_id),
)
.await?
{
tx.commit().await?;
return Ok(SeamOutcome {
already_applied: true,
registrations_touched: 0,
booths_latched_paid: 0,
});
}
let rows = sqlx::query(
r#"UPDATE event.registrations
SET state = 'cancel',
sale_order_state = 'cancel'
WHERE sale_order_id = $1 AND active
RETURNING id"#,
)
.bind(order_id)
.fetch_all(&mut *tx)
.await?;
let touched = rows.len();
record_audit(
&self.rpool(),
"sale_seam_cancelled",
actor,
"sale_order",
order_id,
serde_json::json!({ "verb": "cancel_cascade", "cancelled": touched }),
)
.await;
tx.commit().await?;
Ok(SeamOutcome {
already_applied: false,
registrations_touched: touched,
booths_latched_paid: 0,
})
}
pub async fn on_order_paid(
&self,
consumer_key: &str,
delivery_id: Option<&str>,
order_id: Uuid,
line_ids: &[Uuid],
actor: Option<Uuid>,
) -> Result<SeamOutcome, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
if !Self::claim(
&mut tx,
consumer_key,
&Self::claim_key("paid", delivery_id, order_id),
)
.await?
{
tx.commit().await?;
return Ok(SeamOutcome {
already_applied: true,
registrations_touched: 0,
booths_latched_paid: 0,
});
}
let healed = sqlx::query(
r#"UPDATE event.registrations
SET state = 'open',
sale_order_state = 'sale',
sale_status = 'sold'
WHERE sale_order_id = $1
AND state IN ('draft','cancel')
AND sale_status = 'to_pay'
AND active
RETURNING id"#,
)
.bind(order_id)
.fetch_all(&mut *tx)
.await?;
let touched = healed.len();
sqlx::query(
r#"UPDATE event.registrations
SET sale_status = 'sold'
WHERE sale_order_id = $1 AND sale_status = 'to_pay' AND active"#,
)
.bind(order_id)
.execute(&mut *tx)
.await?;
sqlx::query(
r#"UPDATE event.mails m SET scheduled_date = now(), mail_done = false
WHERE m.interval_kind = 'after_sub'
AND m.event_id IN (
SELECT DISTINCT event_id FROM event.registrations
WHERE sale_order_id = $1)"#,
)
.bind(order_id)
.execute(&mut *tx)
.await?;
let mut latched = 0usize;
if !line_ids.is_empty() {
let rows = sqlx::query(
r#"UPDATE event.booths SET is_paid = true
WHERE sale_order_line_id = ANY($1) AND NOT is_paid
RETURNING id"#,
)
.bind(line_ids)
.fetch_all(&mut *tx)
.await?;
latched = rows.len();
}
record_audit(
&self.rpool(),
"sale_seam_paid",
actor,
"sale_order",
order_id,
serde_json::json!({
"verb": "paid_heal",
"healed": touched,
"booths_latched": latched,
}),
)
.await;
tx.commit().await?;
Ok(SeamOutcome {
already_applied: false,
registrations_touched: touched,
booths_latched_paid: latched,
})
}
}