Skip to main content

fraiseql_server/server/
extensions.rs

1//! Server extensions: relay pagination, Arrow Flight service, and observer runtime
2//! initialization.
3
4use std::sync::Arc;
5
6#[cfg(feature = "arrow")]
7use fraiseql_arrow::FraiseQLFlightService;
8#[cfg(all(feature = "arrow", feature = "auth"))]
9use fraiseql_core::security::OidcValidator;
10use fraiseql_core::{
11    cache::{CacheConfig, CachedDatabaseAdapter, QueryResultCache},
12    db::traits::{DatabaseAdapter, RelayDatabaseAdapter},
13    runtime::{Executor, RuntimeConfig, SubscriptionManager},
14    schema::CompiledSchema,
15};
16#[cfg(feature = "observers")]
17use tokio::sync::RwLock;
18use tracing::info;
19#[cfg(feature = "observers")]
20use tracing::warn;
21
22#[cfg(feature = "arrow")]
23use super::RateLimiter;
24#[cfg(all(feature = "arrow", feature = "auth"))]
25use super::ServerError;
26#[cfg(feature = "observers")]
27use super::{ObserverRuntime, ObserverRuntimeConfig};
28use super::{Result, Server, ServerConfig};
29
30impl<A: DatabaseAdapter + RelayDatabaseAdapter + Clone + Send + Sync + 'static>
31    Server<CachedDatabaseAdapter<A>>
32{
33    /// Create a server with relay pagination support enabled.
34    ///
35    /// The adapter must implement [`RelayDatabaseAdapter`]. Currently, only
36    /// `PostgresAdapter` and `CachedDatabaseAdapter<PostgresAdapter>` satisfy this bound.
37    ///
38    /// Relay queries issued against a server created with [`Server::new`] return a
39    /// `Validation` error at runtime; those issued against a server created with this
40    /// constructor succeed.
41    ///
42    /// # Arguments
43    ///
44    /// * `config` - Server configuration
45    /// * `schema` - Compiled GraphQL schema
46    /// * `adapter` - Database adapter (must implement `RelayDatabaseAdapter`)
47    /// * `db_pool` - Database connection pool (optional, required for observers)
48    ///
49    /// # Errors
50    ///
51    /// Returns error if OIDC validator initialization fails.
52    ///
53    /// # Panics
54    ///
55    /// Panics if the `adapter` `Arc` has been cloned before calling this constructor
56    /// (refcount > 1). The builder must have exclusive ownership to unwrap the adapter
57    /// for `CachedDatabaseAdapter` construction.
58    ///
59    /// # Example
60    ///
61    /// ```text
62    /// // Requires: running PostgreSQL database and compiled schema file.
63    /// let adapter = Arc::new(PostgresAdapter::new(db_url).await?);
64    /// let server = Server::with_relay_pagination(config, schema, adapter, None).await?;
65    /// server.serve().await?;
66    /// ```
67    pub async fn with_relay_pagination(
68        config: ServerConfig,
69        schema: CompiledSchema,
70        adapter: Arc<A>,
71        db_pool: Option<sqlx::PgPool>,
72    ) -> Result<Self> {
73        // Validate cache + RLS safety (mirrors Server::new).
74        if config.cache_enabled && !schema.has_rls_configured() {
75            if schema.is_multi_tenant() {
76                return Err(super::ServerError::ConfigError(
77                    "Cache is enabled in a multi-tenant schema but no Row-Level Security \
78                     policies are declared. This would allow cross-tenant cache hits and \
79                     data leakage. In fraiseql.toml, either disable caching with \
80                     [cache] enabled = false, declare [security.rls] policies, or set \
81                     [security] multi_tenant = false to acknowledge single-tenant mode."
82                        .to_string(),
83                ));
84            }
85            tracing::warn!(
86                "Query-result caching is enabled but no Row-Level Security policies are \
87                 declared in the compiled schema. This is safe for single-tenant deployments."
88            );
89        }
90
91        // Same boot gates as `Server::new` — these must not drift by constructor (H16).
92        // Refuse to boot if any field is marked for at-rest encryption (H12); the write
93        // path does not encrypt, so the data would be stored in plaintext.
94        crate::server::initialization::field_encryption_unsupported_check(&schema)?;
95        // Build the runtime config from the compiled schema (validates format version,
96        // reads the audit flag, applies the #421 page-size ceiling + change-log toggle).
97        let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
98            super::ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
99        })?;
100
101        // Read security configs from compiled schema BEFORE schema is moved.
102        #[cfg(feature = "federation")]
103        let circuit_breaker = schema.federation.as_ref().and_then(
104            crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
105        );
106        #[cfg(not(feature = "federation"))]
107        let circuit_breaker: Option<()> = None;
108        #[cfg(not(feature = "federation"))]
109        let _ = &schema.federation;
110        let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
111        #[cfg(feature = "auth")]
112        let state_encryption = Self::state_encryption_from_schema(&schema)?;
113        #[cfg(not(feature = "auth"))]
114        let state_encryption: Option<
115            std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
116        > = None;
117        #[cfg(feature = "auth")]
118        let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
119        #[cfg(not(feature = "auth"))]
120        let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
121        #[cfg(feature = "auth")]
122        let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
123        #[cfg(not(feature = "auth"))]
124        let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
125        let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
126        let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
127        let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
128        let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
129        let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
130
131        let cache_config = CacheConfig::from(config.cache_enabled);
132        let cache = QueryResultCache::new(cache_config);
133        // Unwrap Arc: refcount is 1 here — adapter has not been cloned since being passed in.
134        let inner = Arc::into_inner(adapter)
135            .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
136        let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
137            .with_ttl_overrides_from_schema(&schema);
138        let executor = Arc::new(Executor::with_config_and_relay(
139            schema.clone(),
140            Arc::new(cached),
141            executor_config,
142        ));
143        let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
144
145        let mut server = Self::from_executor(
146            config,
147            executor,
148            subscription_manager,
149            circuit_breaker,
150            error_sanitizer,
151            state_encryption,
152            pkce_store,
153            oidc_server_client,
154            schema_rate_limiter,
155            api_key_authenticator,
156            revocation_manager,
157            trusted_docs,
158            db_pool,
159            tasks,
160        )
161        .await?;
162
163        server.adapter_cache_enabled = cache_config.enabled;
164
165        // Initialize MCP config from compiled schema when the feature is compiled in.
166        #[cfg(feature = "mcp")]
167        if let Some(ref cfg) = server.executor.schema().mcp_config {
168            if cfg.enabled {
169                let tool_count =
170                    crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
171                info!(
172                    path = %cfg.path,
173                    transport = %cfg.transport,
174                    tools = tool_count,
175                    "MCP server configured"
176                );
177                server.mcp_config = Some(cfg.clone());
178            }
179        }
180
181        // Initialize APQ store when enabled.
182        if server.config.apq_enabled {
183            let apq_store: fraiseql_core::apq::ArcApqStorage =
184                Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
185            server.apq_store = Some(apq_store);
186            info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
187        }
188
189        Ok(server)
190    }
191}
192
193impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
194    /// Create new server with pre-configured Arrow Flight service.
195    ///
196    /// Use this constructor when you want to provide a Flight service with a real database adapter.
197    ///
198    /// # Arguments
199    ///
200    /// * `config` - Server configuration
201    /// * `schema` - Compiled GraphQL schema
202    /// * `adapter` - Database adapter
203    /// * `db_pool` - Database connection pool (optional, required for observers)
204    /// * `flight_service` - Pre-configured Flight service (only available with arrow feature)
205    ///
206    /// # Errors
207    ///
208    /// Returns error if OIDC validator initialization fails.
209    #[cfg(feature = "arrow")]
210    pub async fn with_flight_service(
211        config: ServerConfig,
212        schema: CompiledSchema,
213        adapter: Arc<A>,
214        #[allow(unused_variables)]
215        // Reason: used inside #[cfg(feature = "observers")] block; unused when feature is off
216        db_pool: Option<sqlx::PgPool>,
217        flight_service: Option<FraiseQLFlightService>,
218    ) -> Result<Self> {
219        // Same boot gates as `Server::new` — these must not drift by constructor (H16).
220        // Refuse to boot on at-rest-encryption-marked fields (H12, plaintext write path).
221        crate::server::initialization::field_encryption_unsupported_check(&schema)?;
222        // Build the runtime config from the compiled schema (validates format version,
223        // reads the audit flag, applies the #421 page-size ceiling + change-log toggle).
224        let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
225            super::ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
226        })?;
227
228        // Read security configs from compiled schema BEFORE schema is moved.
229        #[cfg(feature = "federation")]
230        let circuit_breaker = schema.federation.as_ref().and_then(
231            crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
232        );
233        // Non-federation builds construct the `Server` struct literal below with the
234        // `circuit_breaker` field cfg'd out, so there is no placeholder to bind here —
235        // just mark `schema.federation` read to mirror the federation branch.
236        #[cfg(not(feature = "federation"))]
237        let _ = &schema.federation;
238        let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
239        #[cfg(feature = "auth")]
240        let state_encryption = Self::state_encryption_from_schema(&schema)?;
241        #[cfg(not(feature = "auth"))]
242        let _state_encryption: Option<
243            std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
244        > = None;
245        #[cfg(feature = "auth")]
246        let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
247        #[cfg(not(feature = "auth"))]
248        let _pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
249        #[cfg(feature = "auth")]
250        let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
251        #[cfg(not(feature = "auth"))]
252        let _oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
253        let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
254        let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
255        let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
256        let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
257        let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
258
259        let executor = Arc::new(Executor::with_config(schema.clone(), adapter, executor_config));
260        let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
261
262        // Initialize OIDC validator if auth is configured
263        #[cfg(feature = "auth")]
264        let oidc_validator = if let Some(ref auth_config) = config.auth {
265            info!(
266                issuer = %auth_config.issuer,
267                "Initializing OIDC authentication"
268            );
269            let validator = OidcValidator::new(auth_config.clone())
270                .await
271                .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
272            Some(Arc::new(validator))
273        } else {
274            None
275        };
276        #[cfg(not(feature = "auth"))]
277        let oidc_validator: Option<Arc<fraiseql_core::security::OidcValidator>> = None;
278
279        // Initialize HS256 validator if configured (mutually exclusive with OIDC).
280        let hs256_auth = super::builder::build_hs256_auth(&config)?;
281
282        // Initialize rate limiter: compiled schema config takes priority over server config.
283        let rate_limiter = if let Some(rl) = schema_rate_limiter {
284            Some(rl)
285        } else if let Some(ref rate_config) = config.rate_limiting {
286            if rate_config.enabled {
287                info!(
288                    rps_per_ip = rate_config.rps_per_ip,
289                    rps_per_user = rate_config.rps_per_user,
290                    "Initializing rate limiting from server config"
291                );
292                Some(Arc::new(RateLimiter::new(rate_config.clone())))
293            } else {
294                info!("Rate limiting disabled by configuration");
295                None
296            }
297        } else {
298            None
299        };
300
301        // Initialize observer runtime
302        #[cfg(feature = "observers")]
303        let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
304
305        // Warn if PKCE is configured but [auth] is missing.
306        #[cfg(feature = "auth")]
307        if pkce_store.is_some() && oidc_server_client.is_none() {
308            tracing::error!(
309                "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
310                 Auth routes will NOT be mounted."
311            );
312        }
313
314        // Refuse to start if FRAISEQL_REQUIRE_REDIS is set and PKCE store is in-memory.
315        #[cfg(feature = "auth")]
316        Self::check_redis_requirement(pkce_store.as_ref())?;
317
318        // Spawn background PKCE state cleanup task (every 5 minutes).
319        #[cfg(feature = "auth")]
320        Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
321
322        let apq_enabled = config.apq_enabled;
323
324        Ok(Self {
325            config,
326            executor,
327            subscription_manager,
328            subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
329            max_subscriptions_per_connection: None,
330            oidc_validator,
331            hs256_auth,
332            rate_limiter,
333            #[cfg(feature = "secrets")]
334            secrets_manager: None,
335            #[cfg(feature = "federation")]
336            circuit_breaker,
337            error_sanitizer,
338            #[cfg(feature = "auth")]
339            state_encryption,
340            #[cfg(feature = "auth")]
341            pkce_store,
342            #[cfg(feature = "auth")]
343            oidc_server_client,
344            #[cfg(feature = "auth")]
345            social_login: None,
346            #[cfg(feature = "auth")]
347            anon_signup_state: None,
348            #[cfg(feature = "auth")]
349            mfa_state: None,
350            api_key_authenticator,
351            revocation_manager,
352            apq_store: if apq_enabled {
353                Some(Arc::new(fraiseql_core::apq::InMemoryApqStorage::default())
354                    as fraiseql_core::apq::ArcApqStorage)
355            } else {
356                None
357            },
358            trusted_docs,
359            #[cfg(feature = "mcp")]
360            mcp_config: None,
361            pool_tuning_config: None,
362            #[cfg(feature = "observers")]
363            observer_runtime,
364            #[cfg(feature = "auth")]
365            enrichment_pool: db_pool.clone(),
366            #[cfg(feature = "observers")]
367            db_pool,
368            storage_state: None,
369            realtime_state: None,
370            #[cfg(feature = "functions-runtime")]
371            functions_hooks: None,
372            tenant_executor_factory: None,
373            #[cfg(feature = "arrow")]
374            flight_service,
375            adapter_cache_enabled: false,
376            broadcast_manager: None,
377            presence_manager: None,
378            storage_backend: None,
379            storage_max_upload_bytes: 100 * 1024 * 1024, // 100 MiB default
380            #[cfg(feature = "functions")]
381            function_store: None,
382            #[cfg(feature = "functions")]
383            function_runtime: None,
384            usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
385            tasks,
386        })
387    }
388
389    /// Initialize observer runtime from configuration.
390    ///
391    /// # Errors
392    ///
393    /// Returns `ServerError::ConfigError` when a non-Postgres observer transport
394    /// is selected that this binary cannot run (feature not compiled in, or NATS
395    /// without a URL) while in production mode (#350), or when the transport
396    /// configuration is otherwise invalid. In development such a selection is
397    /// downgraded to a warning and the runtime falls back to PostgreSQL.
398    #[cfg(feature = "observers")]
399    pub(super) async fn init_observer_runtime(
400        config: &ServerConfig,
401        pool: Option<&sqlx::PgPool>,
402    ) -> crate::Result<Option<Arc<RwLock<ObserverRuntime>>>> {
403        use fraiseql_observers::config::TransportKind;
404
405        // Check if enabled
406        let observer_config = match &config.observers {
407            Some(cfg) if cfg.enabled => cfg,
408            _ => {
409                info!("Observer runtime disabled");
410                return Ok(None);
411            },
412        };
413
414        let Some(pool) = pool else {
415            warn!("No database pool provided for observers");
416            return Ok(None);
417        };
418
419        info!("Initializing observer runtime");
420
421        // Resolve the event transport from compiled config + env overrides, then
422        // fail loud (#350) on a selection this binary cannot run before validating
423        // the finer NATS/JetStream bounds — never a silent fallback to PostgreSQL.
424        let mut transport = observer_config.runtime.transport.clone().with_env_overrides();
425        let compiled_in = cfg!(feature = "observers-nats");
426        let nats_url_present = !transport.nats.url.is_empty();
427        crate::server::initialization::observer_transport_check(
428            transport.transport,
429            compiled_in,
430            nats_url_present,
431            crate::ServerConfig::is_production_mode(),
432        )?;
433
434        // In production an unrunnable selection already returned above; the only
435        // way past the guard with an unrunnable transport is development, where it
436        // was downgraded to a warning — fall back to PostgreSQL so local boot works.
437        let usable = match transport.transport {
438            TransportKind::Postgres | TransportKind::InMemory => true,
439            TransportKind::Nats => compiled_in && nats_url_present,
440            _ => false,
441        };
442        if !usable {
443            transport.transport = TransportKind::Postgres;
444        }
445
446        transport.validate().map_err(|e| {
447            crate::ServerError::ConfigError(format!("invalid observer transport config: {e}"))
448        })?;
449
450        let runtime_config = ObserverRuntimeConfig::new(pool.clone())
451            .with_poll_interval(observer_config.runtime.poll_interval_ms)
452            .with_batch_size(observer_config.runtime.batch_size)
453            .with_channel_capacity(observer_config.runtime.channel_capacity)
454            .with_max_dlq_size(observer_config.runtime.max_dlq_size)
455            .with_transport(transport)
456            .with_email(observer_config.runtime.email.clone())
457            .with_log_payloads(observer_config.runtime.log_payloads);
458
459        let runtime = ObserverRuntime::new(runtime_config);
460        Ok(Some(Arc::new(RwLock::new(runtime))))
461    }
462}