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 service_account_authenticator =
128            crate::service_account::service_account_authenticator_from_schema(&schema);
129        let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
130        let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
131        let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
132
133        let cache_config = CacheConfig::from(config.cache_enabled);
134        let cache = QueryResultCache::new(cache_config);
135        // Unwrap Arc: refcount is 1 here — adapter has not been cloned since being passed in.
136        let inner = Arc::into_inner(adapter)
137            .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
138        let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
139            .with_ttl_overrides_from_schema(&schema);
140        let executor = Arc::new(Executor::with_config_and_relay(
141            schema.clone(),
142            Arc::new(cached),
143            executor_config,
144        ));
145        let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
146
147        let mut server = Self::from_executor(
148            config,
149            executor,
150            subscription_manager,
151            circuit_breaker,
152            error_sanitizer,
153            state_encryption,
154            pkce_store,
155            oidc_server_client,
156            schema_rate_limiter,
157            api_key_authenticator,
158            service_account_authenticator,
159            revocation_manager,
160            trusted_docs,
161            db_pool,
162            tasks,
163        )
164        .await?;
165
166        server.adapter_cache_enabled = cache_config.enabled;
167
168        // Initialize MCP config from compiled schema when the feature is compiled in.
169        #[cfg(feature = "mcp")]
170        if let Some(ref cfg) = server.executor.schema().mcp_config {
171            if cfg.enabled {
172                let tool_count =
173                    crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
174                info!(
175                    path = %cfg.path,
176                    transport = %cfg.transport,
177                    tools = tool_count,
178                    "MCP server configured"
179                );
180                server.mcp_config = Some(cfg.clone());
181            }
182        }
183
184        // Initialize APQ store when enabled.
185        if server.config.apq_enabled {
186            let apq_store: fraiseql_core::apq::ArcApqStorage =
187                Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
188            server.apq_store = Some(apq_store);
189            info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
190        }
191
192        Ok(server)
193    }
194}
195
196impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
197    /// Create new server with pre-configured Arrow Flight service.
198    ///
199    /// Use this constructor when you want to provide a Flight service with a real database adapter.
200    ///
201    /// # Arguments
202    ///
203    /// * `config` - Server configuration
204    /// * `schema` - Compiled GraphQL schema
205    /// * `adapter` - Database adapter
206    /// * `db_pool` - Database connection pool (optional, required for observers)
207    /// * `flight_service` - Pre-configured Flight service (only available with arrow feature)
208    ///
209    /// # Errors
210    ///
211    /// Returns error if OIDC validator initialization fails.
212    #[cfg(feature = "arrow")]
213    pub async fn with_flight_service(
214        config: ServerConfig,
215        schema: CompiledSchema,
216        adapter: Arc<A>,
217        #[allow(unused_variables)]
218        // Reason: used inside #[cfg(feature = "observers")] block; unused when feature is off
219        db_pool: Option<sqlx::PgPool>,
220        flight_service: Option<FraiseQLFlightService>,
221    ) -> Result<Self> {
222        // Same boot gates as `Server::new` — these must not drift by constructor (H16).
223        // Refuse to boot on at-rest-encryption-marked fields (H12, plaintext write path).
224        crate::server::initialization::field_encryption_unsupported_check(&schema)?;
225        // Build the runtime config from the compiled schema (validates format version,
226        // reads the audit flag, applies the #421 page-size ceiling + change-log toggle).
227        let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
228            super::ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
229        })?;
230
231        // Read security configs from compiled schema BEFORE schema is moved.
232        #[cfg(feature = "federation")]
233        let circuit_breaker = schema.federation.as_ref().and_then(
234            crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
235        );
236        // Non-federation builds construct the `Server` struct literal below with the
237        // `circuit_breaker` field cfg'd out, so there is no placeholder to bind here —
238        // just mark `schema.federation` read to mirror the federation branch.
239        #[cfg(not(feature = "federation"))]
240        let _ = &schema.federation;
241        let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
242        #[cfg(feature = "auth")]
243        let state_encryption = Self::state_encryption_from_schema(&schema)?;
244        #[cfg(not(feature = "auth"))]
245        let _state_encryption: Option<
246            std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
247        > = None;
248        #[cfg(feature = "auth")]
249        let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
250        #[cfg(not(feature = "auth"))]
251        let _pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
252        #[cfg(feature = "auth")]
253        let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
254        #[cfg(not(feature = "auth"))]
255        let _oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
256        let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
257        let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
258        let service_account_authenticator =
259            crate::service_account::service_account_authenticator_from_schema(&schema);
260        let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
261        let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
262        let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
263
264        let executor = Arc::new(Executor::with_config(schema.clone(), adapter, executor_config));
265        let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
266
267        // Initialize OIDC validator if auth is configured
268        #[cfg(feature = "auth")]
269        let oidc_validator = if let Some(ref auth_config) = config.auth {
270            info!(
271                issuer = %auth_config.issuer,
272                "Initializing OIDC authentication"
273            );
274            let validator = OidcValidator::new(auth_config.clone())
275                .await
276                .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
277            Some(Arc::new(validator))
278        } else {
279            None
280        };
281        #[cfg(not(feature = "auth"))]
282        let oidc_validator: Option<Arc<fraiseql_core::security::OidcValidator>> = None;
283
284        // Initialize HS256 validator if configured (mutually exclusive with OIDC).
285        let hs256_auth = super::builder::build_hs256_auth(&config)?;
286
287        // Initialize rate limiter: compiled schema config takes priority over server config.
288        let rate_limiter = if let Some(rl) = schema_rate_limiter {
289            Some(rl)
290        } else if let Some(ref rate_config) = config.rate_limiting {
291            if rate_config.enabled {
292                info!(
293                    rps_per_ip = rate_config.rps_per_ip,
294                    rps_per_user = rate_config.rps_per_user,
295                    "Initializing rate limiting from server config"
296                );
297                Some(Arc::new(RateLimiter::new(rate_config.clone())))
298            } else {
299                info!("Rate limiting disabled by configuration");
300                None
301            }
302        } else {
303            None
304        };
305
306        // Initialize observer runtime
307        #[cfg(feature = "observers")]
308        let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
309
310        // Warn if PKCE is configured but [auth] is missing.
311        #[cfg(feature = "auth")]
312        if pkce_store.is_some() && oidc_server_client.is_none() {
313            tracing::error!(
314                "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
315                 Auth routes will NOT be mounted."
316            );
317        }
318
319        // Refuse to start if FRAISEQL_REQUIRE_REDIS is set and PKCE store is in-memory.
320        #[cfg(feature = "auth")]
321        Self::check_redis_requirement(pkce_store.as_ref())?;
322
323        // Spawn background PKCE state cleanup task (every 5 minutes).
324        #[cfg(feature = "auth")]
325        Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
326
327        let apq_enabled = config.apq_enabled;
328
329        Ok(Self {
330            config,
331            executor,
332            subscription_manager,
333            subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
334            max_subscriptions_per_connection: None,
335            oidc_validator,
336            hs256_auth,
337            rate_limiter,
338            #[cfg(feature = "secrets")]
339            secrets_manager: None,
340            #[cfg(feature = "federation")]
341            circuit_breaker,
342            error_sanitizer,
343            #[cfg(feature = "auth")]
344            state_encryption,
345            #[cfg(feature = "auth")]
346            pkce_store,
347            #[cfg(feature = "auth")]
348            oidc_server_client,
349            #[cfg(feature = "auth")]
350            social_login: None,
351            #[cfg(feature = "auth")]
352            anon_signup_state: None,
353            #[cfg(feature = "auth")]
354            mfa_state: None,
355            api_key_authenticator,
356            service_account_authenticator,
357            revocation_manager,
358            apq_store: if apq_enabled {
359                Some(Arc::new(fraiseql_core::apq::InMemoryApqStorage::default())
360                    as fraiseql_core::apq::ArcApqStorage)
361            } else {
362                None
363            },
364            trusted_docs,
365            #[cfg(feature = "mcp")]
366            mcp_config: None,
367            pool_tuning_config: None,
368            #[cfg(feature = "observers")]
369            observer_runtime,
370            #[cfg(feature = "auth")]
371            enrichment_pool: db_pool.clone(),
372            #[cfg(feature = "observers")]
373            db_pool,
374            storage_state: None,
375            realtime_state: None,
376            #[cfg(feature = "functions-runtime")]
377            functions_hooks: None,
378            tenant_executor_factory: None,
379            #[cfg(feature = "arrow")]
380            flight_service,
381            adapter_cache_enabled: false,
382            broadcast_manager: None,
383            presence_manager: None,
384            storage_backend: None,
385            storage_max_upload_bytes: 100 * 1024 * 1024, // 100 MiB default
386            #[cfg(feature = "functions")]
387            function_store: None,
388            #[cfg(feature = "functions")]
389            function_runtime: None,
390            usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
391            tasks,
392        })
393    }
394
395    /// Initialize observer runtime from configuration.
396    ///
397    /// # Errors
398    ///
399    /// Returns `ServerError::ConfigError` when a non-Postgres observer transport
400    /// is selected that this binary cannot run (feature not compiled in, or NATS
401    /// without a URL) while in production mode (#350), or when the transport
402    /// configuration is otherwise invalid. In development such a selection is
403    /// downgraded to a warning and the runtime falls back to PostgreSQL.
404    #[cfg(feature = "observers")]
405    pub(super) async fn init_observer_runtime(
406        config: &ServerConfig,
407        pool: Option<&sqlx::PgPool>,
408    ) -> crate::Result<Option<Arc<RwLock<ObserverRuntime>>>> {
409        use fraiseql_observers::config::TransportKind;
410
411        // Check if enabled
412        let observer_config = match &config.observers {
413            Some(cfg) if cfg.enabled => cfg,
414            _ => {
415                info!("Observer runtime disabled");
416                return Ok(None);
417            },
418        };
419
420        let Some(pool) = pool else {
421            warn!("No database pool provided for observers");
422            return Ok(None);
423        };
424
425        info!("Initializing observer runtime");
426
427        // Resolve the event transport from compiled config + env overrides, then
428        // fail loud (#350) on a selection this binary cannot run before validating
429        // the finer NATS/JetStream bounds — never a silent fallback to PostgreSQL.
430        let mut transport = observer_config.runtime.transport.clone().with_env_overrides();
431        let compiled_in = cfg!(feature = "observers-nats");
432        let nats_url_present = !transport.nats.url.is_empty();
433        crate::server::initialization::observer_transport_check(
434            transport.transport,
435            compiled_in,
436            nats_url_present,
437            crate::ServerConfig::is_production_mode(),
438        )?;
439
440        // In production an unrunnable selection already returned above; the only
441        // way past the guard with an unrunnable transport is development, where it
442        // was downgraded to a warning — fall back to PostgreSQL so local boot works.
443        let usable = match transport.transport {
444            TransportKind::Postgres | TransportKind::InMemory => true,
445            TransportKind::Nats => compiled_in && nats_url_present,
446            _ => false,
447        };
448        if !usable {
449            transport.transport = TransportKind::Postgres;
450        }
451
452        transport.validate().map_err(|e| {
453            crate::ServerError::ConfigError(format!("invalid observer transport config: {e}"))
454        })?;
455
456        let runtime_config = ObserverRuntimeConfig::new(pool.clone())
457            .with_poll_interval(observer_config.runtime.poll_interval_ms)
458            .with_batch_size(observer_config.runtime.batch_size)
459            .with_channel_capacity(observer_config.runtime.channel_capacity)
460            .with_max_dlq_size(observer_config.runtime.max_dlq_size)
461            .with_transport(transport)
462            .with_email(observer_config.runtime.email.clone())
463            .with_log_payloads(observer_config.runtime.log_payloads);
464
465        let runtime = ObserverRuntime::new(runtime_config);
466        Ok(Some(Arc::new(RwLock::new(runtime))))
467    }
468}