backbone-events 0.5.1

Events core + the sale/booth/crm/sms/desk overlay — registrations under one authoritative seat counter, the sale seam (order-minted registrations, forward-only paid heal, SO-cancel mirror cascade), booth bookings with the DB exclusivity wall, closed-vocabulary lead rules through a host sink port, the dual-channel (mail/sms) self-arming scheduler, and the exact-match desk verb (Odoo event port, core + overlay)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
//! The event repository (hand-written; user-owned; see
//! `metaphor.codegen.yaml`).
//!
//! The transactional SQL of the event verbs: create (with the ONE
//! template-apply — the event type's type_mail rows are copied into
//! event-scoped scheduler rows exactly once, at create, and never
//! re-propagate), the PATCH whitelist application (the FENCE itself
//! is enforced in the service — this repo has no arm that can write
//! `is_published`/`date_publish`), publish/unpublish (the ONLY
//! writers of the fence pair), mark_done (writes the first pipe_end
//! stage by sequence), and the reads the capability surface needs.

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;

/// The event row shape the verbs return.
#[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>>,
}

/// The create input (officer verbs + probes; every field explicit).
#[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>,
}

/// The patch input — every arm OPTIONAL; `None` = leave untouched.
/// The publication fence pair is DELIBERATELY ABSENT: a repo that
/// cannot write the pair cannot be tricked into writing it.
#[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
    }

    /// The database this call runs on: the composer's request pool when one
    /// is bound (a tenant mount, or a relay consumer wrapped by the host),
    /// else the composed pool (ADR-0029 pool law).
    fn rpool(&self) -> PgPool {
        crate::request_pool::current().unwrap_or_else(|| self.pool.clone())
    }

    /// Create an event + apply the type's mail templates EXACTLY ONCE
    /// (never re-propagates — a later type change does not fork
    /// scheduler rows onto existing events).
    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?;

        // THE ONE TEMPLATE-APPLY: copy the type's scheduler templates
        // as event-scoped scheduler rows, and the type's booth rows as
        // event booths (WHITELIST: name + booth_category_id only — the
        // type never templates booking state, contacts, or sale
        // links). Runs ONLY here; never re-propagates.
        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)
    }

    /// Apply a whitelisted patch (COALESCE semantics: absent arms
    /// leave the stored value). The fence pair is not patchable —
    /// not by this repo, not by any caller of it.
    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)
    }

    /// PUBLISH — the only writer that sets `is_published = true`.
    /// `date_publish` stamps now() on FIRST publish and is never
    /// rewritten on republish.
    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)
    }

    /// UNPUBLISH — the only writer that clears `is_published`
    /// (`date_publish` keeps the historical first-publish stamp).
    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)
    }

    /// MARK DONE — the verb: sets `done` AND moves the row to the
    /// FIRST pipe_end stage by sequence (upstream's two-axis close).
    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)
    }

    /// Fetch one event 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)
    }

    /// List events (newest first), capped.
    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)
    }

    /// THE CATALOG READ (EP-2): the distinct product ids linked across
    /// the event family (tickets + booth categories), through the
    /// `event_linked_products` view. Events owns the linkage; the
    /// catalog consumes only this.
    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)
    }

    /// The publication-checked event read for the capability surface:
    /// only a PUBLISHED, non-cancelled row resolves; unpublished,
    /// cancelled AND missing are all the ONE uniform not-published
    /// refusal (no oracle distinguishing members).
    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)
    }

    /// A slot's hours, but ONLY when the slot belongs to the named
    /// event (the ics slot-scope check + window read in one).
    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))))
    }
}