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 if schema.schema_format_version.is_none() {
95 warn!(
96 "Loaded schema has no schema_format_version (pre-v2.1 format). \
97 Re-compile with the current fraiseql-cli for version compatibility checking."
98 );
99 }
100 schema.validate_format_version().map_err(|msg| {
101 ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
102 })?;
103
104 crate::server::initialization::field_encryption_unsupported_check(&schema)?;
108
109 #[cfg(feature = "federation")]
111 let circuit_breaker = schema.federation.as_ref().and_then(
112 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
113 );
114 #[cfg(not(feature = "federation"))]
115 let circuit_breaker: Option<()> = None;
116 #[cfg(not(feature = "federation"))]
117 let _ = &schema.federation;
118 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
119 #[cfg(feature = "auth")]
120 let state_encryption = Self::state_encryption_from_schema(&schema)?;
121 #[cfg(not(feature = "auth"))]
122 let state_encryption: Option<
123 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
124 > = None;
125 #[cfg(feature = "auth")]
126 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
127 #[cfg(not(feature = "auth"))]
128 let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
129 #[cfg(feature = "auth")]
130 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
131 #[cfg(not(feature = "auth"))]
132 let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
133 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
134 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
135 if api_key_authenticator.is_some() {
136 info!("API key authentication enabled");
137 }
138 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
139 if revocation_manager.is_some() {
140 info!("Token revocation enabled");
141 }
142 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
145 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
146
147 if config.cache_enabled && !schema.has_rls_configured() {
152 if schema.is_multi_tenant() {
153 return Err(ServerError::ConfigError(
155 "Cache is enabled in a multi-tenant schema but no Row-Level Security \
156 policies are declared. This would allow cross-tenant cache hits and \
157 data leakage. In fraiseql.toml, either disable caching with \
158 [cache] enabled = false, declare [security.rls] policies, or set \
159 [security] multi_tenant = false to acknowledge single-tenant mode."
160 .to_string(),
161 ));
162 }
163 warn!(
165 "Query-result caching is enabled but no Row-Level Security policies are \
166 declared in the compiled schema. This is safe for single-tenant deployments. \
167 For multi-tenant deployments, declare RLS policies and set \
168 `security.multi_tenant = true` in your schema."
169 );
170 }
171
172 let cache_config = CacheConfig::from(config.cache_enabled);
174 let cache = QueryResultCache::new(cache_config);
175
176 if cache_config.enabled {
178 tracing::info!(
179 max_entries = cache_config.max_entries,
180 ttl_seconds = cache_config.ttl_seconds,
181 rls_enforcement = ?cache_config.rls_enforcement,
182 "Query result cache: active"
183 );
184 } else {
185 tracing::info!("Query result cache: disabled");
186 }
187
188 let subscriptions_config = schema.subscriptions_config.clone();
190
191 let inner = Arc::into_inner(adapter)
193 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
194 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
195 .with_ttl_overrides_from_schema(&schema)
196 .with_rls(schema.has_rls_configured());
197
198 let audit_mutations = schema
201 .security
202 .as_ref()
203 .and_then(|s| s.additional.get("enterprise"))
204 .and_then(|e| e.get("audit_logging_enabled"))
205 .and_then(|v| v.as_bool())
206 .unwrap_or(false);
207 if audit_mutations {
208 info!("Mutation audit logging enabled (target: fraiseql::mutation_audit)");
209 }
210 let max_page_size = page_size_precedence(
213 std::env::var("FRAISEQL_MAX_PAGE_SIZE").ok().as_deref(),
214 schema.validation_config.as_ref().and_then(|v| v.max_page_size),
215 );
216 let changelog_enabled = std::env::var("FRAISEQL_CHANGELOG_ENABLED")
220 .ok()
221 .map(|v| {
222 !matches!(v.trim().to_ascii_lowercase().as_str(), "false" | "0" | "no" | "off")
223 })
224 .or_else(|| schema.changelog.as_ref().map(|c| c.write_enabled))
225 .unwrap_or(true);
226 if !changelog_enabled {
227 info!(
228 "Change-log outbox write disabled (FRAISEQL_CHANGELOG_ENABLED / [changelog] write_enabled)"
229 );
230 }
231 let executor_config = RuntimeConfig {
232 audit_mutations,
233 max_page_size,
234 changelog_enabled,
235 ..RuntimeConfig::default()
236 };
237 let executor =
238 Arc::new(Executor::with_config(schema.clone(), Arc::new(cached), executor_config));
239 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
240
241 let mut server = Self::from_executor(
242 config,
243 executor,
244 subscription_manager,
245 circuit_breaker,
246 error_sanitizer,
247 state_encryption,
248 pkce_store,
249 oidc_server_client,
250 schema_rate_limiter,
251 api_key_authenticator,
252 revocation_manager,
253 trusted_docs,
254 db_pool,
255 tasks,
256 )
257 .await?;
258
259 server.adapter_cache_enabled = cache_config.enabled;
260
261 if let Some(pt) = server.config.pool_tuning.clone() {
263 if pt.enabled {
264 server = server
265 .with_pool_tuning(pt)
266 .map_err(|e| ServerError::ConfigError(format!("pool_tuning: {e}")))?;
267 }
268 }
269
270 #[cfg(feature = "mcp")]
272 if let Some(ref cfg) = server.executor.schema().mcp_config {
273 if cfg.enabled {
274 let tool_count =
275 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
276 info!(
277 path = %cfg.path,
278 transport = %cfg.transport,
279 tools = tool_count,
280 "MCP server configured"
281 );
282 server.mcp_config = Some(cfg.clone());
283 }
284 }
285
286 if server.config.apq_enabled {
288 let apq_store: fraiseql_core::apq::ArcApqStorage =
289 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
290 server.apq_store = Some(apq_store);
291 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
292 }
293
294 if let Some(ref subs) = subscriptions_config {
296 if let Some(max) = subs.max_subscriptions_per_connection {
297 server.max_subscriptions_per_connection = Some(max);
298 }
299 if let Some(lifecycle) = crate::subscriptions::WebhookLifecycle::from_config(subs) {
300 server.subscription_lifecycle = Arc::new(lifecycle);
301 }
302 }
303
304 Ok(server)
305 }
306}
307
308impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
309 #[allow(clippy::too_many_arguments)]
314 #[allow(clippy::cognitive_complexity)] pub(super) async fn from_executor(
318 config: ServerConfig,
319 executor: Arc<Executor<A>>,
320 subscription_manager: Arc<SubscriptionManager>,
321 #[cfg(feature = "federation")] circuit_breaker: Option<
322 Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>,
323 >,
324 #[cfg(not(feature = "federation"))] _circuit_breaker: Option<()>,
325 error_sanitizer: Arc<crate::config::error_sanitization::ErrorSanitizer>,
326 state_encryption: Option<Arc<crate::auth::state_encryption::StateEncryptionService>>,
327 pkce_store: Option<Arc<crate::auth::PkceStateStore>>,
328 oidc_server_client: Option<Arc<crate::auth::OidcServerClient>>,
329 schema_rate_limiter: Option<Arc<RateLimiter>>,
330 api_key_authenticator: Option<Arc<crate::api_key::ApiKeyAuthenticator>>,
331 revocation_manager: Option<Arc<crate::token_revocation::TokenRevocationManager>>,
332 trusted_docs: Option<Arc<crate::trusted_documents::TrustedDocumentStore>>,
333 #[cfg_attr(
335 not(any(feature = "observers", feature = "auth")),
336 allow(unused_variables)
337 )]
338 db_pool: Option<sqlx::PgPool>,
339 mut tasks: tokio::task::JoinSet<()>,
340 ) -> Result<Self> {
341 let oidc_validator = if let Some(ref auth_config) = config.auth {
343 info!(
344 issuer = %auth_config.issuer,
345 "Initializing OIDC authentication"
346 );
347 let validator = OidcValidator::new(auth_config.clone())
348 .await
349 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
350 Some(Arc::new(validator))
351 } else {
352 None
353 };
354
355 let hs256_auth = build_hs256_auth(&config)?;
357
358 let rate_limiter = if let Some(rl) = schema_rate_limiter {
360 Some(rl)
361 } else if let Some(ref rate_config) = config.rate_limiting {
362 if rate_config.enabled {
363 info!(
364 rps_per_ip = rate_config.rps_per_ip,
365 rps_per_user = rate_config.rps_per_user,
366 "Initializing rate limiting from server config"
367 );
368 Some(Arc::new(RateLimiter::new(rate_config.clone())))
369 } else {
370 info!("Rate limiting disabled by configuration");
371 None
372 }
373 } else {
374 None
375 };
376
377 #[cfg(feature = "observers")]
379 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
380
381 #[cfg(feature = "arrow")]
383 let flight_service = {
384 let mut service = FraiseQLFlightService::new();
385 if let Some(ref validator) = oidc_validator {
386 info!("Enabling OIDC authentication for Arrow Flight");
387 service.set_oidc_validator(validator.clone());
388 } else {
389 info!("Arrow Flight initialized without authentication (dev mode)");
390 }
391 Some(service)
392 };
393
394 #[cfg(feature = "auth")]
396 if pkce_store.is_some() && oidc_server_client.is_none() {
397 tracing::error!(
398 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
399 Auth routes (/auth/start, /auth/callback) will NOT be mounted. \
400 Add [auth] with discovery_url, client_id, client_secret_env, and \
401 server_redirect_uri to fraiseql.toml and recompile the schema."
402 );
403 }
404
405 #[cfg(feature = "auth")]
407 Self::check_redis_requirement(pkce_store.as_ref())?;
408
409 #[cfg(feature = "auth")]
411 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
412
413 #[cfg(not(feature = "auth"))]
416 let _ = (state_encryption, pkce_store, oidc_server_client);
417 Ok(Self {
418 config,
419 executor,
420 subscription_manager,
421 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
422 max_subscriptions_per_connection: None,
423 oidc_validator,
424 hs256_auth,
425 rate_limiter,
426 #[cfg(feature = "secrets")]
427 secrets_manager: None,
428 #[cfg(feature = "federation")]
429 circuit_breaker,
430 error_sanitizer,
431 #[cfg(feature = "auth")]
432 state_encryption,
433 #[cfg(feature = "auth")]
434 pkce_store,
435 #[cfg(feature = "auth")]
436 oidc_server_client,
437 #[cfg(feature = "auth")]
438 social_login: None,
439 #[cfg(feature = "auth")]
440 mfa_state: None,
441 #[cfg(feature = "auth")]
442 anon_signup_state: None,
443 api_key_authenticator,
444 revocation_manager,
445 apq_store: None,
446 trusted_docs,
447 #[cfg(feature = "observers")]
448 observer_runtime,
449 #[cfg(feature = "auth")]
450 enrichment_pool: db_pool.clone(),
451 #[cfg(feature = "observers")]
452 db_pool,
453 storage_state: None,
454 realtime_state: None,
455 tenant_executor_factory: None,
456 #[cfg(feature = "arrow")]
457 flight_service,
458 #[cfg(feature = "mcp")]
459 mcp_config: None,
460 pool_tuning_config: None,
461 adapter_cache_enabled: false,
462 broadcast_manager: None,
463 presence_manager: None,
464 storage_backend: None,
465 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
467 function_store: None,
468 #[cfg(feature = "functions")]
469 function_runtime: None,
470 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
471 tasks,
472 })
473 }
474
475 #[cfg(feature = "auth")]
480 pub(super) fn spawn_pkce_cleanup(
481 pkce_store: Option<&Arc<crate::auth::PkceStateStore>>,
482 tasks: &mut tokio::task::JoinSet<()>,
483 ) {
484 use std::time::Duration;
485
486 use tokio::time::MissedTickBehavior;
487
488 if let Some(store) = pkce_store {
489 let store_clone = Arc::clone(store);
490 tasks.spawn(async move {
491 let mut ticker = tokio::time::interval(Duration::from_secs(300));
492 ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
493 loop {
494 ticker.tick().await;
495 store_clone.cleanup_expired().await;
496 }
497 });
498 }
499 }
500
501 #[must_use]
503 pub fn with_subscription_lifecycle(
504 mut self,
505 lifecycle: Arc<dyn crate::subscriptions::SubscriptionLifecycle>,
506 ) -> Self {
507 self.subscription_lifecycle = lifecycle;
508 self
509 }
510
511 #[must_use]
513 pub const fn with_max_subscriptions_per_connection(mut self, max: u32) -> Self {
514 self.max_subscriptions_per_connection = Some(max);
515 self
516 }
517
518 #[must_use]
528 pub fn with_tenant_executor_factory(
529 mut self,
530 factory: crate::tenancy::TenantExecutorFactory<A>,
531 ) -> Self {
532 self.tenant_executor_factory = Some(factory);
533 self
534 }
535
536 #[must_use]
542 pub fn with_realtime(mut self, state: crate::realtime::server::RealtimeState) -> Self {
543 self.realtime_state = Some(state);
544 self
545 }
546
547 #[must_use]
549 pub fn with_broadcast(mut self, config: crate::subscriptions::BroadcastConfig) -> Self {
550 self.broadcast_manager =
551 Some(Arc::new(crate::subscriptions::BroadcastManager::new(config)));
552 self
553 }
554
555 #[must_use]
557 pub fn with_presence(mut self, config: crate::subscriptions::PresenceConfig) -> Self {
558 self.presence_manager = Some(Arc::new(crate::subscriptions::PresenceManager::new(config)));
559 self
560 }
561
562 pub fn with_pool_tuning(
571 mut self,
572 config: crate::config::pool_tuning::PoolPressureMonitorConfig,
573 ) -> std::result::Result<Self, String> {
574 config.validate()?;
575 self.pool_tuning_config = Some(config);
576 Ok(self)
577 }
578
579 #[cfg(feature = "auth")]
593 #[must_use]
594 pub fn with_social_login(
595 mut self,
596 social_login: Arc<crate::auth::social::SocialLoginState>,
597 ) -> Self {
598 self.social_login = Some(social_login);
599 self
600 }
601
602 #[cfg(feature = "auth")]
608 #[must_use]
609 pub fn with_anon_signup(mut self, state: Arc<crate::auth::AnonSignupState>) -> Self {
610 self.anon_signup_state = Some(state);
611 self
612 }
613
614 #[cfg(feature = "auth")]
623 #[must_use]
624 pub fn with_mfa(mut self, mfa_state: Arc<crate::auth::MfaRouteState>) -> Self {
625 self.mfa_state = Some(mfa_state);
626 self
627 }
628
629 #[must_use]
638 pub fn with_storage(mut self, backend: Arc<dyn crate::storage::StorageBackend>) -> Self {
639 self.storage_backend = Some(backend);
640 self
641 }
642
643 #[must_use]
648 pub const fn with_storage_max_upload_bytes(mut self, bytes: usize) -> Self {
649 self.storage_max_upload_bytes = bytes;
650 self
651 }
652
653 #[must_use]
664 pub fn with_storage_state(mut self, state: fraiseql_storage::StorageState) -> Self {
665 self.storage_state = Some(state);
666 self
667 }
668
669 #[must_use]
677 pub fn with_revocation_manager(
678 mut self,
679 manager: Arc<crate::token_revocation::TokenRevocationManager>,
680 ) -> Self {
681 self.revocation_manager = Some(manager);
682 self
683 }
684
685 #[cfg(feature = "functions")]
691 #[must_use]
692 pub fn with_functions(
693 mut self,
694 store: Arc<dyn fraiseql_functions::FunctionStore>,
695 runtime: Arc<dyn fraiseql_functions::runtime::SendFunctionRuntime>,
696 ) -> Self {
697 self.function_store = Some(store);
698 self.function_runtime = Some(runtime);
699 self
700 }
701
702 #[cfg(feature = "secrets")]
706 pub fn set_secrets_manager(&mut self, manager: Arc<crate::secrets_manager::SecretsManager>) {
707 self.secrets_manager = Some(manager);
708 info!("Secrets manager attached to server");
709 }
710
711 #[cfg(feature = "mcp")]
721 pub async fn serve_mcp_stdio(self) -> Result<()> {
722 use rmcp::ServiceExt;
723
724 let mcp_cfg = self.mcp_config.ok_or_else(|| {
725 ServerError::ConfigError(
726 "FRAISEQL_MCP_STDIO=1 but MCP is not configured. \
727 Add [mcp] enabled = true to fraiseql.toml and recompile the schema."
728 .into(),
729 )
730 })?;
731
732 let schema = Arc::new(self.executor.schema().clone());
733 let executor = self.executor.clone();
734
735 let service = crate::mcp::handler::FraiseQLMcpService::new(schema, executor, mcp_cfg)
736 .with_oidc_validator(self.oidc_validator.clone());
737
738 info!("MCP stdio transport starting — reading from stdin, writing to stdout");
739
740 let running = service
741 .serve((tokio::io::stdin(), tokio::io::stdout()))
742 .await
743 .map_err(|e| ServerError::ConfigError(format!("MCP stdio init failed: {e}")))?;
744
745 running
746 .waiting()
747 .await
748 .map_err(|e| ServerError::ConfigError(format!("MCP stdio error: {e}")))?;
749
750 Ok(())
751 }
752}
753
754pub fn page_size_precedence(env: Option<&str>, compiled: Option<u32>) -> Option<u32> {
761 if let Some(raw) = env {
762 let trimmed = raw.trim();
763 if trimmed.eq_ignore_ascii_case("none") || trimmed == "0" {
764 return None;
765 }
766 if let Ok(n) = trimmed.parse::<u32>() {
767 return Some(n);
768 }
769 }
771 compiled.or(RuntimeConfig::default().max_page_size)
772}