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 hs.validate()
28 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize HS256 auth: {e}")))?;
29 let secret = hs
30 .load_secret()
31 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize HS256 auth: {e}")))?;
32 let mut auth_config = AuthConfig::with_hs256(&secret);
33 if let Some(ref iss) = hs.issuer {
34 auth_config = auth_config.with_issuer(iss);
35 }
36 if let Some(ref aud) = hs.audience {
37 auth_config = auth_config.with_audience(aud);
38 }
39 info!(
40 secret_env = %hs.secret_env,
41 issuer = ?hs.issuer,
42 audience = ?hs.audience,
43 "Initializing HS256 authentication (local validation, no network)"
44 );
45 Ok(Some(Arc::new(AuthMiddleware::from_config(auth_config))))
46}
47
48impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<CachedDatabaseAdapter<A>> {
49 #[allow(clippy::cognitive_complexity)] pub async fn new(
87 config: ServerConfig,
88 schema: CompiledSchema,
89 adapter: Arc<A>,
90 db_pool: Option<sqlx::PgPool>,
91 ) -> Result<Self> {
92 let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
99 ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
100 })?;
101
102 crate::server::initialization::field_encryption_unsupported_check(&schema)?;
106
107 #[cfg(feature = "federation")]
109 let circuit_breaker = schema.federation.as_ref().and_then(
110 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
111 );
112 #[cfg(not(feature = "federation"))]
113 let circuit_breaker: Option<()> = None;
114 #[cfg(not(feature = "federation"))]
115 let _ = &schema.federation;
116 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
117 #[cfg(feature = "auth")]
118 let state_encryption = Self::state_encryption_from_schema(&schema)?;
119 #[cfg(not(feature = "auth"))]
120 let state_encryption: Option<
121 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
122 > = None;
123 #[cfg(feature = "auth")]
124 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
125 #[cfg(not(feature = "auth"))]
126 let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
127 #[cfg(feature = "auth")]
128 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
129 #[cfg(not(feature = "auth"))]
130 let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
131 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
132 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
133 if api_key_authenticator.is_some() {
134 info!("API key authentication enabled");
135 }
136 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
137 if revocation_manager.is_some() {
138 info!("Token revocation enabled");
139 }
140 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
143 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
144
145 if config.cache_enabled && !schema.has_rls_configured() {
150 if schema.is_multi_tenant() {
151 return Err(ServerError::ConfigError(
153 "Cache is enabled in a multi-tenant schema but no Row-Level Security \
154 policies are declared. This would allow cross-tenant cache hits and \
155 data leakage. In fraiseql.toml, either disable caching with \
156 [cache] enabled = false, declare [security.rls] policies, or set \
157 [security] multi_tenant = false to acknowledge single-tenant mode."
158 .to_string(),
159 ));
160 }
161 warn!(
163 "Query-result caching is enabled but no Row-Level Security policies are \
164 declared in the compiled schema. This is safe for single-tenant deployments. \
165 For multi-tenant deployments, declare RLS policies and set \
166 `security.multi_tenant = true` in your schema."
167 );
168 }
169
170 let cache_config = CacheConfig::from(config.cache_enabled);
172 let cache = QueryResultCache::new(cache_config);
173
174 if cache_config.enabled {
176 tracing::info!(
177 max_entries = cache_config.max_entries,
178 ttl_seconds = cache_config.ttl_seconds,
179 rls_enforcement = ?cache_config.rls_enforcement,
180 "Query result cache: active"
181 );
182 } else {
183 tracing::info!("Query result cache: disabled");
184 }
185
186 let subscriptions_config = schema.subscriptions_config.clone();
188
189 let inner = Arc::into_inner(adapter)
191 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
192 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
193 .with_ttl_overrides_from_schema(&schema)
194 .with_rls(schema.has_rls_configured());
195
196 let executor =
199 Arc::new(Executor::with_config(schema.clone(), Arc::new(cached), executor_config));
200 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
201
202 let mut server = Self::from_executor(
203 config,
204 executor,
205 subscription_manager,
206 circuit_breaker,
207 error_sanitizer,
208 state_encryption,
209 pkce_store,
210 oidc_server_client,
211 schema_rate_limiter,
212 api_key_authenticator,
213 revocation_manager,
214 trusted_docs,
215 db_pool,
216 tasks,
217 )
218 .await?;
219
220 server.adapter_cache_enabled = cache_config.enabled;
221
222 if let Some(pt) = server.config.pool_tuning.clone() {
224 if pt.enabled {
225 server = server
226 .with_pool_tuning(pt)
227 .map_err(|e| ServerError::ConfigError(format!("pool_tuning: {e}")))?;
228 }
229 }
230
231 #[cfg(feature = "mcp")]
233 if let Some(ref cfg) = server.executor.schema().mcp_config {
234 if cfg.enabled {
235 let tool_count =
236 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
237 info!(
238 path = %cfg.path,
239 transport = %cfg.transport,
240 tools = tool_count,
241 "MCP server configured"
242 );
243 server.mcp_config = Some(cfg.clone());
244 }
245 }
246
247 if server.config.apq_enabled {
249 let apq_store: fraiseql_core::apq::ArcApqStorage =
250 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
251 server.apq_store = Some(apq_store);
252 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
253 }
254
255 if let Some(ref subs) = subscriptions_config {
257 if let Some(max) = subs.max_subscriptions_per_connection {
258 server.max_subscriptions_per_connection = Some(max);
259 }
260 if let Some(lifecycle) = crate::subscriptions::WebhookLifecycle::from_config(subs) {
261 server.subscription_lifecycle = Arc::new(lifecycle);
262 }
263 }
264
265 Ok(server)
266 }
267}
268
269impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
270 #[allow(clippy::too_many_arguments)]
275 #[allow(clippy::cognitive_complexity)] pub(super) async fn from_executor(
279 config: ServerConfig,
280 executor: Arc<Executor<A>>,
281 subscription_manager: Arc<SubscriptionManager>,
282 #[cfg(feature = "federation")] circuit_breaker: Option<
283 Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>,
284 >,
285 #[cfg(not(feature = "federation"))] _circuit_breaker: Option<()>,
286 error_sanitizer: Arc<crate::config::error_sanitization::ErrorSanitizer>,
287 state_encryption: Option<Arc<crate::auth::state_encryption::StateEncryptionService>>,
288 pkce_store: Option<Arc<crate::auth::PkceStateStore>>,
289 oidc_server_client: Option<Arc<crate::auth::OidcServerClient>>,
290 schema_rate_limiter: Option<Arc<RateLimiter>>,
291 api_key_authenticator: Option<Arc<crate::api_key::ApiKeyAuthenticator>>,
292 revocation_manager: Option<Arc<crate::token_revocation::TokenRevocationManager>>,
293 trusted_docs: Option<Arc<crate::trusted_documents::TrustedDocumentStore>>,
294 #[cfg_attr(
296 not(any(feature = "observers", feature = "auth")),
297 allow(unused_variables)
298 )]
299 db_pool: Option<sqlx::PgPool>,
300 mut tasks: tokio::task::JoinSet<()>,
301 ) -> Result<Self> {
302 let oidc_validator = if let Some(ref auth_config) = config.auth {
304 info!(
305 issuer = %auth_config.issuer,
306 "Initializing OIDC authentication"
307 );
308 let validator = OidcValidator::new(auth_config.clone())
309 .await
310 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
311 Some(Arc::new(validator))
312 } else {
313 None
314 };
315
316 let hs256_auth = build_hs256_auth(&config)?;
318
319 let rate_limiter = if let Some(rl) = schema_rate_limiter {
321 Some(rl)
322 } else if let Some(ref rate_config) = config.rate_limiting {
323 if rate_config.enabled {
324 info!(
325 rps_per_ip = rate_config.rps_per_ip,
326 rps_per_user = rate_config.rps_per_user,
327 "Initializing rate limiting from server config"
328 );
329 Some(Arc::new(RateLimiter::new(rate_config.clone())))
330 } else {
331 info!("Rate limiting disabled by configuration");
332 None
333 }
334 } else {
335 None
336 };
337
338 #[cfg(feature = "observers")]
340 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
341
342 #[cfg(feature = "arrow")]
344 let flight_service = {
345 let mut service = FraiseQLFlightService::new();
346 if let Some(ref validator) = oidc_validator {
347 info!("Enabling OIDC authentication for Arrow Flight");
348 service.set_oidc_validator(validator.clone());
349 } else {
350 info!("Arrow Flight initialized without authentication (dev mode)");
351 }
352 Some(service)
353 };
354
355 #[cfg(feature = "auth")]
357 if pkce_store.is_some() && oidc_server_client.is_none() {
358 tracing::error!(
359 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
360 Auth routes (/auth/start, /auth/callback) will NOT be mounted. \
361 Add [auth] with discovery_url, client_id, client_secret_env, and \
362 server_redirect_uri to fraiseql.toml and recompile the schema."
363 );
364 }
365
366 #[cfg(feature = "auth")]
368 Self::check_redis_requirement(pkce_store.as_ref())?;
369
370 #[cfg(feature = "auth")]
372 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
373
374 #[cfg(not(feature = "auth"))]
377 let _ = (state_encryption, pkce_store, oidc_server_client);
378 Ok(Self {
379 config,
380 executor,
381 subscription_manager,
382 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
383 max_subscriptions_per_connection: None,
384 oidc_validator,
385 hs256_auth,
386 rate_limiter,
387 #[cfg(feature = "secrets")]
388 secrets_manager: None,
389 #[cfg(feature = "federation")]
390 circuit_breaker,
391 error_sanitizer,
392 #[cfg(feature = "auth")]
393 state_encryption,
394 #[cfg(feature = "auth")]
395 pkce_store,
396 #[cfg(feature = "auth")]
397 oidc_server_client,
398 #[cfg(feature = "auth")]
399 social_login: None,
400 #[cfg(feature = "auth")]
401 mfa_state: None,
402 #[cfg(feature = "auth")]
403 anon_signup_state: None,
404 api_key_authenticator,
405 revocation_manager,
406 apq_store: None,
407 trusted_docs,
408 #[cfg(feature = "observers")]
409 observer_runtime,
410 #[cfg(feature = "auth")]
411 enrichment_pool: db_pool.clone(),
412 #[cfg(feature = "observers")]
413 db_pool,
414 storage_state: None,
415 realtime_state: None,
416 #[cfg(feature = "functions-runtime")]
417 functions_hooks: None,
418 tenant_executor_factory: None,
419 #[cfg(feature = "arrow")]
420 flight_service,
421 #[cfg(feature = "mcp")]
422 mcp_config: None,
423 pool_tuning_config: None,
424 adapter_cache_enabled: false,
425 broadcast_manager: None,
426 presence_manager: None,
427 storage_backend: None,
428 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
430 function_store: None,
431 #[cfg(feature = "functions")]
432 function_runtime: None,
433 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
434 tasks,
435 })
436 }
437
438 #[cfg(feature = "auth")]
443 pub(super) fn spawn_pkce_cleanup(
444 pkce_store: Option<&Arc<crate::auth::PkceStateStore>>,
445 tasks: &mut tokio::task::JoinSet<()>,
446 ) {
447 use std::time::Duration;
448
449 use tokio::time::MissedTickBehavior;
450
451 if let Some(store) = pkce_store {
452 let store_clone = Arc::clone(store);
453 tasks.spawn(async move {
454 let mut ticker = tokio::time::interval(Duration::from_secs(300));
455 ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
456 loop {
457 ticker.tick().await;
458 store_clone.cleanup_expired().await;
459 }
460 });
461 }
462 }
463
464 #[must_use]
466 pub fn with_subscription_lifecycle(
467 mut self,
468 lifecycle: Arc<dyn crate::subscriptions::SubscriptionLifecycle>,
469 ) -> Self {
470 self.subscription_lifecycle = lifecycle;
471 self
472 }
473
474 #[must_use]
476 pub const fn with_max_subscriptions_per_connection(mut self, max: u32) -> Self {
477 self.max_subscriptions_per_connection = Some(max);
478 self
479 }
480
481 #[must_use]
491 pub fn with_tenant_executor_factory(
492 mut self,
493 factory: crate::tenancy::TenantExecutorFactory<A>,
494 ) -> Self {
495 self.tenant_executor_factory = Some(factory);
496 self
497 }
498
499 #[must_use]
505 pub fn with_realtime(mut self, state: crate::realtime::server::RealtimeState) -> Self {
506 self.realtime_state = Some(state);
507 self
508 }
509
510 #[must_use]
512 pub fn with_broadcast(mut self, config: crate::subscriptions::BroadcastConfig) -> Self {
513 self.broadcast_manager =
514 Some(Arc::new(crate::subscriptions::BroadcastManager::new(config)));
515 self
516 }
517
518 #[must_use]
520 pub fn with_presence(mut self, config: crate::subscriptions::PresenceConfig) -> Self {
521 self.presence_manager = Some(Arc::new(crate::subscriptions::PresenceManager::new(config)));
522 self
523 }
524
525 pub fn with_pool_tuning(
534 mut self,
535 config: crate::config::pool_tuning::PoolPressureMonitorConfig,
536 ) -> std::result::Result<Self, String> {
537 config.validate()?;
538 self.pool_tuning_config = Some(config);
539 Ok(self)
540 }
541
542 #[cfg(feature = "auth")]
556 #[must_use]
557 pub fn with_social_login(
558 mut self,
559 social_login: Arc<crate::auth::social::SocialLoginState>,
560 ) -> Self {
561 self.social_login = Some(social_login);
562 self
563 }
564
565 #[cfg(feature = "auth")]
571 #[must_use]
572 pub fn with_anon_signup(mut self, state: Arc<crate::auth::AnonSignupState>) -> Self {
573 self.anon_signup_state = Some(state);
574 self
575 }
576
577 #[cfg(feature = "auth")]
586 #[must_use]
587 pub fn with_mfa(mut self, mfa_state: Arc<crate::auth::MfaRouteState>) -> Self {
588 self.mfa_state = Some(mfa_state);
589 self
590 }
591
592 #[must_use]
601 pub fn with_storage(mut self, backend: Arc<dyn crate::storage::StorageBackend>) -> Self {
602 self.storage_backend = Some(backend);
603 self
604 }
605
606 #[must_use]
611 pub const fn with_storage_max_upload_bytes(mut self, bytes: usize) -> Self {
612 self.storage_max_upload_bytes = bytes;
613 self
614 }
615
616 #[must_use]
627 pub fn with_storage_state(mut self, state: fraiseql_storage::StorageState) -> Self {
628 self.storage_state = Some(state);
629 self
630 }
631
632 #[must_use]
640 pub fn with_revocation_manager(
641 mut self,
642 manager: Arc<crate::token_revocation::TokenRevocationManager>,
643 ) -> Self {
644 self.revocation_manager = Some(manager);
645 self
646 }
647
648 #[cfg(feature = "functions")]
654 #[must_use]
655 pub fn with_functions(
656 mut self,
657 store: Arc<dyn fraiseql_functions::FunctionStore>,
658 runtime: Arc<dyn fraiseql_functions::runtime::SendFunctionRuntime>,
659 ) -> Self {
660 self.function_store = Some(store);
661 self.function_runtime = Some(runtime);
662 self
663 }
664
665 #[cfg(feature = "secrets")]
669 pub fn set_secrets_manager(&mut self, manager: Arc<crate::secrets_manager::SecretsManager>) {
670 self.secrets_manager = Some(manager);
671 info!("Secrets manager attached to server");
672 }
673
674 #[cfg(feature = "mcp")]
684 pub async fn serve_mcp_stdio(self) -> Result<()> {
685 use rmcp::ServiceExt;
686
687 let mcp_cfg = self.mcp_config.ok_or_else(|| {
688 ServerError::ConfigError(
689 "FRAISEQL_MCP_STDIO=1 but MCP is not configured. \
690 Add [mcp] enabled = true to fraiseql.toml and recompile the schema."
691 .into(),
692 )
693 })?;
694
695 let schema = Arc::new(self.executor.schema().clone());
696 let executor = self.executor.clone();
697
698 let service = crate::mcp::handler::FraiseQLMcpService::new(schema, executor, mcp_cfg)
699 .with_oidc_validator(self.oidc_validator.clone());
700
701 info!("MCP stdio transport starting — reading from stdin, writing to stdout");
702
703 let running = service
704 .serve((tokio::io::stdin(), tokio::io::stdout()))
705 .await
706 .map_err(|e| ServerError::ConfigError(format!("MCP stdio init failed: {e}")))?;
707
708 running
709 .waiting()
710 .await
711 .map_err(|e| ServerError::ConfigError(format!("MCP stdio error: {e}")))?;
712
713 Ok(())
714 }
715}