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
//! Serve-time functions-runtime preparation.
//!
//! Loads the compiled schema's function modules, registers the runtimes, attaches
//! the `send_email` wiring, and stores the before-mutation hooks so
//! `build_app_state` mounts them and after:mutation functions actually fire. Runs
//! once at serve time (async, fail-loud) — the counterpart to the RBAC/inbound
//! schema init already in the serve path.
use std::sync::Arc;
use super::{Server, ServerError};
use crate::subsystems::loader::build_functions_subsystem;
impl Server {
/// Prepare functions-runtime dispatch from the functions section this server
/// was built with.
///
/// When it declares functions, builds the subsystem (modules loaded from
/// `module_dir`, runtimes registered), attaches the `send_email` wiring, and
/// stores the resulting hooks. A no-op (hooks stay `None`) when no functions are
/// declared.
///
/// Reads [`Server::functions_config`], **not** `config.schema_path` (#896). The
/// disk re-read could configure the functions subsystem from a different artifact
/// than the one serving queries, and it made this step impossible on a server
/// built from an in-memory schema — which is why it used to run on one of the two
/// serving entry points.
///
/// # Errors
///
/// Returns [`ServerError::ConfigError`] if a declared function's module is
/// missing/unreadable (fail-loud: a declared function that can never run is a
/// misconfiguration).
pub(super) async fn prepare_functions_runtime(&mut self) -> Result<(), ServerError> {
let Some(functions_config) = self.functions_config.take() else {
return Ok(()); // no `functions` section was supplied
};
if functions_config.definitions.is_empty() {
return Ok(());
}
// Resolve the DLQ store choice (#598) before the config is consumed —
// `FRAISEQL_FUNCTIONS_DLQ_STORE` overrides the compiled `[functions] dlq_store`.
let dlq_store_kind = crate::routes::after_mutation::DlqStoreKind::resolve(
functions_config.dlq_store.as_deref(),
|key| std::env::var(key).ok(),
);
let subsystem = build_functions_subsystem(functions_config).map_err(|error| {
ServerError::ConfigError(format!("functions-runtime setup failed: {error}"))
})?;
// Sign the per-dispatch idempotency token when an HMAC secret is
// configured; unsigned digest otherwise (zero-config default).
let idempotency_key = self.build_idempotency_key();
if idempotency_key.is_some() {
tracing::info!("idempotency tokens are HMAC-signed (VERP-ready send-ids)");
}
let mut hooks =
subsystem.into_before_mutation_hooks().with_idempotency_key(idempotency_key);
// Select the durable Postgres-backed DLQ when requested and a pool exists;
// otherwise the in-memory store (the default) stays. A postgres request
// without a pool is a loud fallback, not a startup failure.
hooks = self.select_function_dlq(hooks, dlq_store_kind).await?;
if let Some((resolver, transport)) = self.build_send_email_wiring().await? {
hooks = hooks.with_email(resolver, transport);
tracing::info!(
"send_email host op enabled (host-owned from, per-connected-account SMTP)"
);
}
let function_count = hooks.module_registry.len();
let hooks = Arc::new(hooks);
self.install_function_seams(Arc::clone(&hooks));
self.functions_hooks = Some(hooks);
tracing::info!(functions = function_count, "functions-runtime dispatch enabled");
Ok(())
}
/// Install both function seams on the executor: the `before:mutation`
/// enforcement gate (#1327) and the function-backed query resolver (#1329).
///
/// The chain is enforcement, so it has to run wherever a mutation runs — which
/// is the engine's write chokepoint, not one HTTP handler. The gate therefore
/// lives on the executor's `RuntimeConfig`, and the executor is rebuilt here
/// through `Executor::rebuild_with` rather than a fourth construction path, so
/// it keeps the backend and the relay dispatch it already had (#750). Nothing
/// has been served at this point in the serve
/// path, so the rebuild discards no warm state; `with_compiled_schema` carries
/// caller-owned config through, so the gate also survives every later hot
/// reload.
fn install_function_seams(&mut self, hooks: Arc<crate::subsystems::BeforeMutationHooks>) {
use crate::routes::{
before_mutation::{BeforeMutationBudget, FunctionChainGate},
query_function::{FunctionQueryResolver, QueryFunctionBudget},
};
let budget = BeforeMutationBudget::from_env();
let gate = Arc::new(FunctionChainGate::new(Arc::clone(&hooks)).with_budget(budget));
// #1329: the read-side seam, installed in the **same** executor rebuild as
// the write-side gate. Two rebuilds would discard the first one's config the
// way `with_compiled_schema` carries caller-owned config forward only once.
let query_budget = QueryFunctionBudget::from_env();
let resolver = Arc::new(FunctionQueryResolver::new(hooks).with_budget(query_budget));
let config = self
.executor
.config()
.clone()
.with_before_mutation_gate(gate)
.with_query_function_resolver(resolver);
let schema = self.executor.schema().clone();
self.executor = Arc::new(self.executor.rebuild_with(schema, config));
if query_budget.is_enforced() {
tracing::info!(
budget_ms = u64::try_from(query_budget.duration().as_millis()).unwrap_or(u64::MAX),
"function-backed root query fields enabled (#1329), read-only and as the caller"
);
} else {
tracing::warn!(
env = QueryFunctionBudget::ENV,
"function-backed root query fields run with NO default latency ceiling — a \
function declaring no timeout_ms holds its request until the query timeout"
);
}
if budget.is_enforced() {
tracing::info!(
budget_ms = u64::try_from(budget.duration().as_millis()).unwrap_or(u64::MAX),
"before:mutation enforcement installed at the mutation chokepoint (every \
transport), with a read-only caller-scoped query bridge"
);
} else {
tracing::warn!(
env = BeforeMutationBudget::ENV,
"before:mutation chains run with NO latency ceiling — a slow or hanging hook \
holds its write's request open indefinitely"
);
}
}
/// Swap in the Postgres-backed function DLQ when selected and a pool exists (#598).
///
/// `dlq_store = "postgres"` makes a dead-lettered dispatch survive a restart; the
/// `_fraiseql_function_dlq` table is created here (async, fail-loud on DDL error).
/// A postgres request with no database pool is a loud fallback to the in-memory
/// store, not a startup failure — the same posture as the send-tracking store.
/// `dlq_store = "memory"` (the default) returns the hooks unchanged.
async fn select_function_dlq(
&self,
hooks: crate::subsystems::BeforeMutationHooks,
kind: crate::routes::after_mutation::DlqStoreKind,
) -> Result<crate::subsystems::BeforeMutationHooks, ServerError> {
use crate::routes::after_mutation::{DispatchDefaults, DlqStoreKind};
if kind != DlqStoreKind::Postgres {
return Ok(hooks);
}
let Some(pool) = self.db_pool.as_ref() else {
tracing::warn!(
"[functions] dlq_store = \"postgres\" but no database pool is available — \
dead-lettered dispatches use the in-memory store and will not survive a restart"
);
return Ok(hooks);
};
// Same retention cap the in-memory store honors, so `dlq_store` only changes
// durability, not the drop-newest policy.
let max_size = DispatchDefaults::from_env().dlq_max_size;
let store = crate::observers::pg_function_dlq::PgFunctionDlq::new(pool.clone(), max_size);
store.init().await.map_err(|error| {
ServerError::ConfigError(format!("failed to initialize function DLQ schema: {error}"))
})?;
tracing::info!("function DLQ: postgres-backed (dead-letters survive restarts)");
Ok(hooks.with_dlq(Arc::new(store)))
}
/// Derive the idempotency-token HMAC subkey from the configured server HMAC
/// secret (`hmac_secret_env` names an env var). `None` → the token stays an
/// unsigned digest (the zero-config default); a signed token is required before
/// it is exposed externally as a VERP Return-Path (P04b). A configured-but-empty
/// secret is a misconfiguration surfaced loudly, not silently signed with "".
fn build_idempotency_key(&self) -> Option<Arc<[u8]>> {
let env_name = self.config.hmac_secret_env.as_deref()?;
if let Some(secret) = std::env::var(env_name).ok().filter(|secret| !secret.is_empty()) {
let subkey = fraiseql_observers::derive_idempotency_subkey(secret.as_bytes());
Some(Arc::from(subkey.as_slice()))
} else {
tracing::warn!(
env = env_name,
"hmac_secret_env is set but the environment variable is empty/unset — \
idempotency tokens stay unsigned and VERP send-correlation is disabled"
);
None
}
}
/// Build the `send_email` wiring — sender-identity resolver + SMTP transport —
/// from config. Returns `None` when no SMTP mailbox is configured, leaving
/// `send_email` fail-loud.
///
/// When a database pool is available, the transport is wired to the
/// delivery-feedback store (`PgSendTracker`, tables created here); when the
/// server HMAC secret is set, the recipient address-hash key is derived so the
/// suppression check is active. The store's tables are created here (async,
/// fail-loud) rather than in the IMAP-worker block, because a send-only mailbox
/// needs them without ever polling.
///
/// # Errors
///
/// Returns [`ServerError::ConfigError`] if the delivery-feedback schema cannot
/// be created.
// Reason: the body awaits only under `inbound-email`; without it, `Ok(None)`.
#[allow(clippy::unused_async)]
async fn build_send_email_wiring(
&self,
) -> Result<
Option<(
Arc<dyn fraiseql_functions::SenderIdentityResolver>,
Arc<dyn fraiseql_functions::EmailTransport>,
)>,
ServerError,
> {
#[cfg(feature = "inbound-email")]
{
// The delivery-feedback store (suppression + send-status + exactly-once)
// needs a database pool; without one the transport still sends, just
// without tracking.
let tracker = match self.db_pool.as_ref() {
Some(pool) => {
let tracker = crate::inbound::email::PgSendTracker::new(pool.clone());
tracker.init().await.map_err(|error| {
ServerError::ConfigError(format!(
"failed to initialize send-tracking schema: {error}"
))
})?;
Some(Arc::new(tracker) as Arc<dyn crate::inbound::email::SendTracker>)
},
None => None,
};
let address_hash_key = self.build_address_hash_key();
let Some(transport) = crate::inbound::email::build_email_transport(
&self.config.mailbox,
|name| std::env::var(name).ok(),
tracker,
address_hash_key,
) else {
return Ok(None);
};
Ok(Some((self.build_sender_resolver(), transport)))
}
#[cfg(not(feature = "inbound-email"))]
{
Ok(None)
}
}
/// Derive the recipient address-hash key from the configured server HMAC secret
/// (domain-separated from the send-id subkey). `None` → no suppression check
/// (the same fail-closed posture as the unsigned idempotency token).
///
/// Shared with the poll-worker correlation path (`lifecycle`), which keys the
/// same recipient hash to write/lift suppressions.
#[cfg(feature = "inbound-email")]
pub(super) fn build_address_hash_key(&self) -> Option<Arc<[u8]>> {
let env_name = self.config.hmac_secret_env.as_deref()?;
let secret = std::env::var(env_name).ok().filter(|secret| !secret.is_empty())?;
let key = fraiseql_observers::derive_address_hash_key(secret.as_bytes());
Some(Arc::from(key.as_slice()))
}
/// The sender-identity resolver: DB-backed on the shared identity primitive
/// when `[identity.sender]` is enabled, else the login-email default.
#[cfg(feature = "inbound-email")]
fn build_sender_resolver(&self) -> Arc<dyn fraiseql_functions::SenderIdentityResolver> {
#[cfg(feature = "auth")]
if let Some(sender) =
self.config.identity.as_ref().and_then(|identity| identity.sender.as_ref())
{
if sender.enabled {
if let Some(pool) = self.enrichment_pool.as_ref() {
let resolver = crate::identity::resolver::IdentityResolver::postgres(
sender.clone(),
pool.clone(),
);
return Arc::new(crate::identity::sender::DbSenderIdentityResolver::new(
resolver,
SENDING_ADDRESS_FIELD,
Some(DISPLAY_NAME_FIELD.to_string()),
));
}
tracing::warn!(
"[identity.sender] is enabled but no auth database pool is available — \
send_email falls back to the login-email sender"
);
}
}
Arc::new(fraiseql_functions::LoginEmailSender)
}
}
/// The resolved sender query's field holding the verified from-address (convention,
/// matching `docs/architecture/enriched-identity-rls.md`).
#[cfg(all(feature = "inbound-email", feature = "auth"))]
const SENDING_ADDRESS_FIELD: &str = "sending_address";
/// The resolved sender query's field holding the sender display name.
#[cfg(all(feature = "inbound-email", feature = "auth"))]
const DISPLAY_NAME_FIELD: &str = "display_name";