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
//! Repository for the sms send queue (hand-written; user-owned).
//!
//! Holds the SQL for [`crate::application::service::SmsWriteService`]: enqueue with
//! uuid correlation (SM-M1/SM-M21 — NO FK, the tracker outlives the sms row's GC),
//! the MMB-4 `FOR UPDATE SKIP LOCKED` drain claim (one claim per row per pass;
//! concurrent workers drain disjoint sets; a crashed worker's claim expires with the
//! transaction), and the state-guarded outcome advances.
use sqlx::{PgConnection, Row};
use uuid::Uuid;
/// A claimed sms row, as the drainer holds it between claim and outcome.
pub struct SmsQueueRow {
pub id: Uuid,
pub uuid: String,
pub number: String,
pub body: String,
pub mail_message_id: Option<Uuid>,
}
/// Hand-written sms queue SQL.
pub struct SmsQueueRepository;
impl SmsQueueRepository {
pub fn new() -> Self {
Self
}
}
impl Default for SmsQueueRepository {
fn default() -> Self {
Self::new()
}
}
impl SmsQueueRepository {
/// Enqueue an sms for send: mint the `sms` row (`state='outgoing'`) and its
/// uuid-correlated `sms_tracker`. Runs on the caller's open transaction so the
/// row, the tracker, and the `SmsCreated` outbox event (TR-SM-1 queue re-arm)
/// commit atomically.
pub async fn enqueue(
conn: &mut PgConnection,
id: Uuid,
sms_uuid: &str,
number: &str,
body: &str,
mail_message_id: Option<Uuid>,
notification_id: Option<Uuid>,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"INSERT INTO messaging.sms (id, uuid, number, body, state, mail_message_id)
VALUES ($1,$2,$3,$4,'outgoing'::sms_state,$5)"#,
)
.bind(id)
.bind(sms_uuid)
.bind(number)
.bind(body)
.bind(mail_message_id)
.execute(&mut *conn)
.await?;
sqlx::query(
r#"INSERT INTO messaging.sms_trackers (sms_uuid, message_id, notification_id, state)
VALUES ($1,$2,$3,'process'::mail_notification_status)
ON CONFLICT (sms_uuid) DO NOTHING"#,
)
.bind(sms_uuid)
.bind(mail_message_id)
.bind(notification_id)
.execute(&mut *conn)
.await?;
Ok(())
}
/// MMB-4 drain claim: `outgoing → process` over the batch, one statement, with
/// `FOR UPDATE SKIP LOCKED` inside the claim's subselect. Concurrent drainers
/// claim disjoint sets; a drainer that crashes mid-claim leaves no phantom
/// claims (the lock and the row-state revert with the transaction).
pub async fn claim_batch_for_drain(
conn: &mut PgConnection,
batch: i64,
) -> Result<Vec<SmsQueueRow>, sqlx::Error> {
let rows = sqlx::query(
r#"UPDATE messaging.sms AS s
SET state = 'process'::sms_state
WHERE s.id IN (
SELECT id FROM messaging.sms
WHERE state = 'outgoing'::sms_state
ORDER BY id
LIMIT $1
FOR UPDATE SKIP LOCKED
)
RETURNING s.id, s.uuid, s.number, s.body, s.mail_message_id"#,
)
.bind(batch)
.fetch_all(&mut *conn)
.await?;
Ok(rows
.iter()
.map(|r| SmsQueueRow {
id: r.get("id"),
uuid: r.get("uuid"),
number: r.get("number"),
body: r.get("body"),
mail_message_id: r.get("mail_message_id"),
})
.collect())
}
/// State-guarded outcome advance for a drained row: `process → pending` (LABEL
/// 'Sent' — accepted, awaiting delivery report), `process → sent`, or
/// `process → error` (failure_type set). The state guard (`AND state='process'`)
/// plus the SM-B6 trigger make replays and regressions impossible; `Ok(false)`
/// = the row was already advanced (a replayed result).
#[allow(clippy::too_many_arguments)]
pub async fn apply_outcome(
conn: &mut PgConnection,
id: Uuid,
target_state: &str,
failure_type: Option<&str>,
error_message: Option<&str>,
iap_status_code: Option<i32>,
) -> Result<bool, sqlx::Error> {
let updated = sqlx::query_scalar::<_, Uuid>(
r#"UPDATE messaging.sms
SET state = $2::sms_state,
failure_type = $3::sms_failure_type,
error_message = $4,
iap_status_code = $5
WHERE id = $1 AND state = 'process'::sms_state
RETURNING id"#,
)
.bind(id)
.bind(target_state)
.bind(failure_type)
.bind(error_message)
.bind(iap_status_code)
.fetch_optional(&mut *conn)
.await?;
Ok(updated.is_some())
}
/// The advance-seam outcome write: same shape as [`Self::apply_outcome`],
/// but the state guard is a legal-source SET instead of the drainer's
/// single `'process'` arm. `advance_state` (the webhook / late-verdict
/// entry) uses this with `['process','pending']` so the DELIVERY REPORT
/// transition `pending → sent` — the documented contract of that seam —
/// actually lands; the drainer itself keeps calling `apply_outcome`, whose
/// single-source guard is exactly right between claim and first outcome.
///
/// `AND state <> $2` keeps a re-assertion of the CURRENT state a replay
/// (`Ok(false)`): a provider retry of a verdict the row already carries
/// matches zero rows, exactly like the drainer's guard does.
#[allow(clippy::too_many_arguments)]
pub async fn apply_outcome_from(
conn: &mut PgConnection,
id: Uuid,
target_state: &str,
failure_type: Option<&str>,
error_message: Option<&str>,
iap_status_code: Option<i32>,
from_states: &[&str],
) -> Result<bool, sqlx::Error> {
let updated = sqlx::query_scalar::<_, Uuid>(
r#"UPDATE messaging.sms
SET state = $2::sms_state,
failure_type = $3::sms_failure_type,
error_message = $4,
iap_status_code = $5
WHERE id = $1
AND state = ANY($6::sms_state[])
AND state <> $2::sms_state
RETURNING id"#,
)
.bind(id)
.bind(target_state)
.bind(failure_type)
.bind(error_message)
.bind(iap_status_code)
.bind(from_states)
.fetch_optional(&mut *conn)
.await?;
Ok(updated.is_some())
}
/// Set the tracker's state to mirror the sms row's (the tracker is the
/// notification pump's bridge; its `state` column starts 'process').
pub async fn mirror_tracker_state(
conn: &mut PgConnection,
sms_uuid: &str,
tracker_state: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"UPDATE messaging.sms_trackers
SET state = $2::mail_notification_status
WHERE sms_uuid = $1"#,
)
.bind(sms_uuid)
.bind(tracker_state)
.execute(&mut *conn)
.await?;
Ok(())
}
/// Count rows still queued (drain-progress reporting).
pub async fn count_outgoing(conn: &mut PgConnection) -> Result<i64, sqlx::Error> {
sqlx::query_scalar::<_, i64>(r#"SELECT COUNT(*) FROM messaging.sms WHERE state = 'outgoing'::sms_state"#)
.fetch_one(&mut *conn)
.await
}
}