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
//! `SubscriptionWriteService` — the opt-out split verbs (hand-authored,
//! user-owned; see `metaphor.codegen.yaml`).
//!
//! The split, restated: `opt_out` is a FLIP, never a DELETE. Opting out
//! stamps `opt_out_datetime` on the DATABASE clock (optionally a reason);
//! re-subscribing flips back and clears both. Rows live forever — the
//! audience keeps its history, and the send engine's cross-list
//! OPT-OUT-WINS probe (one opt-out anywhere suppresses the contact
//! everywhere) reads the same rows.
//!
//! The archive verb carries the R-M7 guard: an audience cited by any
//! in-flight mailing's canonical domain refuses 409.
use uuid::Uuid;
use crate::infrastructure::persistence::subscription_repository::{
SubscriptionRepository, SubscriptionRow,
};
/// Typed error surface (house style).
#[derive(Debug, thiserror::Error)]
pub enum SubscriptionWriteError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("not found: {0}")]
NotFound(String),
#[error("conflict: {0}")]
Conflict(String),
#[error("invalid: {0}")]
Invalid(String),
#[error("audience is not open to self-service subscription")]
AudienceNotPublic,
}
impl SubscriptionWriteError {
pub fn code(&self) -> &'static str {
match self {
Self::Db(_) => "mailing_db_error",
Self::NotFound(_) => "not_found",
Self::Conflict(_) => "audience_in_use",
Self::Invalid(_) => "invalid_input",
// Deliberately ONE code for unknown, archived, and private —
// the self-service edge must not let a caller distinguish the
// three (the enumeration-oracle posture the whole public
// subscription surface carries).
Self::AudienceNotPublic => "audience_not_public",
}
}
pub fn http_status(&self) -> u16 {
match self {
Self::Db(_) => 500,
Self::NotFound(_) => 404,
Self::Conflict(_) => 409,
Self::Invalid(_) => 422,
Self::AudienceNotPublic => 422,
}
}
}
/// The subscription/audience verbs. Stateless over a pool; one transaction
/// per verb.
pub struct SubscriptionWriteService {
pool: sqlx::PgPool,
}
impl SubscriptionWriteService {
pub fn new(pool: sqlx::PgPool) -> Self {
Self { pool }
}
/// Create an audience.
pub async fn create_audience(&self, name: &str, is_public: bool) -> Result<Uuid, SubscriptionWriteError> {
if name.trim().is_empty() {
return Err(SubscriptionWriteError::Invalid("audience name is required".into()));
}
let id = Uuid::new_v4();
let mut tx = self.pool.begin().await?;
SubscriptionRepository::insert_audience(&mut tx, id, name, is_public).await?;
tx.commit().await?;
Ok(id)
}
/// Subscribe a contact (upserting the contact row by email) to an
/// audience. Idempotent — converges onto the existing live row AND
/// clears any standing opt-out (a re-subscribe is the only sanctioned
/// way out of an opt-out).
///
/// The ONE self-service guard: only PUBLIC audiences are subscribable
/// through this verb. The audience model declares the contract
/// ("self-service subscriptions only create onto public audiences" —
/// `is_public` gates the recipient-facing subscription-preferences
/// surface); this is where it is enforced, on the single sanctioned
/// subscribe verb, so every composed route inherits the same fence
/// instead of each path re-deriving it. Officer-side membership
/// management rides the gated generic CRUD, not this verb. A missing
/// or archived audience id refuses identically to a private one — the
/// caller learns nothing about WHICH ids exist.
pub async fn subscribe(
&self,
email: &str,
audience_id: Uuid,
name: Option<&str>,
) -> Result<SubscriptionRow, SubscriptionWriteError> {
let email = email.trim();
if !email.contains('@') {
return Err(SubscriptionWriteError::Invalid(format!(
"email must be an address: {email}"
)));
}
let mut tx = self.pool.begin().await?;
let is_public = SubscriptionRepository::audience_is_public(&mut tx, audience_id)
.await?
.unwrap_or(false);
if !is_public {
tx.rollback().await?;
return Err(SubscriptionWriteError::AudienceNotPublic);
}
SubscriptionRepository::ensure_default_reasons(&mut tx).await?;
let contact_id = SubscriptionRepository::upsert_contact(&mut tx, email, name).await?;
let row = SubscriptionRepository::subscribe(&mut tx, contact_id, audience_id).await?;
tx.commit().await?;
Ok(row)
}
/// The opt-out flip: stamps `opt_out_datetime` on the DB clock,
/// optionally with a reason. Returns the flipped row; a repeated
/// unsubscribe returns the standing row WITHOUT restamping (the FIRST
/// opt-out moment is the durable fact) — `already_out` tells them apart.
pub async fn opt_out(
&self,
contact_id: Uuid,
audience_id: Uuid,
reason_id: Option<Uuid>,
) -> Result<(SubscriptionRow, bool), SubscriptionWriteError> {
let mut tx = self.pool.begin().await?;
let flipped = SubscriptionRepository::opt_out(&mut tx, contact_id, audience_id, reason_id)
.await?;
tx.commit().await?;
match flipped {
Some(row) => Ok((row, true)),
None => {
let standing = self.find(contact_id, audience_id).await?;
match standing {
Some(row) => Ok((row, false)),
None => Err(SubscriptionWriteError::NotFound(format!(
"subscription for contact {contact_id} on audience {audience_id}"
))),
}
}
}
}
/// The unsubscribe ROUTE seam: arrives with an email (+ the audience
/// context from the trace). Falls back to the default reason when none
/// is picked.
///
/// Anti-oracle posture: every refusal arm — unknown contact, contact
/// with no subscription on this audience — answers the SAME typed
/// error. Distinguishable messages here would hand any caller a probe
/// for both "is this email a known contact" and "is this email on this
/// audience" — the mailing-list membership enumeration the public
/// subscription surface must never expose. The distinction stays in
/// the server log, never the answer.
pub async fn unsubscribe_by_email(
&self,
email: &str,
audience_id: Uuid,
reason_id: Option<Uuid>,
) -> Result<(SubscriptionRow, bool), SubscriptionWriteError> {
const REFUSAL: &str = "no subscription to unsubscribe";
let mut tx = self.pool.begin().await?;
SubscriptionRepository::ensure_default_reasons(&mut tx).await?;
let contact_id = match SubscriptionRepository::contact_id_by_email(&mut tx, email).await? {
Some(id) => id,
None => {
tx.rollback().await?;
tracing::debug!(email = %email, "unsubscribe refused: no such contact");
return Err(SubscriptionWriteError::NotFound(REFUSAL.to_string()));
}
};
let reason = match reason_id {
Some(r) => Some(r),
None => SubscriptionRepository::default_reason_id(&mut tx).await?,
};
let flipped = SubscriptionRepository::opt_out(&mut tx, contact_id, audience_id, reason)
.await?;
if flipped.is_none() {
let standing = SubscriptionRepository::find_live(&mut tx, contact_id, audience_id).await?;
if standing.is_none() {
tx.rollback().await?;
tracing::debug!(
email = %email,
audience_id = %audience_id,
"unsubscribe refused: no subscription on this audience"
);
return Err(SubscriptionWriteError::NotFound(REFUSAL.to_string()));
}
tx.commit().await?;
return Ok((standing.expect("checked above"), false));
}
tx.commit().await?;
Ok((flipped.expect("checked above"), true))
}
/// One live subscription row (the read side).
pub async fn find(
&self,
contact_id: Uuid,
audience_id: Uuid,
) -> Result<Option<SubscriptionRow>, SubscriptionWriteError> {
let mut tx = self.pool.begin().await?;
let row = SubscriptionRepository::find_live(&mut tx, contact_id, audience_id).await?;
tx.commit().await?;
Ok(row)
}
/// Admin removal — soft-deletes the membership (distinct from the
/// opt-out flip: removes the row from the audience entirely; audit
/// trail stays).
pub async fn remove(
&self,
contact_id: Uuid,
audience_id: Uuid,
) -> Result<(), SubscriptionWriteError> {
let mut tx = self.pool.begin().await?;
let removed = SubscriptionRepository::remove(&mut tx, contact_id, audience_id).await?;
tx.commit().await?;
if !removed {
return Err(SubscriptionWriteError::NotFound(format!(
"subscription for contact {contact_id} on audience {audience_id}"
)));
}
Ok(())
}
/// Archive an audience — R-M7 guarded: refuses 409 while any in-flight
/// mailing's canonical domain cites the audience. Returns the live
/// member count it archived with.
pub async fn archive_audience(
&self,
audience_id: Uuid,
) -> Result<i64, SubscriptionWriteError> {
let mut tx = self.pool.begin().await?;
let in_use = SubscriptionRepository::audience_in_use_by_mailings(&mut tx, audience_id)
.await?;
if !in_use.is_empty() {
tx.rollback().await?;
return Err(SubscriptionWriteError::Conflict(format!(
"audience {audience_id} is cited by {} in-flight mailing(s): {}",
in_use.len(),
in_use
.iter()
.map(|id| id.to_string())
.collect::<Vec<_>>()
.join(", ")
)));
}
let members = SubscriptionRepository::active_member_count(&mut tx, audience_id).await?;
let archived = SubscriptionRepository::soft_delete_audience(&mut tx, audience_id).await?;
tx.commit().await?;
if !archived {
return Err(SubscriptionWriteError::NotFound(format!(
"audience {audience_id}"
)));
}
Ok(members)
}
}
#[cfg(test)]
mod tests {
// Behavior coverage (fresh-DB subscribe/opt-out flips, OPT-OUT-WINS at
// send, R-M7 refusal) lives in tests/subscription_cases.rs and
// tests/behavior/ — the pure parse/unit surface here is thin by design.
}