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
//! The email [`IngestSink`] — the durable, transactional half of email ingress.
//!
//! The generic source envelope ([`fraiseql_functions::run_source_once`]) hands each
//! polled batch here. In one transaction this emits every message onto the durable
//! spine (dedup by `Message-ID`) and advances the cursor — atomically, so a crash
//! leaves writes and watermark consistent. Once committed, it runs
//! the email-specific post-ingest work for each *genuinely new* message: delivery
//! correlation (bounces / challenges → send-status / suppression) first — it is not
//! idempotent, so a redelivery must not re-run it — then the `after:ingest:email`
//! function dispatch.
use std::sync::Arc;
use fraiseql_error::{FraiseQLError, Result};
use fraiseql_functions::{InboundMessage, IngestSink, PullBatch};
use fraiseql_observers::{CursorSnapshot, PostgresSourceCursorStore};
use sqlx::PgPool;
use tracing::warn;
use super::tracking::SendCorrelator;
use crate::{
inbound::spine::emit_in_tx,
routes::after_mutation::{plan_after_ingest_dispatch, spawn_after_ingest},
subsystems::BeforeMutationHooks,
};
#[cfg(test)]
mod tests;
/// The durable email ingest sink for one mailbox.
pub struct EmailIngestSink {
/// Stable mailbox identity — for log context.
mailbox_key: String,
/// Pool the ingest transaction is opened on.
pool: PgPool,
/// The generic cursor store (advanced in the ingest transaction).
cursor_store: PostgresSourceCursorStore,
/// Function-dispatch hooks; `None` ingests without firing `after:ingest`.
hooks: Option<Arc<BeforeMutationHooks>>,
/// The `fraiseql_query` bridge builder (#594) for `after:ingest` functions —
/// the request-path executor factory. `None` → an after:ingest function's
/// `fraiseql_query` fails loud (pre-#594 behavior).
query_executor_factory: Option<crate::routes::after_mutation::QueryExecutorFactory>,
/// Delivery-feedback correlator; `None` ingests without correlating inbound
/// bounces / challenges / replies to their send.
correlator: Option<Arc<dyn SendCorrelator>>,
/// The recipient address-hash key for suppression writes (needs the server HMAC
/// secret); `None` transitions status without writing suppressions.
address_hash_key: Option<Arc<[u8]>>,
/// The per-recipient unanswered-challenge suppression threshold (`N`).
challenge_suppress_after: u32,
}
impl EmailIngestSink {
/// Assemble a sink from a pool and its resolved collaborators.
#[must_use]
// Reason: a sink's fixed collaborators; a params struct would relocate them without
// reducing coupling (mirrors the poller/cron constructors).
#[allow(clippy::too_many_arguments)]
pub fn new(
mailbox_key: impl Into<String>,
pool: PgPool,
hooks: Option<Arc<BeforeMutationHooks>>,
query_executor_factory: Option<crate::routes::after_mutation::QueryExecutorFactory>,
correlator: Option<Arc<dyn SendCorrelator>>,
address_hash_key: Option<Arc<[u8]>>,
challenge_suppress_after: u32,
) -> Self {
Self {
mailbox_key: mailbox_key.into(),
cursor_store: PostgresSourceCursorStore::new(pool.clone()),
pool,
hooks,
query_executor_factory,
correlator,
address_hash_key,
challenge_suppress_after,
}
}
/// The attached `fraiseql_query` bridge factory, if any (test observability for
/// the #594 after:ingest wiring; the dispatch path reads the field directly).
#[cfg(test)]
#[must_use]
pub(crate) const fn query_executor_factory(
&self,
) -> Option<&crate::routes::after_mutation::QueryExecutorFactory> {
self.query_executor_factory.as_ref()
}
/// Correlate an inbound bounce / challenge / reply back to its send.
///
/// Best-effort: a correlation failure is logged, not propagated — the message
/// is already ingested, and a redelivery would be deduplicated by the spine, so
/// it would never re-correlate. A no-op when no correlator is wired.
async fn correlate(&self, message: &InboundMessage) {
let Some(correlator) = self.correlator.as_ref() else {
return;
};
let result = super::correlation::correlate(
correlator.as_ref(),
self.address_hash_key.as_deref(),
self.challenge_suppress_after,
chrono::Utc::now(),
message,
)
.await;
if let Err(error) = result {
warn!(
mailbox = %self.mailbox_key,
%error,
"delivery correlation failed; send-status left unchanged"
);
}
}
/// Fire the `after:ingest:email` functions for a persisted message.
fn dispatch(&self, message: &InboundMessage) {
let Some(ref hooks) = self.hooks else {
return;
};
let plans = plan_after_ingest_dispatch(hooks, message);
if !plans.is_empty() {
// #594: after:ingest functions write back under their `run_as` ceiling via
// the request-path executor factory threaded onto the sink at mount time.
spawn_after_ingest(hooks, plans, self.query_executor_factory.clone());
}
}
}
impl IngestSink for EmailIngestSink {
async fn ingest(
&self,
source_name: &str,
batch: PullBatch,
from: &CursorSnapshot,
) -> Result<bool> {
let mut tx = self
.pool
.begin()
.await
.map_err(|error| FraiseQLError::database(format!("email ingest begin: {error}")))?;
// Emit every message onto the spine; collect the genuinely-new ones for
// post-commit correlation + dispatch (a redelivery is deduped and skipped).
let mut fresh: Vec<InboundMessage> = Vec::new();
for message in &batch.messages {
if emit_in_tx(&mut tx, message).await?.is_new() {
fresh.push(message.clone());
}
}
// Advance the cursor in the same transaction (atomic with the emits).
let advanced = self
.cursor_store
.advance_in_tx(&mut tx, source_name, from, batch.next_cursor)
.await
.map_err(|error| FraiseQLError::database(format!("email cursor advance: {error}")))?;
if !advanced {
// Another replica moved the cursor on across a lease-boundary race:
// roll back rather than double-ingest.
if let Err(error) = tx.rollback().await {
warn!(mailbox = %self.mailbox_key, %error, "email ingest cursor-race rollback failed");
}
return Ok(false);
}
tx.commit()
.await
.map_err(|error| FraiseQLError::database(format!("email ingest commit: {error}")))?;
// Post-commit, per genuinely-new message: correlate first (not idempotent),
// then dispatch the app `after:ingest:email` functions.
for message in &fresh {
self.correlate(message).await;
self.dispatch(message);
}
Ok(true)
}
}