use backbone_orm::{company_scope, org_scope};
use chrono::{DateTime, Utc};
use rand::RngCore;
use sqlx::PgPool;
use uuid::Uuid;
use crate::application::service::event_error::EventError;
#[derive(Debug, Clone)]
pub struct RegisterCommand {
pub event_id: Uuid,
pub event_slot_id: Option<Uuid>,
pub event_ticket_id: Option<Uuid>,
pub name: String,
pub email: String,
pub phone: Option<String>,
pub company_name: Option<String>,
pub partner_id: Option<Uuid>,
pub actor: Option<Uuid>,
pub lead_rule_skip: bool,
}
#[derive(Debug, Clone)]
pub struct SaleLink {
pub sale_order_id: Uuid,
pub sale_order_state: String,
pub sale_status: String,
pub initial_state: String,
}
#[derive(Debug, Clone, serde::Serialize, sqlx::FromRow)]
pub struct RegistrationRow {
pub id: Uuid,
pub event_id: Uuid,
pub event_slot_id: Option<Uuid>,
pub event_ticket_id: Option<Uuid>,
pub name: String,
pub email: String,
pub phone: Option<String>,
pub company_name: Option<String>,
pub partner_id: Option<Uuid>,
pub state: String,
pub date_closed: Option<DateTime<Utc>>,
pub sale_order_id: Option<Uuid>,
pub sale_order_state: Option<String>,
pub sale_status: Option<String>,
pub active: bool,
pub barcode: String,
}
#[derive(Debug, Clone, Copy)]
pub struct SeatCounts {
pub limited: bool,
pub capacity: i64,
pub taken: i64,
}
impl SeatCounts {
pub fn available(&self) -> i64 {
if !self.limited || self.capacity == 0 {
i64::MAX
} else {
(self.capacity - self.taken).max(0)
}
}
}
pub fn mint_barcode() -> String {
let mut bytes = [0u8; 8];
rand::thread_rng().fill_bytes(&mut bytes);
u64::from_le_bytes(bytes).to_string()
}
pub async fn record_audit(
pool: &PgPool,
kind: &str,
actor: Option<Uuid>,
subject_type: &str,
subject_id: Uuid,
detail: serde_json::Value,
) {
let _ = crate::infrastructure::persistence::audit::record_audit_on_pool(
pool,
kind,
actor,
subject_type,
Some(subject_id),
detail)
.await;
}
pub struct SeatRepository {
pool: PgPool,
}
#[derive(Debug, Clone, sqlx::FromRow)]
struct LockedEvent {
#[sqlx(rename = "id")]
_id: Uuid,
is_multi_slots: bool,
seats_limited: bool,
seats_max: i32,
kanban_state: String,
}
impl SeatRepository {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
pub fn pool(&self) -> &PgPool {
&self.pool
}
pub fn rpool(&self) -> PgPool {
crate::request_pool::current().unwrap_or_else(|| self.pool.clone())
}
pub async fn register(&self, cmd: &RegisterCommand) -> Result<RegistrationRow, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let row = Self::register_core(&mut tx, cmd, None).await?;
tx.commit().await?;
Ok(row)
}
pub async fn register_sale_linked(
&self,
cmd: &RegisterCommand,
link: &SaleLink,
) -> Result<RegistrationRow, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let row = Self::register_core(&mut tx, cmd, Some(link)).await?;
tx.commit().await?;
Ok(row)
}
pub async fn register_core(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
cmd: &RegisterCommand,
link: Option<&SaleLink>,
) -> Result<RegistrationRow, EventError> {
let event = sqlx::query_as::<_, LockedEvent>(
r#"SELECT id, is_multi_slots, seats_limited, seats_max, kanban_state::text AS kanban_state
FROM event.events WHERE id = $1 FOR UPDATE"#,
)
.bind(cmd.event_id)
.fetch_optional(&mut **tx)
.await?
.ok_or(EventError::EventNotFound)?;
if event.kanban_state == "cancel" {
return Err(EventError::Validation(
"event is cancelled — registrations refused".to_string(),
));
}
let slot_id = if event.is_multi_slots {
let slot = cmd.event_slot_id.ok_or(EventError::EventSlotRequired {
event_id: cmd.event_id,
})?;
let belongs = sqlx::query_scalar::<_, Uuid>(
"SELECT id FROM event.slots WHERE id = $1 AND event_id = $2 FOR UPDATE",
)
.bind(slot)
.bind(cmd.event_id)
.fetch_optional(&mut **tx)
.await?;
if belongs.is_none() {
return Err(EventError::EventSlotNotOfEvent {
event_slot_id: slot,
});
}
Some(slot)
} else {
cmd.event_slot_id
};
if let Some(ticket) = cmd.event_ticket_id {
let window = sqlx::query_as::<_, (Option<DateTime<Utc>>, Option<DateTime<Utc>>)>(
"SELECT start_sale_datetime, end_sale_datetime FROM event.tickets WHERE id = $1 AND event_id = $2",
)
.bind(ticket)
.bind(cmd.event_id)
.fetch_optional(&mut **tx)
.await?
.ok_or(EventError::EventTicketNotOfEvent { event_ticket_id: ticket })?;
let (start, end) = window;
let now = Utc::now();
let shut =
start.map(|s| now < s).unwrap_or(false) || end.map(|e| now > e).unwrap_or(false);
if shut {
return Err(EventError::EventSaleWindowClosed {
event_ticket_id: ticket,
});
}
}
let taken: i64 = match slot_id {
Some(slot) => {
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations
WHERE event_id = $1 AND event_slot_id = $2
AND state IN ('open','done') AND active"#,
)
.bind(cmd.event_id)
.bind(slot)
.fetch_one(&mut **tx)
.await?
}
None => {
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations
WHERE event_id = $1
AND state IN ('open','done') AND active"#,
)
.bind(cmd.event_id)
.fetch_one(&mut **tx)
.await?
}
};
if event.seats_limited && event.seats_max > 0 && taken >= event.seats_max as i64 {
return Err(EventError::EventSeatsExhausted {
event_id: cmd.event_id,
});
}
let born_state = link.map(|l| l.initial_state.as_str()).unwrap_or("open");
let id = Uuid::new_v4();
let barcode = mint_barcode();
let org_axis: bool = sqlx::query_scalar::<_, bool>(
"SELECT EXISTS (SELECT 1 FROM information_schema.columns \
WHERE table_schema = 'event' AND table_name = 'events' \
AND column_name = 'org_unit_id')",
)
.fetch_one(&mut **tx)
.await?;
let row = if org_axis {
sqlx::query_as::<_, RegistrationRow>(
r#"INSERT INTO event.registrations
(id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state, active, barcode, org_unit_id,
sale_order_id, sale_order_state, sale_status)
SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9,
$11::event_registration_state, true, $10, e.org_unit_id,
$12, $13::event_sale_order_state, $14::event_sale_status
FROM event.events e WHERE e.id = $2
RETURNING id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state::text AS state, date_closed,
sale_order_id, sale_order_state::text AS sale_order_state,
sale_status::text AS sale_status, active, barcode"#,
)
.bind(id)
.bind(cmd.event_id)
.bind(slot_id)
.bind(cmd.event_ticket_id)
.bind(&cmd.name)
.bind(&cmd.email)
.bind(&cmd.phone)
.bind(&cmd.company_name)
.bind(cmd.partner_id)
.bind(&barcode)
.bind(born_state)
.bind(link.map(|l| l.sale_order_id))
.bind(link.map(|l| l.sale_order_state.as_str()))
.bind(link.map(|l| l.sale_status.as_str()))
.fetch_one(&mut **tx)
.await?
} else {
sqlx::query_as::<_, RegistrationRow>(
r#"INSERT INTO event.registrations
(id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state, active, barcode,
sale_order_id, sale_order_state, sale_status)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $11::event_registration_state,
true, $10, $12, $13::event_sale_order_state, $14::event_sale_status)
RETURNING id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state::text AS state, date_closed,
sale_order_id, sale_order_state::text AS sale_order_state,
sale_status::text AS sale_status, active, barcode"#,
)
.bind(id)
.bind(cmd.event_id)
.bind(slot_id)
.bind(cmd.event_ticket_id)
.bind(&cmd.name)
.bind(&cmd.email)
.bind(&cmd.phone)
.bind(&cmd.company_name)
.bind(cmd.partner_id)
.bind(&barcode)
.bind(born_state)
.bind(link.map(|l| l.sale_order_id))
.bind(link.map(|l| l.sale_order_state.as_str()))
.bind(link.map(|l| l.sale_status.as_str()))
.fetch_one(&mut **tx)
.await?
};
if born_state == "open" {
sqlx::query(
r#"UPDATE event.mails SET scheduled_date = now(), mail_done = false
WHERE event_id = $1 AND interval_kind = 'after_sub'"#,
)
.bind(cmd.event_id)
.execute(&mut **tx)
.await?;
}
if !cmd.lead_rule_skip {
sqlx::query(
r#"INSERT INTO event.lead_requests (event_id)
SELECT $1 WHERE EXISTS (
SELECT 1 FROM event.lead_rules lr
WHERE lr.active
AND (lr.event_id IS NULL OR lr.event_id = $1)
AND (lr.on_create OR lr.on_confirm))
ON CONFLICT (event_id) DO UPDATE SET done = false WHERE lead_requests.done"#,
)
.bind(cmd.event_id)
.execute(&mut **tx)
.await?;
}
crate::infrastructure::persistence::audit::record_audit(
&mut **tx,
"registration_created",
cmd.actor,
"registration",
Some(row.id),
serde_json::json!({
"event_id": cmd.event_id,
"email": cmd.email,
"state": born_state,
"sale_minted": link.is_some(),
}))
.await?;
Ok(row)
}
pub async fn seat_counts(
&self,
event_id: Uuid,
slot_id: Option<Uuid>,
) -> Result<SeatCounts, EventError> {
let event = company_scope::fetch_optional_scoped(&self.rpool(), sqlx::query_as::<_, LockedEvent>(
r#"SELECT id, is_multi_slots, seats_limited, seats_max, kanban_state::text AS kanban_state
FROM event.events WHERE id = $1"#,
)
.bind(event_id))
.await?
.ok_or(EventError::EventNotFound)?;
let taken: i64 = match slot_id {
Some(slot) => {
company_scope::fetch_one_scalar_scoped(
&self.rpool(),
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations
WHERE event_id = $1 AND event_slot_id = $2
AND state IN ('open','done') AND active"#,
)
.bind(event_id)
.bind(slot),
)
.await?
}
None => {
company_scope::fetch_one_scalar_scoped(
&self.rpool(),
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations
WHERE event_id = $1 AND state IN ('open','done') AND active"#,
)
.bind(event_id),
)
.await?
}
};
Ok(SeatCounts {
limited: event.seats_limited,
capacity: event.seats_max as i64,
taken,
})
}
pub async fn transition(
&self,
registration_id: Uuid,
to: &str,
actor: Option<Uuid>,
) -> Result<(String, String, bool), EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let outcome = sqlx::query_as::<_, (String, String, bool)>(
r#"WITH prev AS (
SELECT state FROM event.registrations WHERE id = $1 FOR UPDATE
)
UPDATE event.registrations r
SET state = $2::event_registration_state,
date_closed = CASE
WHEN $2 = 'done' AND r.date_closed IS NULL THEN now()
ELSE r.date_closed END
FROM prev
WHERE r.id = $1
RETURNING prev.state::text AS before_state,
r.state::text AS after_state,
r.active"#,
)
.bind(registration_id)
.bind(to)
.fetch_optional(&mut *tx)
.await?
.ok_or(EventError::RegistrationNotFound)?;
let (before, after, active) = &outcome;
let takes_seat = *active
&& (*after == "open" || *after == "done")
&& (*before == "draft" || *before == "cancel");
if takes_seat {
let (event_id, slot_id) = sqlx::query_as::<_, (Uuid, Option<Uuid>)>(
"SELECT event_id, event_slot_id FROM event.registrations WHERE id = $1",
)
.bind(registration_id)
.fetch_one(&mut *tx)
.await?;
let event = sqlx::query_as::<_, LockedEvent>(
r#"SELECT id, is_multi_slots, seats_limited, seats_max, kanban_state::text AS kanban_state
FROM event.events WHERE id = $1 FOR UPDATE"#,
)
.bind(event_id)
.fetch_optional(&mut *tx)
.await?
.ok_or(EventError::EventNotFound)?;
let taken: i64 = match slot_id {
Some(slot) => {
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations
WHERE event_id = $1 AND event_slot_id = $2
AND id != $3
AND state IN ('open','done') AND active"#,
)
.bind(event_id)
.bind(slot)
.bind(registration_id)
.fetch_one(&mut *tx)
.await?
}
None => {
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations
WHERE event_id = $1
AND id != $2
AND state IN ('open','done') AND active"#,
)
.bind(event_id)
.bind(registration_id)
.fetch_one(&mut *tx)
.await?
}
};
if event.seats_limited && event.seats_max > 0 && taken >= event.seats_max as i64 {
return Err(EventError::EventSeatsExhausted { event_id });
}
}
if *active && after == "open" && (before == "draft" || before == "cancel") {
sqlx::query(
r#"UPDATE event.mails m SET scheduled_date = now(), mail_done = false
WHERE m.event_id = (SELECT event_id FROM event.registrations WHERE id = $1)
AND m.interval_kind = 'after_sub'"#,
)
.bind(registration_id)
.execute(&mut *tx)
.await?;
}
if *active && (after == "open" || after == "done") && before != after {
let axis = if after == "done" {
"on_done"
} else {
"on_confirm"
};
sqlx::query(
r#"INSERT INTO event.lead_requests (event_id)
SELECT r.event_id FROM event.registrations r
WHERE r.id = $1
AND EXISTS (
SELECT 1 FROM event.lead_rules lr
WHERE lr.active
AND (lr.event_id IS NULL OR lr.event_id = r.event_id)
AND (CASE WHEN $2 = 'on_done' THEN lr.on_done ELSE lr.on_confirm END))
ON CONFLICT (event_id) DO UPDATE SET done = false WHERE lead_requests.done"#,
)
.bind(registration_id)
.bind(axis)
.execute(&mut *tx)
.await?;
}
if before != after {
crate::infrastructure::persistence::audit::record_audit(
&mut *tx,
"registration_state_changed",
actor,
"registration",
Some(registration_id),
serde_json::json!({ "before": before, "after": after }))
.await?;
}
tx.commit().await?;
Ok(outcome)
}
pub async fn sync_from_partner(
&self,
registration_id: Uuid,
partner_id: Uuid,
name: Option<&str>,
phone: Option<&str>,
company_name: Option<&str>,
actor: Option<Uuid>,
) -> Result<RegistrationRow, EventError> {
let row = company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, RegistrationRow>(
r#"UPDATE event.registrations SET
partner_id = COALESCE(partner_id, $2),
name = COALESCE(name, $3),
phone = COALESCE(phone, $4),
company_name = COALESCE(company_name, $5)
WHERE id = $1 AND active
RETURNING id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state::text AS state, date_closed,
sale_order_id, sale_order_state::text AS sale_order_state,
sale_status::text AS sale_status, active, barcode"#,
)
.bind(registration_id)
.bind(partner_id)
.bind(name)
.bind(phone)
.bind(company_name),
)
.await?
.ok_or(EventError::RegistrationNotFound)?;
record_audit(
&self.rpool(),
"registration_updated",
actor,
"registration",
registration_id,
serde_json::json!({ "verb": "sync_from_partner", "partner_id": partner_id }),
)
.await;
Ok(row)
}
pub async fn find_by_barcode(
&self,
barcode: &str,
) -> Result<Option<RegistrationRow>, EventError> {
company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, RegistrationRow>(
r#"SELECT id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state::text AS state, date_closed,
sale_order_id, sale_order_state::text AS sale_order_state,
sale_status::text AS sale_status, active, barcode
FROM event.registrations WHERE barcode = $1"#,
)
.bind(barcode),
)
.await
.map_err(EventError::from)
}
pub async fn find_registration(
&self,
registration_id: Uuid,
) -> Result<RegistrationRow, EventError> {
company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, RegistrationRow>(
r#"SELECT id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state::text AS state, date_closed,
sale_order_id, sale_order_state::text AS sale_order_state,
sale_status::text AS sale_status, active, barcode
FROM event.registrations WHERE id = $1"#,
)
.bind(registration_id),
)
.await?
.ok_or(EventError::RegistrationNotFound)
}
pub async fn list_registrations(
&self,
event_id: Uuid,
limit: i64,
after: Option<Uuid>,
) -> Result<Vec<RegistrationRow>, EventError> {
company_scope::fetch_all_scoped(
&self.rpool(),
sqlx::query_as::<_, RegistrationRow>(
r#"SELECT id, event_id, event_slot_id, event_ticket_id, name, email, phone,
company_name, partner_id, state::text AS state, date_closed,
sale_order_id, sale_order_state::text AS sale_order_state,
sale_status::text AS sale_status, active, barcode
FROM event.registrations
WHERE event_id = $1 AND ($2::uuid IS NULL OR id > $2)
ORDER BY id LIMIT $3"#,
)
.bind(event_id)
.bind(after)
.bind(limit),
)
.await
.map_err(EventError::from)
}
}