use backbone_orm::company_scope;
use chrono::{DateTime, Utc};
use sqlx::PgPool;
use uuid::Uuid;
use crate::application::service::event_error::EventError;
use super::seat_repository::record_audit;
#[derive(Debug, Clone, serde::Serialize, sqlx::FromRow)]
pub struct EventRow {
pub id: Uuid,
pub name: String,
pub event_type_id: Option<Uuid>,
pub stage_id: Uuid,
pub kanban_state: String,
pub date_begin: DateTime<Utc>,
pub date_end: DateTime<Utc>,
pub date_tz: String,
pub is_multi_slots: bool,
pub event_slot_count: i32,
pub seats_limited: bool,
pub seats_max: i32,
pub badge_format: String,
pub is_published: bool,
pub date_publish: Option<DateTime<Utc>>,
}
#[derive(Debug, Clone, Default)]
pub struct CreateEventInput {
pub name: String,
pub event_type_id: Option<Uuid>,
pub date_begin: DateTime<Utc>,
pub date_end: DateTime<Utc>,
pub date_tz: Option<String>,
pub is_multi_slots: bool,
pub event_slot_count: i32,
pub seats_limited: bool,
pub seats_max: i32,
pub organizer_id: Option<Uuid>,
pub user_id: Option<Uuid>,
pub address_id: Option<Uuid>,
pub event_url: Option<String>,
pub badge_format: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct PatchEventInput {
pub name: Option<String>,
pub event_type_id: Option<Uuid>,
pub stage_id: Option<Uuid>,
pub date_begin: Option<DateTime<Utc>>,
pub date_end: Option<DateTime<Utc>>,
pub date_tz: Option<String>,
pub is_multi_slots: Option<bool>,
pub event_slot_count: Option<i32>,
pub seats_limited: Option<bool>,
pub seats_max: Option<i32>,
pub organizer_id: Option<Uuid>,
pub user_id: Option<Uuid>,
pub address_id: Option<Uuid>,
pub event_url: Option<String>,
pub badge_format: Option<String>,
}
const EVENT_COLUMNS: &str =
"id, name, event_type_id, stage_id, kanban_state::text AS kanban_state, \
date_begin, date_end, date_tz, is_multi_slots, event_slot_count, seats_limited, seats_max, \
badge_format::text AS badge_format, is_published, date_publish";
pub struct EventCommandRepository {
pool: PgPool,
}
impl EventCommandRepository {
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 create(
&self,
input: &CreateEventInput,
actor: Option<Uuid>,
) -> Result<EventRow, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let stage_id = match sqlx::query_scalar::<_, Uuid>(
"SELECT id FROM event.stages ORDER BY sequence, id LIMIT 1",
)
.fetch_optional(&mut *tx)
.await?
{
Some(id) => id,
None => {
sqlx::query_scalar::<_, Uuid>(
"INSERT INTO event.stages (name, sequence) VALUES ('New', 1) RETURNING id",
)
.fetch_one(&mut *tx)
.await?
}
};
let id = Uuid::new_v4();
let row = sqlx::query_as::<_, EventRow>(&format!(
r#"INSERT INTO event.events
(id, name, event_type_id, stage_id, date_begin, date_end, date_tz,
is_multi_slots, event_slot_count, seats_limited, seats_max,
organizer_id, user_id, address_id, event_url, badge_format)
VALUES ($1, $2, $3, $4, $5, $6, COALESCE($7, 'UTC'), $8, $9, $10, $11,
$12, $13, $14, $15, COALESCE($16::event_badge_format, 'a4_french_fold'))
RETURNING {EVENT_COLUMNS}"#
))
.bind(id)
.bind(&input.name)
.bind(input.event_type_id)
.bind(stage_id)
.bind(input.date_begin)
.bind(input.date_end)
.bind(input.date_tz.as_deref())
.bind(input.is_multi_slots)
.bind(input.event_slot_count)
.bind(input.seats_limited)
.bind(input.seats_max)
.bind(input.organizer_id)
.bind(input.user_id)
.bind(input.address_id)
.bind(&input.event_url)
.bind(input.badge_format.as_deref())
.fetch_one(&mut *tx)
.await?;
if let Some(type_id) = input.event_type_id {
sqlx::query(
r#"INSERT INTO event.mails
(event_id, interval_nbr, interval_unit, interval_kind,
notification_channel, scheduled_date, template_ref, template_kind)
SELECT $1, tm.interval_nbr, tm.interval_unit, tm.interval_kind,
tm.notification_channel, now(), tm.template_ref, tm.template_kind
FROM event.type_mails tm WHERE tm.event_type_id = $2"#,
)
.bind(id)
.bind(type_id)
.execute(&mut *tx)
.await?;
sqlx::query(
r#"INSERT INTO event.booths (id, event_id, booth_category_id, name)
SELECT gen_random_uuid(), $1, tb.booth_category_id, tb.name
FROM event.type_booths tb WHERE tb.event_type_id = $2"#,
)
.bind(id)
.bind(type_id)
.execute(&mut *tx)
.await?;
}
crate::infrastructure::persistence::audit::record_audit(
&mut *tx,
"event_created",
actor,
"event",
Some(id),
serde_json::json!({ "name": input.name, "event_type_id": input.event_type_id }))
.await?;
tx.commit().await?;
Ok(row)
}
pub async fn patch(
&self,
id: Uuid,
patch: &PatchEventInput,
actor: Option<Uuid>,
) -> Result<EventRow, EventError> {
let row = company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, EventRow>(&format!(
r#"UPDATE event.events SET
name = COALESCE($2, name),
event_type_id = COALESCE($3, event_type_id),
stage_id = COALESCE($4, stage_id),
date_begin = COALESCE($5, date_begin),
date_end = COALESCE($6, date_end),
date_tz = COALESCE($7, date_tz),
is_multi_slots = COALESCE($8, is_multi_slots),
event_slot_count = COALESCE($9, event_slot_count),
seats_limited = COALESCE($10, seats_limited),
seats_max = COALESCE($11, seats_max),
organizer_id = COALESCE($12, organizer_id),
user_id = COALESCE($13, user_id),
address_id = COALESCE($14, address_id),
event_url = COALESCE($15, event_url),
badge_format = COALESCE($16::event_badge_format, badge_format)
WHERE id = $1
RETURNING {EVENT_COLUMNS}"#
))
.bind(id)
.bind(&patch.name)
.bind(patch.event_type_id)
.bind(patch.stage_id)
.bind(patch.date_begin)
.bind(patch.date_end)
.bind(patch.date_tz.as_deref())
.bind(patch.is_multi_slots)
.bind(patch.event_slot_count)
.bind(patch.seats_limited)
.bind(patch.seats_max)
.bind(patch.organizer_id)
.bind(patch.user_id)
.bind(patch.address_id)
.bind(&patch.event_url)
.bind(patch.badge_format.as_deref()),
)
.await?
.ok_or(EventError::EventNotFound)?;
record_audit(
&self.rpool(),
"event_updated",
actor,
"event",
id,
serde_json::json!({ "verb": "patch", "fenced": ["is_published", "date_publish"] }),
)
.await;
Ok(row)
}
pub async fn publish(&self, id: Uuid, actor: Option<Uuid>) -> Result<EventRow, EventError> {
let row = company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, EventRow>(&format!(
r#"UPDATE event.events
SET is_published = true,
date_publish = COALESCE(date_publish, now())
WHERE id = $1
RETURNING {EVENT_COLUMNS}"#
))
.bind(id),
)
.await?
.ok_or(EventError::EventNotFound)?;
record_audit(
&self.rpool(),
"event_published",
actor,
"event",
id,
serde_json::json!({}),
)
.await;
Ok(row)
}
pub async fn unpublish(&self, id: Uuid, actor: Option<Uuid>) -> Result<EventRow, EventError> {
let row = company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, EventRow>(&format!(
r#"UPDATE event.events SET is_published = false
WHERE id = $1
RETURNING {EVENT_COLUMNS}"#
))
.bind(id),
)
.await?
.ok_or(EventError::EventNotFound)?;
record_audit(
&self.rpool(),
"event_unpublished",
actor,
"event",
id,
serde_json::json!({}),
)
.await;
Ok(row)
}
pub async fn mark_done(&self, id: Uuid, actor: Option<Uuid>) -> Result<EventRow, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let pipe_end = sqlx::query_scalar::<_, Uuid>(
"SELECT id FROM event.stages WHERE pipe_end ORDER BY sequence, id LIMIT 1",
)
.fetch_optional(&mut *tx)
.await?;
let row = sqlx::query_as::<_, EventRow>(&format!(
r#"UPDATE event.events
SET kanban_state = 'done',
stage_id = COALESCE($2, stage_id)
WHERE id = $1
RETURNING {EVENT_COLUMNS}"#
))
.bind(id)
.bind(pipe_end)
.fetch_optional(&mut *tx)
.await?
.ok_or(EventError::EventNotFound)?;
crate::infrastructure::persistence::audit::record_audit(
&mut *tx,
"event_mark_done",
actor,
"event",
Some(id),
serde_json::json!({ "verb": "mark_done" }))
.await?;
tx.commit().await?;
Ok(row)
}
pub async fn find(&self, id: Uuid) -> Result<EventRow, EventError> {
company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, EventRow>(&format!(
"SELECT {EVENT_COLUMNS} FROM event.events WHERE id = $1"
))
.bind(id),
)
.await?
.ok_or(EventError::EventNotFound)
}
pub async fn list(&self, limit: i64) -> Result<Vec<EventRow>, EventError> {
company_scope::fetch_all_scoped(
&self.rpool(),
sqlx::query_as::<_, EventRow>(&format!(
"SELECT {EVENT_COLUMNS} FROM event.events ORDER BY date_begin DESC, id LIMIT $1"
))
.bind(limit),
)
.await
.map_err(EventError::from)
}
pub async fn list_linked_products(&self) -> Result<Vec<Uuid>, EventError> {
company_scope::fetch_all_scoped(
&self.rpool(),
sqlx::query_as::<_, (Uuid,)>(
"SELECT product_id FROM event.event_linked_products ORDER BY product_id",
),
)
.await
.map(|rows| rows.into_iter().map(|r| r.0).collect())
.map_err(EventError::from)
}
pub async fn find_published(&self, id: Uuid) -> Result<EventRow, EventError> {
let row = match self.find(id).await {
Ok(row) => row,
Err(EventError::EventNotFound) => return Err(EventError::EventNotPublished),
Err(other) => return Err(other),
};
if !row.is_published || row.kanban_state == "cancel" {
return Err(EventError::EventNotPublished);
}
Ok(row)
}
pub async fn slot_window_of_event(
&self,
slot_id: Uuid,
event_id: Uuid,
) -> Result<Option<(Uuid, DateTime<Utc>, DateTime<Utc>)>, EventError> {
company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, (Uuid, DateTime<Utc>, Option<DateTime<Utc>>)>(
"SELECT id, date_begin, date_end FROM event.slots WHERE id = $1 AND event_id = $2",
)
.bind(slot_id)
.bind(event_id),
)
.await
.map_err(EventError::from)
.map(|row| row.map(|(id, begin, end)| (id, begin, end.unwrap_or(begin))))
}
}