Skip to main content

fraiseql_server/subsystems/
mod.rs

1//! Server subsystem assembly and lifecycle management.
2//!
3//! [`ServerSubsystems`] bundles the three optional platform extensions —
4//! object storage, serverless functions, and realtime entity streams — into a
5//! single coherent struct that the server can query, route-mount, and shut down
6//! in a controlled order.
7//!
8//! # Assembly
9//!
10//! Use [`builder::ServerSubsystemsBuilder`] to assemble subsystems from their
11//! pre-built parts. The builder validates cross-subsystem dependencies before
12//! returning the final [`ServerSubsystems`]:
13//!
14//! ```rust,ignore
15//! let subsystems = ServerSubsystemsBuilder::new()
16//!     .with_storage(storage_subsystem)
17//!     .with_functions(functions_subsystem)
18//!     .with_realtime(realtime_subsystem)
19//!     .build()?;
20//! ```
21//!
22//! # Shutdown order
23//!
24//! Shutdown proceeds in reverse initialization order:
25//! 1. Stop the cron scheduler (functions)
26//! 2. Drain the realtime event channel
27//! 3. Drop the realtime server (closes all `WebSocket` connections)
28//! 4. Drop the functions observer (stops dispatching events)
29//! 5. Drop the storage backend (flushes any pending writes)
30
31pub mod builder;
32pub mod validator;
33
34/// Loads function modules from disk and assembles the functions-runtime subsystem.
35#[cfg(feature = "functions-runtime")]
36pub mod loader;
37
38#[cfg(test)]
39mod tests;
40
41use std::sync::Arc;
42
43pub use builder::{ServerSubsystemsBuilder, SubsystemBuildError};
44use fraiseql_functions::{FunctionObserver, triggers::TriggerRegistry};
45pub use validator::{SubsystemConfigWarning, validate_subsystems_config};
46
47use crate::{
48    realtime::{
49        observer::RealtimeBroadcastObserver, routes::RealtimeSchemaConfig, server::RealtimeServer,
50    },
51    schema::loader::{FunctionsConfig, SchemaStorageConfig},
52};
53
54// ── Subsystem structs ─────────────────────────────────────────────────────────
55
56/// Storage subsystem: backend, metadata repository, RLS evaluator, and bucket config.
57///
58/// Assembled at server startup from the `[storage]` section of the compiled schema
59/// and the server's `PgPool`. Use the [`fraiseql_storage::storage_router`] with the
60/// contained `state` to mount the `/storage/v1` route tree.
61pub struct StorageSubsystem {
62    /// Runtime storage state (backend + metadata repo + RLS + bucket config map).
63    ///
64    /// Pass this to [`fraiseql_storage::storage_router`] to mount the HTTP routes.
65    pub state: fraiseql_storage::StorageState,
66
67    /// Schema-level bucket definitions from the compiled schema.
68    pub schema_config: SchemaStorageConfig,
69}
70
71/// Functions subsystem: observer and trigger registry.
72///
73/// Assembled at server startup from the `[functions]` section of the compiled schema.
74/// The observer dispatches events to function runtimes; the registry maps triggers to
75/// function definitions and provides HTTP route matchers.
76pub struct FunctionsSubsystem {
77    /// Observer that dispatches trigger events to the appropriate function runtime.
78    pub observer: Arc<FunctionObserver>,
79
80    /// Registry mapping trigger types to function definitions.
81    pub trigger_registry: TriggerRegistry,
82
83    /// Loaded function modules keyed by function name.
84    ///
85    /// Populated at server startup by reading source files from `config.module_dir`.
86    /// Used by the before-mutation chain and the after-mutation dispatcher.
87    pub module_registry: std::collections::HashMap<String, fraiseql_functions::FunctionModule>,
88
89    /// Schema-level functions configuration (definitions + module directory).
90    pub config: FunctionsConfig,
91}
92
93/// Realtime subsystem: `WebSocket` broadcast server and event observer.
94///
95/// Assembled at server startup from the `[realtime]` section of the compiled schema.
96/// The server handles `WebSocket` connections; the observer receives mutation events
97/// from the observer pipeline and forwards them to connected clients.
98pub struct RealtimeSubsystem {
99    /// The `WebSocket` broadcast server.
100    ///
101    /// Pass this to [`crate::realtime::routes::realtime_router`] to mount `/realtime/v1`.
102    pub server: Arc<RealtimeServer>,
103
104    /// Observer that forwards mutation events into the realtime delivery pipeline.
105    pub observer: RealtimeBroadcastObserver,
106
107    /// Schema-level realtime configuration (enabled flag, entity list, capacity overrides).
108    pub schema_config: RealtimeSchemaConfig,
109}
110
111// ── Aggregated container ──────────────────────────────────────────────────────
112
113/// All optional platform subsystems assembled from the compiled schema.
114///
115/// Each field is `None` when the corresponding section is absent from or disabled
116/// in the compiled schema. Callers can use [`is_storage_enabled`][Self::is_storage_enabled]
117/// etc. to check at a glance, or match directly on the `Option` fields.
118#[allow(missing_debug_implementations)] // Reason: inner types (RealtimeBroadcastObserver) don't implement Debug
119pub struct ServerSubsystems {
120    /// Object storage subsystem, present when the schema's `"storage"` key is set.
121    pub storage: Option<StorageSubsystem>,
122
123    /// Serverless functions subsystem, present when the schema's `"functions"` key is set.
124    pub functions: Option<FunctionsSubsystem>,
125
126    /// Realtime broadcast subsystem, present when the schema's `"realtime"` key is set
127    /// and `enabled` is `true`.
128    pub realtime: Option<RealtimeSubsystem>,
129}
130
131impl ServerSubsystems {
132    /// Create an empty `ServerSubsystems` with all subsystems disabled.
133    ///
134    /// Equivalent to `ServerSubsystemsBuilder::new().build().unwrap()`.
135    #[must_use]
136    pub const fn none() -> Self {
137        Self {
138            storage:   None,
139            functions: None,
140            realtime:  None,
141        }
142    }
143
144    /// Returns `true` if the storage subsystem is present.
145    #[must_use]
146    pub const fn is_storage_enabled(&self) -> bool {
147        self.storage.is_some()
148    }
149
150    /// Returns `true` if the functions subsystem is present.
151    #[must_use]
152    pub const fn is_functions_enabled(&self) -> bool {
153        self.functions.is_some()
154    }
155
156    /// Returns `true` if the realtime subsystem is present.
157    #[must_use]
158    pub const fn is_realtime_enabled(&self) -> bool {
159        self.realtime.is_some()
160    }
161}
162
163// ── Before-mutation hook bundle ───────────────────────────────────────────────
164
165/// Shared bundle of before-mutation state passed into `AppState` for handler access.
166///
167/// This is a lightweight, cloneable snapshot of the parts of [`FunctionsSubsystem`]
168/// that are needed on the hot path for before-mutation checks. It is extracted once
169/// at server startup and stored in `AppState` via an `Arc`.
170pub struct BeforeMutationHooks {
171    /// Registry of all loaded triggers, keyed by trigger type and mutation name.
172    pub trigger_registry: TriggerRegistry,
173    /// Loaded function modules keyed by function name.
174    pub module_registry:  std::collections::HashMap<String, fraiseql_functions::FunctionModule>,
175    /// Observer that dispatches events to the appropriate function runtime.
176    pub observer:         std::sync::Arc<FunctionObserver>,
177
178    /// Dead-letter queue for durable after:mutation dispatch: an invocation that
179    /// exhausts its retries (or fails permanently) is pushed here rather than
180    /// silently lost.
181    #[cfg(feature = "functions-runtime")]
182    pub dlq: std::sync::Arc<dyn fraiseql_observers::DeadLetterQueue>,
183
184    /// Per-function dispatch settings (re-runnable flag + retry policy) resolved
185    /// from the compiled schema, keyed by function name. Functions absent from
186    /// the map fall back to the durable `FunctionDispatchSetting::default`.
187    #[cfg(feature = "functions-runtime")]
188    pub dispatch_settings:
189        std::collections::HashMap<String, crate::routes::after_mutation::FunctionDispatchSetting>,
190
191    /// Host-owned sender-identity resolver for the `send_email` op — resolves the
192    /// `from` from the authenticated context (the #539 seam). `None` → `send_email`
193    /// is unconfigured and fails loud. Set together with `email_transport` via
194    /// [`with_email`](Self::with_email).
195    #[cfg(feature = "functions-runtime")]
196    pub sender_resolver: Option<std::sync::Arc<dyn fraiseql_functions::SenderIdentityResolver>>,
197
198    /// Email transport for the `send_email` op. `None` → `send_email` fails loud.
199    #[cfg(feature = "functions-runtime")]
200    pub email_transport: Option<std::sync::Arc<dyn fraiseql_functions::EmailTransport>>,
201
202    /// HMAC subkey for the per-dispatch idempotency token, derived from the server
203    /// HMAC secret. `Some` → the token is signed (unforgeable, required before it is
204    /// exposed in a VERP Return-Path); `None` → an unsigned digest (zero-config
205    /// default). Set via [`with_idempotency_key`](Self::with_idempotency_key).
206    #[cfg(feature = "functions-runtime")]
207    pub idempotency_key: Option<std::sync::Arc<[u8]>>,
208
209    /// Per-function `run_as` authority ceilings (#594), keyed by function name.
210    /// A function absent from the map has no ceiling ⇒ its `fraiseql_query` bridge
211    /// runs fail-closed (anonymous `system_job`; RLS/field-authz deny writes).
212    /// Populated from the compiled schema's function definitions.
213    #[cfg(feature = "functions-runtime")]
214    pub run_as: std::collections::HashMap<String, fraiseql_functions::RunAs>,
215}
216
217impl BeforeMutationHooks {
218    /// Create a hook bundle with default durable-dispatch wiring: an unbounded
219    /// in-memory dead-letter queue and no per-function overrides (every function
220    /// uses the durable default).
221    ///
222    /// For the full compiled-schema resolution (per-function settings +
223    /// `FRAISEQL_FUNCTIONS_*` env overrides), use
224    /// [`FunctionsSubsystem::into_before_mutation_hooks`] instead.
225    #[must_use]
226    pub fn new(
227        trigger_registry: TriggerRegistry,
228        module_registry: std::collections::HashMap<String, fraiseql_functions::FunctionModule>,
229        observer: Arc<FunctionObserver>,
230    ) -> Self {
231        Self {
232            trigger_registry,
233            module_registry,
234            observer,
235            #[cfg(feature = "functions-runtime")]
236            dlq: Arc::new(crate::observers::runtime::InMemoryDlq::new_with_max(None)),
237            #[cfg(feature = "functions-runtime")]
238            dispatch_settings: std::collections::HashMap::new(),
239            #[cfg(feature = "functions-runtime")]
240            sender_resolver: None,
241            #[cfg(feature = "functions-runtime")]
242            email_transport: None,
243            #[cfg(feature = "functions-runtime")]
244            idempotency_key: None,
245            #[cfg(feature = "functions-runtime")]
246            run_as: std::collections::HashMap::new(),
247        }
248    }
249
250    /// Attach the HMAC subkey that signs the per-dispatch idempotency token.
251    ///
252    /// Derived once from the server HMAC secret
253    /// ([`fraiseql_observers::derive_idempotency_subkey`]). `None` leaves the token
254    /// as an unsigned digest — the zero-config default; a signed token is required
255    /// before it is exposed externally as a VERP Return-Path (P04b).
256    #[cfg(feature = "functions-runtime")]
257    #[must_use]
258    pub fn with_idempotency_key(mut self, key: Option<std::sync::Arc<[u8]>>) -> Self {
259        self.idempotency_key = key;
260        self
261    }
262
263    /// Enable the `send_email` host op for dispatched functions by attaching a
264    /// sender-identity resolver (the host-owned `from`) and an email transport.
265    ///
266    /// Without both, `send_email` fails loud (mirrors the `sql_query`
267    /// fail-loud-until-wired stance). The resolver is the #539 seam —
268    /// `LoginEmailSender` (from = login email) by default, a DB-backed resolver
269    /// where the sending mailbox differs; the transport is the per-connected-account
270    /// SMTP relay ([`SmtpMailboxTransport`](crate::inbound::email::SmtpMailboxTransport)).
271    #[cfg(feature = "functions-runtime")]
272    #[must_use]
273    pub fn with_email(
274        mut self,
275        sender_resolver: Arc<dyn fraiseql_functions::SenderIdentityResolver>,
276        email_transport: Arc<dyn fraiseql_functions::EmailTransport>,
277    ) -> Self {
278        self.sender_resolver = Some(sender_resolver);
279        self.email_transport = Some(email_transport);
280        self
281    }
282
283    /// Replace the dead-letter store (#598).
284    ///
285    /// The default from [`FunctionsSubsystem::into_before_mutation_hooks`] is the
286    /// in-memory store; the serve path swaps in the Postgres-backed
287    /// [`PgFunctionDlq`](crate::observers::pg_function_dlq::PgFunctionDlq) when
288    /// `[functions] dlq_store = "postgres"` and a database pool is available, so a
289    /// dead-lettered dispatch survives a restart.
290    #[cfg(feature = "functions-runtime")]
291    #[must_use]
292    pub fn with_dlq(mut self, dlq: Arc<dyn fraiseql_observers::DeadLetterQueue>) -> Self {
293        self.dlq = dlq;
294        self
295    }
296}
297
298#[cfg(feature = "functions-runtime")]
299impl FunctionsSubsystem {
300    /// Assemble the before-mutation hook bundle for `AppState`.
301    ///
302    /// Resolves each function's durable-dispatch settings (re-runnable flag +
303    /// retry policy) from the compiled schema, layering the
304    /// `FRAISEQL_FUNCTIONS_*` environment overrides via `DispatchDefaults::from_env`,
305    /// and creates the shared dead-letter queue (capped by
306    /// `FRAISEQL_FUNCTIONS_DLQ_MAX_SIZE`). Consumes the subsystem because the hook
307    /// bundle takes ownership of its trigger registry, modules, and observer.
308    #[must_use]
309    pub fn into_before_mutation_hooks(self) -> BeforeMutationHooks {
310        use crate::routes::after_mutation::{DispatchDefaults, resolve_dispatch_settings};
311
312        let defaults = DispatchDefaults::from_env();
313        let dispatch_settings = resolve_dispatch_settings(&self.config.definitions, &defaults);
314        let dlq = std::sync::Arc::new(crate::observers::runtime::InMemoryDlq::new_with_max(
315            defaults.dlq_max_size,
316        ));
317
318        // #594: collect each function's `run_as` ceiling, keyed by name. A function
319        // with no `run_as` is simply absent (fail-closed at dispatch time).
320        let run_as = self
321            .config
322            .definitions
323            .iter()
324            .filter_map(|def| def.run_as.clone().map(|ceiling| (def.name.clone(), ceiling)))
325            .collect();
326
327        BeforeMutationHooks {
328            trigger_registry: self.trigger_registry,
329            module_registry: self.module_registry,
330            observer: self.observer,
331            dlq,
332            dispatch_settings,
333            sender_resolver: None,
334            email_transport: None,
335            idempotency_key: None,
336            run_as,
337        }
338    }
339}