1use std::sync::Arc;
4
5#[cfg(feature = "arrow")]
6use fraiseql_arrow::FraiseQLFlightService;
7use fraiseql_core::{
8 cache::{CacheConfig, CachedDatabaseAdapter, QueryResultCache},
9 db::traits::DatabaseAdapter,
10 runtime::{Executor, RuntimeConfig, SubscriptionManager},
11 schema::CompiledSchema,
12 security::{AuthConfig, AuthMiddleware, OidcValidator},
13};
14use tracing::{info, warn};
15
16use super::{RateLimiter, Result, Server, ServerConfig, ServerError};
17
18pub(super) fn build_hs256_auth(config: &ServerConfig) -> Result<Option<Arc<AuthMiddleware>>> {
20 let Some(ref hs) = config.auth_hs256 else {
21 return Ok(None);
22 };
23 let secret = hs
24 .load_secret()
25 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize HS256 auth: {e}")))?;
26 let mut auth_config = AuthConfig::with_hs256(&secret);
27 if let Some(ref iss) = hs.issuer {
28 auth_config = auth_config.with_issuer(iss);
29 }
30 if let Some(ref aud) = hs.audience {
31 auth_config = auth_config.with_audience(aud);
32 }
33 info!(
34 secret_env = %hs.secret_env,
35 issuer = ?hs.issuer,
36 audience = ?hs.audience,
37 "Initializing HS256 authentication (local validation, no network)"
38 );
39 Ok(Some(Arc::new(AuthMiddleware::from_config(auth_config))))
40}
41
42impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<CachedDatabaseAdapter<A>> {
43 #[allow(clippy::cognitive_complexity)] pub async fn new(
81 config: ServerConfig,
82 schema: CompiledSchema,
83 adapter: Arc<A>,
84 db_pool: Option<sqlx::PgPool>,
85 ) -> Result<Self> {
86 if schema.schema_format_version.is_none() {
89 warn!(
90 "Loaded schema has no schema_format_version (pre-v2.1 format). \
91 Re-compile with the current fraiseql-cli for version compatibility checking."
92 );
93 }
94 schema.validate_format_version().map_err(|msg| {
95 ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
96 })?;
97
98 #[cfg(feature = "federation")]
100 let circuit_breaker = schema.federation.as_ref().and_then(
101 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
102 );
103 #[cfg(not(feature = "federation"))]
104 let circuit_breaker: Option<()> = None;
105 #[cfg(not(feature = "federation"))]
106 let _ = &schema.federation;
107 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
108 #[cfg(feature = "auth")]
109 let state_encryption = Self::state_encryption_from_schema(&schema)?;
110 #[cfg(not(feature = "auth"))]
111 let state_encryption: Option<
112 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
113 > = None;
114 #[cfg(feature = "auth")]
115 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await;
116 #[cfg(not(feature = "auth"))]
117 let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
118 #[cfg(feature = "auth")]
119 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
120 #[cfg(not(feature = "auth"))]
121 let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
122 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await;
123 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
124 if api_key_authenticator.is_some() {
125 info!("API key authentication enabled");
126 }
127 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema);
128 if revocation_manager.is_some() {
129 info!("Token revocation enabled");
130 }
131 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
134 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
135
136 if config.cache_enabled && !schema.has_rls_configured() {
141 if schema.is_multi_tenant() {
142 return Err(ServerError::ConfigError(
144 "Cache is enabled in a multi-tenant schema but no Row-Level Security \
145 policies are declared. This would allow cross-tenant cache hits and \
146 data leakage. In fraiseql.toml, either disable caching with \
147 [cache] enabled = false, declare [security.rls] policies, or set \
148 [security] multi_tenant = false to acknowledge single-tenant mode."
149 .to_string(),
150 ));
151 }
152 warn!(
154 "Query-result caching is enabled but no Row-Level Security policies are \
155 declared in the compiled schema. This is safe for single-tenant deployments. \
156 For multi-tenant deployments, declare RLS policies and set \
157 `security.multi_tenant = true` in your schema."
158 );
159 }
160
161 let cache_config = CacheConfig::from(config.cache_enabled);
163 let cache = QueryResultCache::new(cache_config);
164
165 if cache_config.enabled {
167 tracing::info!(
168 max_entries = cache_config.max_entries,
169 ttl_seconds = cache_config.ttl_seconds,
170 rls_enforcement = ?cache_config.rls_enforcement,
171 "Query result cache: active"
172 );
173 } else {
174 tracing::info!("Query result cache: disabled");
175 }
176
177 let subscriptions_config = schema.subscriptions_config.clone();
179
180 let inner = Arc::into_inner(adapter)
182 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
183 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
184 .with_ttl_overrides_from_schema(&schema)
185 .with_rls(schema.has_rls_configured());
186
187 let audit_mutations = schema
190 .security
191 .as_ref()
192 .and_then(|s| s.additional.get("enterprise"))
193 .and_then(|e| e.get("audit_logging_enabled"))
194 .and_then(|v| v.as_bool())
195 .unwrap_or(false);
196 if audit_mutations {
197 info!("Mutation audit logging enabled (target: fraiseql::mutation_audit)");
198 }
199 let executor_config = RuntimeConfig {
200 audit_mutations,
201 ..RuntimeConfig::default()
202 };
203 let executor =
204 Arc::new(Executor::with_config(schema.clone(), Arc::new(cached), executor_config));
205 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
206
207 let mut server = Self::from_executor(
208 config,
209 executor,
210 subscription_manager,
211 circuit_breaker,
212 error_sanitizer,
213 state_encryption,
214 pkce_store,
215 oidc_server_client,
216 schema_rate_limiter,
217 api_key_authenticator,
218 revocation_manager,
219 trusted_docs,
220 db_pool,
221 tasks,
222 )
223 .await?;
224
225 server.adapter_cache_enabled = cache_config.enabled;
226
227 if let Some(pt) = server.config.pool_tuning.clone() {
229 if pt.enabled {
230 server = server
231 .with_pool_tuning(pt)
232 .map_err(|e| ServerError::ConfigError(format!("pool_tuning: {e}")))?;
233 }
234 }
235
236 #[cfg(feature = "mcp")]
238 if let Some(ref cfg) = server.executor.schema().mcp_config {
239 if cfg.enabled {
240 let tool_count =
241 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
242 info!(
243 path = %cfg.path,
244 transport = %cfg.transport,
245 tools = tool_count,
246 "MCP server configured"
247 );
248 server.mcp_config = Some(cfg.clone());
249 }
250 }
251
252 if server.config.apq_enabled {
254 let apq_store: fraiseql_core::apq::ArcApqStorage =
255 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
256 server.apq_store = Some(apq_store);
257 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
258 }
259
260 if let Some(ref subs) = subscriptions_config {
262 if let Some(max) = subs.max_subscriptions_per_connection {
263 server.max_subscriptions_per_connection = Some(max);
264 }
265 if let Some(lifecycle) = crate::subscriptions::WebhookLifecycle::from_config(subs) {
266 server.subscription_lifecycle = Arc::new(lifecycle);
267 }
268 }
269
270 Ok(server)
271 }
272}
273
274impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
275 #[allow(clippy::too_many_arguments)]
280 #[allow(clippy::cognitive_complexity)] pub(super) async fn from_executor(
284 config: ServerConfig,
285 executor: Arc<Executor<A>>,
286 subscription_manager: Arc<SubscriptionManager>,
287 #[cfg(feature = "federation")] circuit_breaker: Option<
288 Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>,
289 >,
290 #[cfg(not(feature = "federation"))] _circuit_breaker: Option<()>,
291 error_sanitizer: Arc<crate::config::error_sanitization::ErrorSanitizer>,
292 state_encryption: Option<Arc<crate::auth::state_encryption::StateEncryptionService>>,
293 pkce_store: Option<Arc<crate::auth::PkceStateStore>>,
294 oidc_server_client: Option<Arc<crate::auth::OidcServerClient>>,
295 schema_rate_limiter: Option<Arc<RateLimiter>>,
296 api_key_authenticator: Option<Arc<crate::api_key::ApiKeyAuthenticator>>,
297 revocation_manager: Option<Arc<crate::token_revocation::TokenRevocationManager>>,
298 trusted_docs: Option<Arc<crate::trusted_documents::TrustedDocumentStore>>,
299 #[cfg_attr(
301 not(any(feature = "observers", feature = "auth")),
302 allow(unused_variables)
303 )]
304 db_pool: Option<sqlx::PgPool>,
305 mut tasks: tokio::task::JoinSet<()>,
306 ) -> Result<Self> {
307 let oidc_validator = if let Some(ref auth_config) = config.auth {
309 info!(
310 issuer = %auth_config.issuer,
311 "Initializing OIDC authentication"
312 );
313 let validator = OidcValidator::new(auth_config.clone())
314 .await
315 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
316 Some(Arc::new(validator))
317 } else {
318 None
319 };
320
321 let hs256_auth = build_hs256_auth(&config)?;
323
324 let rate_limiter = if let Some(rl) = schema_rate_limiter {
326 Some(rl)
327 } else if let Some(ref rate_config) = config.rate_limiting {
328 if rate_config.enabled {
329 info!(
330 rps_per_ip = rate_config.rps_per_ip,
331 rps_per_user = rate_config.rps_per_user,
332 "Initializing rate limiting from server config"
333 );
334 Some(Arc::new(RateLimiter::new(rate_config.clone())))
335 } else {
336 info!("Rate limiting disabled by configuration");
337 None
338 }
339 } else {
340 None
341 };
342
343 #[cfg(feature = "observers")]
345 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await;
346
347 #[cfg(feature = "arrow")]
349 let flight_service = {
350 let mut service = FraiseQLFlightService::new();
351 if let Some(ref validator) = oidc_validator {
352 info!("Enabling OIDC authentication for Arrow Flight");
353 service.set_oidc_validator(validator.clone());
354 } else {
355 info!("Arrow Flight initialized without authentication (dev mode)");
356 }
357 Some(service)
358 };
359
360 #[cfg(feature = "auth")]
362 if pkce_store.is_some() && oidc_server_client.is_none() {
363 tracing::error!(
364 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
365 Auth routes (/auth/start, /auth/callback) will NOT be mounted. \
366 Add [auth] with discovery_url, client_id, client_secret_env, and \
367 server_redirect_uri to fraiseql.toml and recompile the schema."
368 );
369 }
370
371 #[cfg(feature = "auth")]
373 Self::check_redis_requirement(pkce_store.as_ref())?;
374
375 #[cfg(feature = "auth")]
377 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
378
379 #[cfg(not(feature = "auth"))]
382 let _ = (state_encryption, pkce_store, oidc_server_client);
383 Ok(Self {
384 config,
385 executor,
386 subscription_manager,
387 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
388 max_subscriptions_per_connection: None,
389 oidc_validator,
390 hs256_auth,
391 rate_limiter,
392 #[cfg(feature = "secrets")]
393 secrets_manager: None,
394 #[cfg(feature = "federation")]
395 circuit_breaker,
396 error_sanitizer,
397 #[cfg(feature = "auth")]
398 state_encryption,
399 #[cfg(feature = "auth")]
400 pkce_store,
401 #[cfg(feature = "auth")]
402 oidc_server_client,
403 #[cfg(feature = "auth")]
404 social_login: None,
405 #[cfg(feature = "auth")]
406 mfa_state: None,
407 #[cfg(feature = "auth")]
408 anon_signup_state: None,
409 api_key_authenticator,
410 revocation_manager,
411 apq_store: None,
412 trusted_docs,
413 #[cfg(feature = "observers")]
414 observer_runtime,
415 #[cfg(feature = "auth")]
416 enrichment_pool: db_pool.clone(),
417 #[cfg(feature = "observers")]
418 db_pool,
419 storage_state: None,
420 realtime_state: None,
421 #[cfg(feature = "arrow")]
422 flight_service,
423 #[cfg(feature = "mcp")]
424 mcp_config: None,
425 pool_tuning_config: None,
426 adapter_cache_enabled: false,
427 broadcast_manager: None,
428 presence_manager: None,
429 storage_backend: None,
430 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
432 function_store: None,
433 #[cfg(feature = "functions")]
434 function_runtime: None,
435 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
436 tasks,
437 })
438 }
439
440 #[cfg(feature = "auth")]
445 pub(super) fn spawn_pkce_cleanup(
446 pkce_store: Option<&Arc<crate::auth::PkceStateStore>>,
447 tasks: &mut tokio::task::JoinSet<()>,
448 ) {
449 use std::time::Duration;
450
451 use tokio::time::MissedTickBehavior;
452
453 if let Some(store) = pkce_store {
454 let store_clone = Arc::clone(store);
455 tasks.spawn(async move {
456 let mut ticker = tokio::time::interval(Duration::from_secs(300));
457 ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
458 loop {
459 ticker.tick().await;
460 store_clone.cleanup_expired().await;
461 }
462 });
463 }
464 }
465
466 #[must_use]
468 pub fn with_subscription_lifecycle(
469 mut self,
470 lifecycle: Arc<dyn crate::subscriptions::SubscriptionLifecycle>,
471 ) -> Self {
472 self.subscription_lifecycle = lifecycle;
473 self
474 }
475
476 #[must_use]
478 pub const fn with_max_subscriptions_per_connection(mut self, max: u32) -> Self {
479 self.max_subscriptions_per_connection = Some(max);
480 self
481 }
482
483 #[must_use]
489 pub fn with_realtime(mut self, state: crate::realtime::server::RealtimeState) -> Self {
490 self.realtime_state = Some(state);
491 self
492 }
493
494 #[must_use]
496 pub fn with_broadcast(mut self, config: crate::subscriptions::BroadcastConfig) -> Self {
497 self.broadcast_manager =
498 Some(Arc::new(crate::subscriptions::BroadcastManager::new(config)));
499 self
500 }
501
502 #[must_use]
504 pub fn with_presence(mut self, config: crate::subscriptions::PresenceConfig) -> Self {
505 self.presence_manager = Some(Arc::new(crate::subscriptions::PresenceManager::new(config)));
506 self
507 }
508
509 pub fn with_pool_tuning(
518 mut self,
519 config: crate::config::pool_tuning::PoolPressureMonitorConfig,
520 ) -> std::result::Result<Self, String> {
521 config.validate()?;
522 self.pool_tuning_config = Some(config);
523 Ok(self)
524 }
525
526 #[cfg(feature = "auth")]
540 #[must_use]
541 pub fn with_social_login(
542 mut self,
543 social_login: Arc<crate::auth::social::SocialLoginState>,
544 ) -> Self {
545 self.social_login = Some(social_login);
546 self
547 }
548
549 #[cfg(feature = "auth")]
555 #[must_use]
556 pub fn with_anon_signup(mut self, state: Arc<crate::auth::AnonSignupState>) -> Self {
557 self.anon_signup_state = Some(state);
558 self
559 }
560
561 #[cfg(feature = "auth")]
570 #[must_use]
571 pub fn with_mfa(mut self, mfa_state: Arc<crate::auth::MfaRouteState>) -> Self {
572 self.mfa_state = Some(mfa_state);
573 self
574 }
575
576 #[must_use]
585 pub fn with_storage(mut self, backend: Arc<dyn crate::storage::StorageBackend>) -> Self {
586 self.storage_backend = Some(backend);
587 self
588 }
589
590 #[must_use]
595 pub const fn with_storage_max_upload_bytes(mut self, bytes: usize) -> Self {
596 self.storage_max_upload_bytes = bytes;
597 self
598 }
599
600 #[cfg(feature = "functions")]
606 #[must_use]
607 pub fn with_functions(
608 mut self,
609 store: Arc<dyn fraiseql_functions::FunctionStore>,
610 runtime: Arc<dyn fraiseql_functions::runtime::SendFunctionRuntime>,
611 ) -> Self {
612 self.function_store = Some(store);
613 self.function_runtime = Some(runtime);
614 self
615 }
616
617 #[cfg(feature = "secrets")]
621 pub fn set_secrets_manager(&mut self, manager: Arc<crate::secrets_manager::SecretsManager>) {
622 self.secrets_manager = Some(manager);
623 info!("Secrets manager attached to server");
624 }
625
626 #[cfg(feature = "mcp")]
636 pub async fn serve_mcp_stdio(self) -> Result<()> {
637 use rmcp::ServiceExt;
638
639 let mcp_cfg = self.mcp_config.ok_or_else(|| {
640 ServerError::ConfigError(
641 "FRAISEQL_MCP_STDIO=1 but MCP is not configured. \
642 Add [mcp] enabled = true to fraiseql.toml and recompile the schema."
643 .into(),
644 )
645 })?;
646
647 let schema = Arc::new(self.executor.schema().clone());
648 let executor = self.executor.clone();
649
650 let service = crate::mcp::handler::FraiseQLMcpService::new(schema, executor, mcp_cfg);
651
652 info!("MCP stdio transport starting — reading from stdin, writing to stdout");
653
654 let running = service
655 .serve((tokio::io::stdin(), tokio::io::stdout()))
656 .await
657 .map_err(|e| ServerError::ConfigError(format!("MCP stdio init failed: {e}")))?;
658
659 running
660 .waiting()
661 .await
662 .map_err(|e| ServerError::ConfigError(format!("MCP stdio error: {e}")))?;
663
664 Ok(())
665 }
666}