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 #[cfg(feature = "federation")]
106 let circuit_breaker = schema.federation.as_ref().and_then(
107 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
108 );
109 #[cfg(not(feature = "federation"))]
110 let circuit_breaker: Option<()> = None;
111 #[cfg(not(feature = "federation"))]
112 let _ = &schema.federation;
113 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
114 #[cfg(feature = "auth")]
115 let state_encryption = Self::state_encryption_from_schema(&schema)?;
116 #[cfg(not(feature = "auth"))]
117 let state_encryption: Option<
118 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
119 > = None;
120 #[cfg(feature = "auth")]
121 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
122 #[cfg(not(feature = "auth"))]
123 let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
124 #[cfg(feature = "auth")]
125 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
126 #[cfg(not(feature = "auth"))]
127 let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
128 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
129 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
130 if api_key_authenticator.is_some() {
131 info!("API key authentication enabled");
132 }
133 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
134 if revocation_manager.is_some() {
135 info!("Token revocation enabled");
136 }
137 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
140 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
141
142 if config.cache_enabled && !schema.has_rls_configured() {
147 if schema.is_multi_tenant() {
148 return Err(ServerError::ConfigError(
150 "Cache is enabled in a multi-tenant schema but no Row-Level Security \
151 policies are declared. This would allow cross-tenant cache hits and \
152 data leakage. In fraiseql.toml, either disable caching with \
153 [cache] enabled = false, declare [security.rls] policies, or set \
154 [security] multi_tenant = false to acknowledge single-tenant mode."
155 .to_string(),
156 ));
157 }
158 warn!(
160 "Query-result caching is enabled but no Row-Level Security policies are \
161 declared in the compiled schema. This is safe for single-tenant deployments. \
162 For multi-tenant deployments, declare RLS policies and set \
163 `security.multi_tenant = true` in your schema."
164 );
165 }
166
167 let cache_config = CacheConfig::from(config.cache_enabled);
169 let cache = QueryResultCache::new(cache_config);
170
171 if cache_config.enabled {
173 tracing::info!(
174 max_entries = cache_config.max_entries,
175 ttl_seconds = cache_config.ttl_seconds,
176 rls_enforcement = ?cache_config.rls_enforcement,
177 "Query result cache: active"
178 );
179 } else {
180 tracing::info!("Query result cache: disabled");
181 }
182
183 let subscriptions_config = schema.subscriptions_config.clone();
185
186 let inner = Arc::into_inner(adapter)
188 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
189 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
190 .with_ttl_overrides_from_schema(&schema)
191 .with_rls(schema.has_rls_configured());
192
193 let audit_mutations = schema
196 .security
197 .as_ref()
198 .and_then(|s| s.additional.get("enterprise"))
199 .and_then(|e| e.get("audit_logging_enabled"))
200 .and_then(|v| v.as_bool())
201 .unwrap_or(false);
202 if audit_mutations {
203 info!("Mutation audit logging enabled (target: fraiseql::mutation_audit)");
204 }
205 let max_page_size = page_size_precedence(
208 std::env::var("FRAISEQL_MAX_PAGE_SIZE").ok().as_deref(),
209 schema.validation_config.as_ref().and_then(|v| v.max_page_size),
210 );
211 let changelog_enabled = std::env::var("FRAISEQL_CHANGELOG_ENABLED")
215 .ok()
216 .map(|v| {
217 !matches!(v.trim().to_ascii_lowercase().as_str(), "false" | "0" | "no" | "off")
218 })
219 .or_else(|| schema.changelog.as_ref().map(|c| c.write_enabled))
220 .unwrap_or(true);
221 if !changelog_enabled {
222 info!(
223 "Change-log outbox write disabled (FRAISEQL_CHANGELOG_ENABLED / [changelog] write_enabled)"
224 );
225 }
226 let executor_config = RuntimeConfig {
227 audit_mutations,
228 max_page_size,
229 changelog_enabled,
230 ..RuntimeConfig::default()
231 };
232 let executor =
233 Arc::new(Executor::with_config(schema.clone(), Arc::new(cached), executor_config));
234 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
235
236 let mut server = Self::from_executor(
237 config,
238 executor,
239 subscription_manager,
240 circuit_breaker,
241 error_sanitizer,
242 state_encryption,
243 pkce_store,
244 oidc_server_client,
245 schema_rate_limiter,
246 api_key_authenticator,
247 revocation_manager,
248 trusted_docs,
249 db_pool,
250 tasks,
251 )
252 .await?;
253
254 server.adapter_cache_enabled = cache_config.enabled;
255
256 if let Some(pt) = server.config.pool_tuning.clone() {
258 if pt.enabled {
259 server = server
260 .with_pool_tuning(pt)
261 .map_err(|e| ServerError::ConfigError(format!("pool_tuning: {e}")))?;
262 }
263 }
264
265 #[cfg(feature = "mcp")]
267 if let Some(ref cfg) = server.executor.schema().mcp_config {
268 if cfg.enabled {
269 let tool_count =
270 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
271 info!(
272 path = %cfg.path,
273 transport = %cfg.transport,
274 tools = tool_count,
275 "MCP server configured"
276 );
277 server.mcp_config = Some(cfg.clone());
278 }
279 }
280
281 if server.config.apq_enabled {
283 let apq_store: fraiseql_core::apq::ArcApqStorage =
284 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
285 server.apq_store = Some(apq_store);
286 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
287 }
288
289 if let Some(ref subs) = subscriptions_config {
291 if let Some(max) = subs.max_subscriptions_per_connection {
292 server.max_subscriptions_per_connection = Some(max);
293 }
294 if let Some(lifecycle) = crate::subscriptions::WebhookLifecycle::from_config(subs) {
295 server.subscription_lifecycle = Arc::new(lifecycle);
296 }
297 }
298
299 Ok(server)
300 }
301}
302
303impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
304 #[allow(clippy::too_many_arguments)]
309 #[allow(clippy::cognitive_complexity)] pub(super) async fn from_executor(
313 config: ServerConfig,
314 executor: Arc<Executor<A>>,
315 subscription_manager: Arc<SubscriptionManager>,
316 #[cfg(feature = "federation")] circuit_breaker: Option<
317 Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>,
318 >,
319 #[cfg(not(feature = "federation"))] _circuit_breaker: Option<()>,
320 error_sanitizer: Arc<crate::config::error_sanitization::ErrorSanitizer>,
321 state_encryption: Option<Arc<crate::auth::state_encryption::StateEncryptionService>>,
322 pkce_store: Option<Arc<crate::auth::PkceStateStore>>,
323 oidc_server_client: Option<Arc<crate::auth::OidcServerClient>>,
324 schema_rate_limiter: Option<Arc<RateLimiter>>,
325 api_key_authenticator: Option<Arc<crate::api_key::ApiKeyAuthenticator>>,
326 revocation_manager: Option<Arc<crate::token_revocation::TokenRevocationManager>>,
327 trusted_docs: Option<Arc<crate::trusted_documents::TrustedDocumentStore>>,
328 #[cfg_attr(
330 not(any(feature = "observers", feature = "auth")),
331 allow(unused_variables)
332 )]
333 db_pool: Option<sqlx::PgPool>,
334 mut tasks: tokio::task::JoinSet<()>,
335 ) -> Result<Self> {
336 let oidc_validator = if let Some(ref auth_config) = config.auth {
338 info!(
339 issuer = %auth_config.issuer,
340 "Initializing OIDC authentication"
341 );
342 let validator = OidcValidator::new(auth_config.clone())
343 .await
344 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
345 Some(Arc::new(validator))
346 } else {
347 None
348 };
349
350 let hs256_auth = build_hs256_auth(&config)?;
352
353 let rate_limiter = if let Some(rl) = schema_rate_limiter {
355 Some(rl)
356 } else if let Some(ref rate_config) = config.rate_limiting {
357 if rate_config.enabled {
358 info!(
359 rps_per_ip = rate_config.rps_per_ip,
360 rps_per_user = rate_config.rps_per_user,
361 "Initializing rate limiting from server config"
362 );
363 Some(Arc::new(RateLimiter::new(rate_config.clone())))
364 } else {
365 info!("Rate limiting disabled by configuration");
366 None
367 }
368 } else {
369 None
370 };
371
372 #[cfg(feature = "observers")]
374 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
375
376 #[cfg(feature = "arrow")]
378 let flight_service = {
379 let mut service = FraiseQLFlightService::new();
380 if let Some(ref validator) = oidc_validator {
381 info!("Enabling OIDC authentication for Arrow Flight");
382 service.set_oidc_validator(validator.clone());
383 } else {
384 info!("Arrow Flight initialized without authentication (dev mode)");
385 }
386 Some(service)
387 };
388
389 #[cfg(feature = "auth")]
391 if pkce_store.is_some() && oidc_server_client.is_none() {
392 tracing::error!(
393 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
394 Auth routes (/auth/start, /auth/callback) will NOT be mounted. \
395 Add [auth] with discovery_url, client_id, client_secret_env, and \
396 server_redirect_uri to fraiseql.toml and recompile the schema."
397 );
398 }
399
400 #[cfg(feature = "auth")]
402 Self::check_redis_requirement(pkce_store.as_ref())?;
403
404 #[cfg(feature = "auth")]
406 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
407
408 #[cfg(not(feature = "auth"))]
411 let _ = (state_encryption, pkce_store, oidc_server_client);
412 Ok(Self {
413 config,
414 executor,
415 subscription_manager,
416 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
417 max_subscriptions_per_connection: None,
418 oidc_validator,
419 hs256_auth,
420 rate_limiter,
421 #[cfg(feature = "secrets")]
422 secrets_manager: None,
423 #[cfg(feature = "federation")]
424 circuit_breaker,
425 error_sanitizer,
426 #[cfg(feature = "auth")]
427 state_encryption,
428 #[cfg(feature = "auth")]
429 pkce_store,
430 #[cfg(feature = "auth")]
431 oidc_server_client,
432 #[cfg(feature = "auth")]
433 social_login: None,
434 #[cfg(feature = "auth")]
435 mfa_state: None,
436 #[cfg(feature = "auth")]
437 anon_signup_state: None,
438 api_key_authenticator,
439 revocation_manager,
440 apq_store: None,
441 trusted_docs,
442 #[cfg(feature = "observers")]
443 observer_runtime,
444 #[cfg(feature = "auth")]
445 enrichment_pool: db_pool.clone(),
446 #[cfg(feature = "observers")]
447 db_pool,
448 storage_state: None,
449 realtime_state: None,
450 tenant_executor_factory: None,
451 #[cfg(feature = "arrow")]
452 flight_service,
453 #[cfg(feature = "mcp")]
454 mcp_config: None,
455 pool_tuning_config: None,
456 adapter_cache_enabled: false,
457 broadcast_manager: None,
458 presence_manager: None,
459 storage_backend: None,
460 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
462 function_store: None,
463 #[cfg(feature = "functions")]
464 function_runtime: None,
465 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
466 tasks,
467 })
468 }
469
470 #[cfg(feature = "auth")]
475 pub(super) fn spawn_pkce_cleanup(
476 pkce_store: Option<&Arc<crate::auth::PkceStateStore>>,
477 tasks: &mut tokio::task::JoinSet<()>,
478 ) {
479 use std::time::Duration;
480
481 use tokio::time::MissedTickBehavior;
482
483 if let Some(store) = pkce_store {
484 let store_clone = Arc::clone(store);
485 tasks.spawn(async move {
486 let mut ticker = tokio::time::interval(Duration::from_secs(300));
487 ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
488 loop {
489 ticker.tick().await;
490 store_clone.cleanup_expired().await;
491 }
492 });
493 }
494 }
495
496 #[must_use]
498 pub fn with_subscription_lifecycle(
499 mut self,
500 lifecycle: Arc<dyn crate::subscriptions::SubscriptionLifecycle>,
501 ) -> Self {
502 self.subscription_lifecycle = lifecycle;
503 self
504 }
505
506 #[must_use]
508 pub const fn with_max_subscriptions_per_connection(mut self, max: u32) -> Self {
509 self.max_subscriptions_per_connection = Some(max);
510 self
511 }
512
513 #[must_use]
523 pub fn with_tenant_executor_factory(
524 mut self,
525 factory: crate::tenancy::TenantExecutorFactory<A>,
526 ) -> Self {
527 self.tenant_executor_factory = Some(factory);
528 self
529 }
530
531 #[must_use]
537 pub fn with_realtime(mut self, state: crate::realtime::server::RealtimeState) -> Self {
538 self.realtime_state = Some(state);
539 self
540 }
541
542 #[must_use]
544 pub fn with_broadcast(mut self, config: crate::subscriptions::BroadcastConfig) -> Self {
545 self.broadcast_manager =
546 Some(Arc::new(crate::subscriptions::BroadcastManager::new(config)));
547 self
548 }
549
550 #[must_use]
552 pub fn with_presence(mut self, config: crate::subscriptions::PresenceConfig) -> Self {
553 self.presence_manager = Some(Arc::new(crate::subscriptions::PresenceManager::new(config)));
554 self
555 }
556
557 pub fn with_pool_tuning(
566 mut self,
567 config: crate::config::pool_tuning::PoolPressureMonitorConfig,
568 ) -> std::result::Result<Self, String> {
569 config.validate()?;
570 self.pool_tuning_config = Some(config);
571 Ok(self)
572 }
573
574 #[cfg(feature = "auth")]
588 #[must_use]
589 pub fn with_social_login(
590 mut self,
591 social_login: Arc<crate::auth::social::SocialLoginState>,
592 ) -> Self {
593 self.social_login = Some(social_login);
594 self
595 }
596
597 #[cfg(feature = "auth")]
603 #[must_use]
604 pub fn with_anon_signup(mut self, state: Arc<crate::auth::AnonSignupState>) -> Self {
605 self.anon_signup_state = Some(state);
606 self
607 }
608
609 #[cfg(feature = "auth")]
618 #[must_use]
619 pub fn with_mfa(mut self, mfa_state: Arc<crate::auth::MfaRouteState>) -> Self {
620 self.mfa_state = Some(mfa_state);
621 self
622 }
623
624 #[must_use]
633 pub fn with_storage(mut self, backend: Arc<dyn crate::storage::StorageBackend>) -> Self {
634 self.storage_backend = Some(backend);
635 self
636 }
637
638 #[must_use]
643 pub const fn with_storage_max_upload_bytes(mut self, bytes: usize) -> Self {
644 self.storage_max_upload_bytes = bytes;
645 self
646 }
647
648 #[must_use]
659 pub fn with_storage_state(mut self, state: fraiseql_storage::StorageState) -> Self {
660 self.storage_state = Some(state);
661 self
662 }
663
664 #[must_use]
672 pub fn with_revocation_manager(
673 mut self,
674 manager: Arc<crate::token_revocation::TokenRevocationManager>,
675 ) -> Self {
676 self.revocation_manager = Some(manager);
677 self
678 }
679
680 #[cfg(feature = "functions")]
686 #[must_use]
687 pub fn with_functions(
688 mut self,
689 store: Arc<dyn fraiseql_functions::FunctionStore>,
690 runtime: Arc<dyn fraiseql_functions::runtime::SendFunctionRuntime>,
691 ) -> Self {
692 self.function_store = Some(store);
693 self.function_runtime = Some(runtime);
694 self
695 }
696
697 #[cfg(feature = "secrets")]
701 pub fn set_secrets_manager(&mut self, manager: Arc<crate::secrets_manager::SecretsManager>) {
702 self.secrets_manager = Some(manager);
703 info!("Secrets manager attached to server");
704 }
705
706 #[cfg(feature = "mcp")]
716 pub async fn serve_mcp_stdio(self) -> Result<()> {
717 use rmcp::ServiceExt;
718
719 let mcp_cfg = self.mcp_config.ok_or_else(|| {
720 ServerError::ConfigError(
721 "FRAISEQL_MCP_STDIO=1 but MCP is not configured. \
722 Add [mcp] enabled = true to fraiseql.toml and recompile the schema."
723 .into(),
724 )
725 })?;
726
727 let schema = Arc::new(self.executor.schema().clone());
728 let executor = self.executor.clone();
729
730 let service = crate::mcp::handler::FraiseQLMcpService::new(schema, executor, mcp_cfg)
731 .with_oidc_validator(self.oidc_validator.clone());
732
733 info!("MCP stdio transport starting — reading from stdin, writing to stdout");
734
735 let running = service
736 .serve((tokio::io::stdin(), tokio::io::stdout()))
737 .await
738 .map_err(|e| ServerError::ConfigError(format!("MCP stdio init failed: {e}")))?;
739
740 running
741 .waiting()
742 .await
743 .map_err(|e| ServerError::ConfigError(format!("MCP stdio error: {e}")))?;
744
745 Ok(())
746 }
747}
748
749pub fn page_size_precedence(env: Option<&str>, compiled: Option<u32>) -> Option<u32> {
756 if let Some(raw) = env {
757 let trimmed = raw.trim();
758 if trimmed.eq_ignore_ascii_case("none") || trimmed == "0" {
759 return None;
760 }
761 if let Ok(n) = trimmed.parse::<u32>() {
762 return Some(n);
763 }
764 }
766 compiled.or(RuntimeConfig::default().max_page_size)
767}