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 service_account_authenticator =
137 crate::service_account::service_account_authenticator_from_schema(&schema);
138 if service_account_authenticator.is_some() {
139 info!("Service-account authentication enabled");
140 }
141 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
142 if revocation_manager.is_some() {
143 info!("Token revocation enabled");
144 }
145 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
148 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
149
150 if config.cache_enabled && !schema.has_rls_configured() {
155 if schema.is_multi_tenant() {
156 return Err(ServerError::ConfigError(
158 "Cache is enabled in a multi-tenant schema but no Row-Level Security \
159 policies are declared. This would allow cross-tenant cache hits and \
160 data leakage. In fraiseql.toml, either disable caching with \
161 [cache] enabled = false, declare [security.rls] policies, or set \
162 [security] multi_tenant = false to acknowledge single-tenant mode."
163 .to_string(),
164 ));
165 }
166 warn!(
168 "Query-result caching is enabled but no Row-Level Security policies are \
169 declared in the compiled schema. This is safe for single-tenant deployments. \
170 For multi-tenant deployments, declare RLS policies and set \
171 `security.multi_tenant = true` in your schema."
172 );
173 }
174
175 let cache_config = CacheConfig::from(config.cache_enabled);
177 let cache = QueryResultCache::new(cache_config);
178
179 if cache_config.enabled {
181 tracing::info!(
182 max_entries = cache_config.max_entries,
183 ttl_seconds = cache_config.ttl_seconds,
184 rls_enforcement = ?cache_config.rls_enforcement,
185 "Query result cache: active"
186 );
187 } else {
188 tracing::info!("Query result cache: disabled");
189 }
190
191 let subscriptions_config = schema.subscriptions_config.clone();
193
194 let inner = Arc::into_inner(adapter)
196 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
197 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
198 .with_ttl_overrides_from_schema(&schema)
199 .with_rls(schema.has_rls_configured());
200
201 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 service_account_authenticator,
219 revocation_manager,
220 trusted_docs,
221 db_pool,
222 tasks,
223 )
224 .await?;
225
226 server.adapter_cache_enabled = cache_config.enabled;
227
228 if let Some(pt) = server.config.pool_tuning.clone() {
230 if pt.enabled {
231 server = server
232 .with_pool_tuning(pt)
233 .map_err(|e| ServerError::ConfigError(format!("pool_tuning: {e}")))?;
234 }
235 }
236
237 #[cfg(feature = "mcp")]
239 if let Some(ref cfg) = server.executor.schema().mcp_config {
240 if cfg.enabled {
241 let tool_count =
242 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
243 info!(
244 path = %cfg.path,
245 transport = %cfg.transport,
246 tools = tool_count,
247 "MCP server configured"
248 );
249 server.mcp_config = Some(cfg.clone());
250 }
251 }
252
253 if server.config.apq_enabled {
255 let apq_store: fraiseql_core::apq::ArcApqStorage =
256 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
257 server.apq_store = Some(apq_store);
258 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
259 }
260
261 if let Some(ref subs) = subscriptions_config {
263 if let Some(max) = subs.max_subscriptions_per_connection {
264 server.max_subscriptions_per_connection = Some(max);
265 }
266 if let Some(lifecycle) = crate::subscriptions::WebhookLifecycle::from_config(subs) {
267 server.subscription_lifecycle = Arc::new(lifecycle);
268 }
269 }
270
271 Ok(server)
272 }
273}
274
275impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
276 #[allow(clippy::too_many_arguments)]
281 #[allow(clippy::cognitive_complexity)] pub(super) async fn from_executor(
285 config: ServerConfig,
286 executor: Arc<Executor<A>>,
287 subscription_manager: Arc<SubscriptionManager>,
288 #[cfg(feature = "federation")] circuit_breaker: Option<
289 Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>,
290 >,
291 #[cfg(not(feature = "federation"))] _circuit_breaker: Option<()>,
292 error_sanitizer: Arc<crate::config::error_sanitization::ErrorSanitizer>,
293 state_encryption: Option<Arc<crate::auth::state_encryption::StateEncryptionService>>,
294 pkce_store: Option<Arc<crate::auth::PkceStateStore>>,
295 oidc_server_client: Option<Arc<crate::auth::OidcServerClient>>,
296 schema_rate_limiter: Option<Arc<RateLimiter>>,
297 api_key_authenticator: Option<Arc<crate::api_key::ApiKeyAuthenticator>>,
298 service_account_authenticator: Option<
299 Arc<crate::service_account::ServiceAccountAuthenticator>,
300 >,
301 revocation_manager: Option<Arc<crate::token_revocation::TokenRevocationManager>>,
302 trusted_docs: Option<Arc<crate::trusted_documents::TrustedDocumentStore>>,
303 #[cfg_attr(
305 not(any(feature = "observers", feature = "auth")),
306 allow(unused_variables)
307 )]
308 db_pool: Option<sqlx::PgPool>,
309 mut tasks: tokio::task::JoinSet<()>,
310 ) -> Result<Self> {
311 let oidc_validator = if let Some(ref auth_config) = config.auth {
313 info!(
314 issuer = %auth_config.issuer,
315 "Initializing OIDC authentication"
316 );
317 let validator = OidcValidator::new(auth_config.clone())
318 .await
319 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
320 Some(Arc::new(validator))
321 } else {
322 None
323 };
324
325 let hs256_auth = build_hs256_auth(&config)?;
327
328 let rate_limiter = if let Some(rl) = schema_rate_limiter {
330 Some(rl)
331 } else if let Some(ref rate_config) = config.rate_limiting {
332 if rate_config.enabled {
333 info!(
334 rps_per_ip = rate_config.rps_per_ip,
335 rps_per_user = rate_config.rps_per_user,
336 "Initializing rate limiting from server config"
337 );
338 Some(Arc::new(RateLimiter::new(rate_config.clone())))
339 } else {
340 info!("Rate limiting disabled by configuration");
341 None
342 }
343 } else {
344 None
345 };
346
347 #[cfg(feature = "observers")]
349 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
350
351 #[cfg(feature = "arrow")]
353 let flight_service = {
354 let mut service = FraiseQLFlightService::new();
355 if let Some(ref validator) = oidc_validator {
356 info!("Enabling OIDC authentication for Arrow Flight");
357 service.set_oidc_validator(validator.clone());
358 } else {
359 info!("Arrow Flight initialized without authentication (dev mode)");
360 }
361 Some(service)
362 };
363
364 #[cfg(feature = "auth")]
366 if pkce_store.is_some() && oidc_server_client.is_none() {
367 tracing::error!(
368 "pkce.enabled = true but no OIDC client is available. Auth routes \
369 (/auth/start, /auth/callback) will NOT be mounted. Building an \
370 OidcServerClient from the compiled schema's [auth] block is not yet \
371 functional (the compiled schema carries no auth/auth_endpoints) — \
372 tracked in #621."
373 );
374 }
375
376 #[cfg(feature = "auth")]
378 Self::check_redis_requirement(pkce_store.as_ref())?;
379
380 #[cfg(feature = "auth")]
382 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
383
384 #[cfg(not(feature = "auth"))]
387 let _ = (state_encryption, pkce_store, oidc_server_client);
388 Ok(Self {
389 config,
390 executor,
391 subscription_manager,
392 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
393 max_subscriptions_per_connection: None,
394 oidc_validator,
395 hs256_auth,
396 rate_limiter,
397 #[cfg(feature = "secrets")]
398 secrets_manager: None,
399 #[cfg(feature = "federation")]
400 circuit_breaker,
401 error_sanitizer,
402 #[cfg(feature = "auth")]
403 state_encryption,
404 #[cfg(feature = "auth")]
405 pkce_store,
406 #[cfg(feature = "auth")]
407 oidc_server_client,
408 #[cfg(feature = "auth")]
409 social_login: None,
410 #[cfg(feature = "auth")]
411 mfa_state: None,
412 #[cfg(feature = "auth")]
413 anon_signup_state: None,
414 api_key_authenticator,
415 service_account_authenticator,
416 revocation_manager,
417 apq_store: None,
418 trusted_docs,
419 #[cfg(feature = "observers")]
420 observer_runtime,
421 #[cfg(feature = "auth")]
422 enrichment_pool: db_pool.clone(),
423 #[cfg(feature = "observers")]
424 db_pool,
425 storage_state: None,
426 #[cfg(feature = "functions-runtime")]
427 functions_hooks: None,
428 tenant_executor_factory: None,
429 #[cfg(feature = "arrow")]
430 flight_service,
431 #[cfg(feature = "mcp")]
432 mcp_config: None,
433 pool_tuning_config: None,
434 adapter_cache_enabled: false,
435 storage_backend: None,
436 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
438 function_store: None,
439 #[cfg(feature = "functions")]
440 function_runtime: None,
441 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
442 tasks,
443 })
444 }
445
446 #[cfg(feature = "auth")]
451 pub(super) fn spawn_pkce_cleanup(
452 pkce_store: Option<&Arc<crate::auth::PkceStateStore>>,
453 tasks: &mut tokio::task::JoinSet<()>,
454 ) {
455 use std::time::Duration;
456
457 use tokio::time::MissedTickBehavior;
458
459 if let Some(store) = pkce_store {
460 let store_clone = Arc::clone(store);
461 tasks.spawn(async move {
462 let mut ticker = tokio::time::interval(Duration::from_secs(300));
463 ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
464 loop {
465 ticker.tick().await;
466 store_clone.cleanup_expired().await;
467 }
468 });
469 }
470 }
471
472 #[must_use]
474 pub fn with_subscription_lifecycle(
475 mut self,
476 lifecycle: Arc<dyn crate::subscriptions::SubscriptionLifecycle>,
477 ) -> Self {
478 self.subscription_lifecycle = lifecycle;
479 self
480 }
481
482 #[must_use]
484 pub const fn with_max_subscriptions_per_connection(mut self, max: u32) -> Self {
485 self.max_subscriptions_per_connection = Some(max);
486 self
487 }
488
489 #[must_use]
499 pub fn with_tenant_executor_factory(
500 mut self,
501 factory: crate::tenancy::TenantExecutorFactory<A>,
502 ) -> Self {
503 self.tenant_executor_factory = Some(factory);
504 self
505 }
506
507 pub fn with_pool_tuning(
516 mut self,
517 config: crate::config::pool_tuning::PoolPressureMonitorConfig,
518 ) -> std::result::Result<Self, String> {
519 config.validate()?;
520 self.pool_tuning_config = Some(config);
521 Ok(self)
522 }
523
524 #[cfg(feature = "auth")]
538 #[must_use]
539 pub fn with_social_login(
540 mut self,
541 social_login: Arc<crate::auth::social::SocialLoginState>,
542 ) -> Self {
543 self.social_login = Some(social_login);
544 self
545 }
546
547 #[cfg(feature = "auth")]
553 #[must_use]
554 pub fn with_anon_signup(mut self, state: Arc<crate::auth::AnonSignupState>) -> Self {
555 self.anon_signup_state = Some(state);
556 self
557 }
558
559 #[cfg(feature = "auth")]
568 #[must_use]
569 pub fn with_mfa(mut self, mfa_state: Arc<crate::auth::MfaRouteState>) -> Self {
570 self.mfa_state = Some(mfa_state);
571 self
572 }
573
574 #[must_use]
583 pub fn with_storage(mut self, backend: Arc<dyn crate::storage::StorageBackend>) -> Self {
584 self.storage_backend = Some(backend);
585 self
586 }
587
588 #[must_use]
593 pub const fn with_storage_max_upload_bytes(mut self, bytes: usize) -> Self {
594 self.storage_max_upload_bytes = bytes;
595 self
596 }
597
598 #[must_use]
609 pub fn with_storage_state(mut self, state: fraiseql_storage::StorageState) -> Self {
610 self.storage_state = Some(state);
611 self
612 }
613
614 #[must_use]
622 pub fn with_revocation_manager(
623 mut self,
624 manager: Arc<crate::token_revocation::TokenRevocationManager>,
625 ) -> Self {
626 self.revocation_manager = Some(manager);
627 self
628 }
629
630 #[cfg(feature = "functions")]
636 #[must_use]
637 pub fn with_functions(
638 mut self,
639 store: Arc<dyn fraiseql_functions::FunctionStore>,
640 runtime: Arc<dyn fraiseql_functions::runtime::SendFunctionRuntime>,
641 ) -> Self {
642 self.function_store = Some(store);
643 self.function_runtime = Some(runtime);
644 self
645 }
646
647 #[cfg(feature = "secrets")]
651 pub fn set_secrets_manager(&mut self, manager: Arc<crate::secrets_manager::SecretsManager>) {
652 self.secrets_manager = Some(manager);
653 info!("Secrets manager attached to server");
654 }
655
656 #[cfg(feature = "mcp")]
666 pub async fn serve_mcp_stdio(self) -> Result<()> {
667 use rmcp::ServiceExt;
668
669 let mcp_cfg = self.mcp_config.ok_or_else(|| {
670 ServerError::ConfigError(
671 "FRAISEQL_MCP_STDIO=1 but MCP is not configured. \
672 Add [mcp] enabled = true to fraiseql.toml and recompile the schema."
673 .into(),
674 )
675 })?;
676
677 let schema = Arc::new(self.executor.schema().clone());
678 let executor = self.executor.clone();
679
680 let service = crate::mcp::handler::FraiseQLMcpService::new(schema, executor, mcp_cfg)
681 .with_oidc_validator(self.oidc_validator.clone());
682
683 info!("MCP stdio transport starting — reading from stdin, writing to stdout");
684
685 let running = service
686 .serve((tokio::io::stdin(), tokio::io::stdout()))
687 .await
688 .map_err(|e| ServerError::ConfigError(format!("MCP stdio init failed: {e}")))?;
689
690 running
691 .waiting()
692 .await
693 .map_err(|e| ServerError::ConfigError(format!("MCP stdio error: {e}")))?;
694
695 Ok(())
696 }
697}