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
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
//! The scheduler repository (hand-written; user-owned; see
//! `metaphor.codegen.yaml`).
//!
//! The transactional SQL of the self-arming communication scheduler
//! (ADR-0020): the SKIP LOCKED claim domain, the LAZY receipt
//! materialization (receipts are created per pass, chunked and
//! capped, for registrations entering the eligible set), the receipt
//! walk, the receipt-truth completion recompute, the cancellation
//! propagation, and the daily mark-done sweep.
//!
//! `mail_done` is NOT a hand-set flag: it is the RECEIPT-TRUTH
//! recompute — "every eligible registration has a sent (or visibly
//! dropped) receipt" — so a LATE registrant re-opens the scheduler
//! automatically (EBB-2). Eligibility is always the pair
//! `state IN ('open','done') AND active` (EVM2-4).
//!
//! `'sent' = queued`: a successful enqueue sets `mail_sent` from
//! this module's point of view; delivery belongs to the transport.
use backbone_orm::{company_scope, org_scope};
use chrono::{DateTime, Utc};
use sqlx::PgPool;
use uuid::Uuid;
use crate::application::service::event_error::EventError;
use super::seat_repository::record_audit;
/// One claimed scheduler row.
#[derive(Debug, Clone, serde::Serialize, sqlx::FromRow)]
pub struct SchedulerRow {
pub id: Uuid,
pub event_id: Uuid,
pub interval_nbr: i32,
pub interval_unit: String,
pub interval_kind: String,
pub notification_channel: String,
pub template_ref: Option<Uuid>,
pub template_kind: Option<String>,
pub mail_done: bool,
pub last_registration_id: Option<Uuid>,
}
/// One due receipt joined to its registration (the render context's
/// arms — the phone arms the sms channel; the mail arm reads the
/// email).
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct DueReceipt {
pub receipt_id: Uuid,
pub registration_id: Uuid,
pub attendee_name: String,
pub attendee_email: String,
pub attendee_phone: Option<String>,
pub barcode: String,
pub scheduled_date: Option<DateTime<Utc>>,
}
/// The scheduler repository.
pub struct SchedulerRepository {
pool: PgPool,
}
impl SchedulerRepository {
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())
}
/// The claim domain: scheduler rows NOT done, DUE, on non-cancel
/// events, with WORK REMAINING (an eligible registration without
/// a sent receipt).
///
/// `FOR UPDATE SKIP LOCKED` selects a disjoint set for two passes
/// that overlap IN THIS STATEMENT, and nothing beyond it: the row
/// locks end with the statement, so the pass that walks the row
/// afterwards holds no claim on it. That is deliberate, because
/// the walk enqueues mail through a host port and must not run
/// inside a transaction held open across it. What actually keeps
/// a concurrent pass from repeating work is the receipt grain:
/// `UNIQUE(scheduler_id, registration_id)` on materialization,
/// and every send marked `AND NOT mail_sent`. The contract is
/// at-least-once, so two passes racing the same row can enqueue a
/// duplicate; they cannot lose one, and they cannot double-count
/// a receipt.
pub async fn claim_due(&self, limit: i64) -> Result<Vec<SchedulerRow>, EventError> {
company_scope::fetch_all_scoped(
&self.rpool(),
sqlx::query_as::<_, SchedulerRow>(
r#"WITH claimed AS (
SELECT m.id FROM event.mails m
JOIN event.events e ON e.id = m.event_id
WHERE NOT m.mail_done
AND m.scheduled_date <= now()
AND e.kanban_state::text <> 'cancel'
AND EXISTS (
SELECT 1 FROM event.registrations r
WHERE r.event_id = m.event_id
AND r.state IN ('open','done') AND r.active
AND NOT EXISTS (
SELECT 1 FROM event.mail_registrations mr
WHERE mr.scheduler_id = m.id
AND mr.registration_id = r.id
AND mr.mail_sent))
ORDER BY m.scheduled_date
LIMIT $1
FOR UPDATE SKIP LOCKED
)
SELECT m.id, m.event_id, m.interval_nbr, m.interval_unit::text AS interval_unit,
m.interval_kind::text AS interval_kind,
m.notification_channel::text AS notification_channel,
m.template_ref, m.template_kind,
m.mail_done, m.last_registration_id
FROM event.mails m JOIN claimed c ON c.id = m.id"#,
)
.bind(limit),
)
.await
.map_err(EventError::from)
}
/// LAZY receipt materialization: create the MISSING receipts for
/// eligible registrations, capped per pass. The idempotence
/// predicate is the ANTI-JOIN (`NOT EXISTS` a receipt for the
/// pair) backed by the UNIQUE(scheduler_id, registration_id)
/// constraint — materialized rows drop out of the missing set, so
/// the anti-join self-paginates under the cap. The
/// `last_registration_id` column is kept as a WATERMARK of the
/// highest materialized id and is NEVER a filter: registration ids
/// are random UUIDs, so a LATE registrant can sort below any
/// watermark (a cursor filter would strand it forever).
/// `scheduled_date` derives from the registration's creation stamp
/// + the row's interval (`now` = due immediately).
pub async fn materialize_receipts(
&self,
scheduler: &SchedulerRow,
cap: i64,
) -> Result<Vec<Uuid>, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let ids = sqlx::query_scalar::<_, Uuid>(
r#"INSERT INTO event.mail_registrations (scheduler_id, registration_id, scheduled_date)
SELECT m.id, r.id,
(r.metadata->>'created_at')::timestamptz +
CASE WHEN m.interval_unit::text = 'now'
THEN '0 seconds'::interval
ELSE (m.interval_nbr::text || ' ' || m.interval_unit::text)::interval
END
FROM event.mails m
JOIN event.registrations r ON r.event_id = m.event_id
WHERE m.id = $1
AND r.state IN ('open','done') AND r.active
AND NOT EXISTS (
SELECT 1 FROM event.mail_registrations x
WHERE x.scheduler_id = m.id AND x.registration_id = r.id)
ORDER BY r.id
LIMIT $2
RETURNING registration_id"#,
)
.bind(scheduler.id)
.bind(cap)
.fetch_all(&mut *tx)
.await?;
if let Some(last) = ids.last() {
sqlx::query(
"UPDATE event.mails SET last_registration_id = GREATEST(last_registration_id, $2) WHERE id = $1",
)
.bind(scheduler.id)
.bind(last)
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
Ok(ids)
}
/// How many eligible registrations still lack ANY receipt for
/// this scheduler (the overflow signal for the re-arm rule).
pub async fn unmaterialized_count(&self, scheduler_id: Uuid) -> Result<i64, EventError> {
company_scope::fetch_one_scalar_scoped(
&self.rpool(),
sqlx::query_scalar::<_, i64>(
r#"SELECT count(*) FROM event.registrations r
WHERE r.event_id = (SELECT event_id FROM event.mails WHERE id = $1)
AND r.state IN ('open','done') AND r.active
AND NOT EXISTS (
SELECT 1 FROM event.mail_registrations x
WHERE x.scheduler_id = $1 AND x.registration_id = r.id)"#,
)
.bind(scheduler_id),
)
.await
.map_err(EventError::from)
}
/// The due receipts (scheduled, unsent, eligible). A NULL
/// scheduled_date counts as due (the `now` unit).
pub async fn due_receipts(
&self,
scheduler_id: Uuid,
limit: i64,
) -> Result<Vec<DueReceipt>, EventError> {
company_scope::fetch_all_scoped(&self.rpool(), sqlx::query_as::<_, DueReceipt>(
r#"SELECT mr.id AS receipt_id, mr.registration_id, r.name AS attendee_name,
r.email AS attendee_email, r.phone AS attendee_phone, r.barcode, mr.scheduled_date
FROM event.mail_registrations mr
JOIN event.registrations r ON r.id = mr.registration_id
WHERE mr.scheduler_id = $1
AND NOT mr.mail_sent
AND (mr.scheduled_date IS NULL OR mr.scheduled_date <= now())
AND r.state IN ('open','done') AND r.active
ORDER BY mr.scheduled_date NULLS FIRST, mr.registration_id
LIMIT $2"#,
)
.bind(scheduler_id)
.bind(limit))
.await
.map_err(EventError::from)
}
/// Mark one receipt `'sent' = queued` (the idempotent send guard:
/// only the first writer flips the flag).
pub async fn mark_receipt_queued(&self, receipt_id: Uuid) -> Result<(), EventError> {
org_scope::execute_scoped(&self.rpool(), sqlx::query(
"UPDATE event.mail_registrations SET mail_sent = true, outcome = 'queued' WHERE id = $1 AND NOT mail_sent",
)
.bind(receipt_id))
.await?;
Ok(())
}
/// The VISIBLE drop (the window closed before the send): the
/// receipt closes with the drop outcome rather than silently
/// disappearing. Terminal — never retried.
pub async fn mark_receipt_dropped_window_closed(
&self,
receipt_id: Uuid,
) -> Result<(), EventError> {
org_scope::execute_scoped(&self.rpool(), sqlx::query(
"UPDATE event.mail_registrations SET mail_sent = true, outcome = 'dropped_window_closed' WHERE id = $1 AND NOT mail_sent",
)
.bind(receipt_id))
.await?;
Ok(())
}
/// Record a typed scheduler failure (the family: template
/// unresolved / renderer not composed / render failed / enqueue
/// refused / recipient invalid). Recorded and CONTINUED — never a
/// registration blocker, never a throttle.
pub async fn record_failure(&self, scheduler_id: Uuid, kind: &str) -> Result<(), EventError> {
org_scope::execute_scoped(&self.rpool(), sqlx::query(
"UPDATE event.mails SET error_kind = $2::event_mail_error_kind, error_datetime = now() WHERE id = $1",
)
.bind(scheduler_id)
.bind(kind))
.await?;
Ok(())
}
/// Clear the failure marker after a fully successful pass.
pub async fn clear_failure(&self, scheduler_id: Uuid) -> Result<(), EventError> {
org_scope::execute_scoped(
&self.rpool(),
sqlx::query(
"UPDATE event.mails SET error_kind = NULL, error_datetime = NULL WHERE id = $1",
)
.bind(scheduler_id),
)
.await?;
Ok(())
}
/// THE RECEIPT-TRUTH COMPLETION: `mail_done` is true iff every
/// eligible registration carries a sent (or dropped) receipt. A
/// late registrant re-opens the row automatically.
pub async fn recompute_mail_done(&self, scheduler_id: Uuid) -> Result<bool, EventError> {
let done: bool = company_scope::fetch_optional_scalar_scoped(
&self.rpool(),
sqlx::query_scalar::<_, bool>(
r#"UPDATE event.mails m
SET mail_done = NOT EXISTS (
SELECT 1 FROM event.registrations r
WHERE r.event_id = m.event_id
AND r.state IN ('open','done') AND r.active
AND NOT EXISTS (
SELECT 1 FROM event.mail_registrations mr
WHERE mr.scheduler_id = m.id
AND mr.registration_id = r.id
AND mr.mail_sent))
WHERE m.id = $1
RETURNING mail_done"#,
)
.bind(scheduler_id),
)
.await?
.unwrap_or(true);
Ok(done)
}
/// The overflow re-arm: a pass hit its cap with work left.
pub async fn rearm(&self, scheduler_id: Uuid) -> Result<(), EventError> {
org_scope::execute_scoped(
&self.rpool(),
sqlx::query(
"UPDATE event.mails SET scheduled_date = now() WHERE id = $1 AND NOT mail_done",
)
.bind(scheduler_id),
)
.await?;
Ok(())
}
/// Cancellation propagation: pending (unsent) receipts of
/// cancelled registrations are deleted on the next pass — the
/// audit row is the durable trace.
pub async fn propagate_cancellations(&self, limit: i64) -> Result<i64, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let deleted = sqlx::query_scalar::<_, Uuid>(
r#"DELETE FROM event.mail_registrations mr
WHERE mr.id IN (
SELECT x.id FROM event.mail_registrations x
JOIN event.registrations r ON r.id = x.registration_id
WHERE r.state = 'cancel' AND NOT x.mail_sent
LIMIT $1)
RETURNING mr.scheduler_id"#,
)
.bind(limit)
.fetch_all(&mut *tx)
.await?;
let n = deleted.len() as i64;
if n > 0 {
record_audit(
&self.rpool(),
"scheduler_run",
None,
"mail_scheduler",
Uuid::nil(),
serde_json::json!({
"verb": "cancellation_propagation",
"deleted_pending_receipts": n,
}),
)
.await;
}
tx.commit().await?;
Ok(n)
}
/// THE TEMPLATE CASCADE (EVM2-2 + ESM-1, collapsed into ONE verb):
/// a deleted template takes its dependent scheduler rows AND the
/// type-level template rows of the same (kind, ref) pair in one
/// set-based transaction. Covers BOTH channels — the pair is
/// (template_kind, template_ref) and the channel derives from the
/// kind, so mail and sms deps fall to the same verb (the declared
/// deviation from upstream's literal DB-level ON DELETE CASCADE is
/// recorded in docs/spec-overlay.md: the module holds no FK across
/// the template store boundary, so the sweep is a verb the template
/// side declares and calls, not a constraint that fires untyped).
/// Returns (scheduler rows deleted, type template rows deleted).
pub async fn cascade_delete_template_dependents(
&self,
template_kind: &str,
template_ref: Uuid,
actor: Option<Uuid>,
) -> Result<(i64, i64), EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let mails = sqlx::query_scalar::<_, Uuid>(
r#"DELETE FROM event.mails
WHERE template_ref = $1 AND template_kind = $2
RETURNING id"#,
)
.bind(template_ref)
.bind(template_kind)
.fetch_all(&mut *tx)
.await?;
let type_mails = sqlx::query_scalar::<_, Uuid>(
r#"DELETE FROM event.type_mails
WHERE template_ref = $1 AND template_kind = $2
RETURNING id"#,
)
.bind(template_ref)
.bind(template_kind)
.fetch_all(&mut *tx)
.await?;
let counts = (mails.len() as i64, type_mails.len() as i64);
record_audit(
&self.rpool(),
"template_cascade",
actor,
"mail_scheduler",
template_ref,
serde_json::json!({
"template_kind": template_kind,
"mails_deleted": counts.0,
"type_mails_deleted": counts.1,
}),
)
.await;
tx.commit().await?;
Ok(counts)
}
/// Officer read: one scheduler row.
pub async fn find(&self, scheduler_id: Uuid) -> Result<SchedulerRow, EventError> {
company_scope::fetch_optional_scoped(
&self.rpool(),
sqlx::query_as::<_, SchedulerRow>(
r#"SELECT id, event_id, interval_nbr, interval_unit::text AS interval_unit,
interval_kind::text AS interval_kind,
notification_channel::text AS notification_channel,
template_ref, template_kind,
mail_done, last_registration_id
FROM event.mails WHERE id = $1"#,
)
.bind(scheduler_id),
)
.await?
.ok_or(EventError::RegistrationNotFound)
.map_err(|e| match e {
EventError::RegistrationNotFound => EventError::EventNotFound,
other => other,
})
}
/// Officer read: an event's scheduler rows.
pub async fn list_for_event(&self, event_id: Uuid) -> Result<Vec<SchedulerRow>, EventError> {
company_scope::fetch_all_scoped(
&self.rpool(),
sqlx::query_as::<_, SchedulerRow>(
r#"SELECT id, event_id, interval_nbr, interval_unit::text AS interval_unit,
interval_kind::text AS interval_kind,
notification_channel::text AS notification_channel,
template_ref, template_kind,
mail_done, last_registration_id
FROM event.mails WHERE event_id = $1 ORDER BY id"#,
)
.bind(event_id),
)
.await
.map_err(EventError::from)
}
/// The daily done sweep (bounded batches, FOR UPDATE SKIP LOCKED
/// claims): past events not already done/cancelled and not
/// resting in a pipe_end stage move to `done` — the verb, run as
/// a sweep. Returns the swept ids.
pub async fn sweep_mark_done(&self, limit: i64) -> Result<Vec<Uuid>, EventError> {
let mut tx = self.rpool().begin().await?;
super::relay_ambient_scope(&mut tx).await?;
let ids = sqlx::query_scalar::<_, Uuid>(
r#"WITH claimed AS (
SELECT e.id FROM event.events e
WHERE e.kanban_state::text NOT IN ('done','cancel')
AND e.date_end < now()
AND NOT EXISTS (
SELECT 1 FROM event.stages s
WHERE s.id = e.stage_id AND s.pipe_end)
ORDER BY e.date_end
LIMIT $1
FOR UPDATE SKIP LOCKED
)
UPDATE event.events ev SET kanban_state = 'done'
FROM claimed c WHERE ev.id = c.id
RETURNING ev.id"#,
)
.bind(limit)
.fetch_all(&mut *tx)
.await?;
for id in &ids {
crate::infrastructure::persistence::audit::record_audit(
&mut *tx,
"event_mark_done",
None,
"event",
Some(*id),
serde_json::json!({ "verb": "done_sweep" }))
.await?;
}
tx.commit().await?;
Ok(ids)
}
}