1use std::sync::Arc;
5
6#[cfg(feature = "arrow")]
7use fraiseql_arrow::FraiseQLFlightService;
8#[cfg(all(feature = "arrow", feature = "auth"))]
9use fraiseql_core::security::OidcValidator;
10use fraiseql_core::{
11 cache::{CacheConfig, CachedDatabaseAdapter, QueryResultCache},
12 db::traits::{DatabaseAdapter, RelayDatabaseAdapter},
13 runtime::{Executor, RuntimeConfig, SubscriptionManager},
14 schema::CompiledSchema,
15};
16#[cfg(feature = "observers")]
17use tokio::sync::RwLock;
18use tracing::info;
19#[cfg(feature = "observers")]
20use tracing::warn;
21
22#[cfg(feature = "arrow")]
23use super::RateLimiter;
24#[cfg(all(feature = "arrow", feature = "auth"))]
25use super::ServerError;
26#[cfg(feature = "observers")]
27use super::{ObserverRuntime, ObserverRuntimeConfig};
28use super::{Result, Server, ServerConfig};
29
30impl<A: DatabaseAdapter + RelayDatabaseAdapter + Clone + Send + Sync + 'static>
31 Server<CachedDatabaseAdapter<A>>
32{
33 pub async fn with_relay_pagination(
68 config: ServerConfig,
69 schema: CompiledSchema,
70 adapter: Arc<A>,
71 db_pool: Option<sqlx::PgPool>,
72 ) -> Result<Self> {
73 if config.cache_enabled && !schema.has_rls_configured() {
75 if schema.is_multi_tenant() {
76 return Err(super::ServerError::ConfigError(
77 "Cache is enabled in a multi-tenant schema but no Row-Level Security \
78 policies are declared. This would allow cross-tenant cache hits and \
79 data leakage. In fraiseql.toml, either disable caching with \
80 [cache] enabled = false, declare [security.rls] policies, or set \
81 [security] multi_tenant = false to acknowledge single-tenant mode."
82 .to_string(),
83 ));
84 }
85 tracing::warn!(
86 "Query-result caching is enabled but no Row-Level Security policies are \
87 declared in the compiled schema. This is safe for single-tenant deployments."
88 );
89 }
90
91 crate::server::initialization::field_encryption_unsupported_check(&schema)?;
95 let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
98 super::ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
99 })?;
100
101 #[cfg(feature = "federation")]
103 let circuit_breaker = schema.federation.as_ref().and_then(
104 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
105 );
106 #[cfg(not(feature = "federation"))]
107 let circuit_breaker: Option<()> = None;
108 #[cfg(not(feature = "federation"))]
109 let _ = &schema.federation;
110 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
111 #[cfg(feature = "auth")]
112 let state_encryption = Self::state_encryption_from_schema(&schema)?;
113 #[cfg(not(feature = "auth"))]
114 let state_encryption: Option<
115 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
116 > = None;
117 #[cfg(feature = "auth")]
118 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
119 #[cfg(not(feature = "auth"))]
120 let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
121 #[cfg(feature = "auth")]
122 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
123 #[cfg(not(feature = "auth"))]
124 let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
125 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
126 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
127 let service_account_authenticator =
128 crate::service_account::service_account_authenticator_from_schema(&schema);
129 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
130 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
131 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
132
133 let cache_config = CacheConfig::from(config.cache_enabled);
134 let cache = QueryResultCache::new(cache_config);
135 let inner = Arc::into_inner(adapter)
137 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
138 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
139 .with_ttl_overrides_from_schema(&schema);
140 let executor = Arc::new(Executor::with_config_and_relay(
141 schema.clone(),
142 Arc::new(cached),
143 executor_config,
144 ));
145 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
146
147 let mut server = Self::from_executor(
148 config,
149 executor,
150 subscription_manager,
151 circuit_breaker,
152 error_sanitizer,
153 state_encryption,
154 pkce_store,
155 oidc_server_client,
156 schema_rate_limiter,
157 api_key_authenticator,
158 service_account_authenticator,
159 revocation_manager,
160 trusted_docs,
161 db_pool,
162 tasks,
163 )
164 .await?;
165
166 server.adapter_cache_enabled = cache_config.enabled;
167
168 #[cfg(feature = "mcp")]
170 if let Some(ref cfg) = server.executor.schema().mcp_config {
171 if cfg.enabled {
172 let tool_count =
173 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
174 info!(
175 path = %cfg.path,
176 transport = %cfg.transport,
177 tools = tool_count,
178 "MCP server configured"
179 );
180 server.mcp_config = Some(cfg.clone());
181 }
182 }
183
184 if server.config.apq_enabled {
186 let apq_store: fraiseql_core::apq::ArcApqStorage =
187 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
188 server.apq_store = Some(apq_store);
189 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
190 }
191
192 Ok(server)
193 }
194}
195
196impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
197 #[cfg(feature = "arrow")]
213 pub async fn with_flight_service(
214 config: ServerConfig,
215 schema: CompiledSchema,
216 adapter: Arc<A>,
217 #[allow(unused_variables)]
218 db_pool: Option<sqlx::PgPool>,
220 flight_service: Option<FraiseQLFlightService>,
221 ) -> Result<Self> {
222 crate::server::initialization::field_encryption_unsupported_check(&schema)?;
225 let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
228 super::ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
229 })?;
230
231 #[cfg(feature = "federation")]
233 let circuit_breaker = schema.federation.as_ref().and_then(
234 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
235 );
236 #[cfg(not(feature = "federation"))]
240 let _ = &schema.federation;
241 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
242 #[cfg(feature = "auth")]
243 let state_encryption = Self::state_encryption_from_schema(&schema)?;
244 #[cfg(not(feature = "auth"))]
245 let _state_encryption: Option<
246 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
247 > = None;
248 #[cfg(feature = "auth")]
249 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
250 #[cfg(not(feature = "auth"))]
251 let _pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
252 #[cfg(feature = "auth")]
253 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
254 #[cfg(not(feature = "auth"))]
255 let _oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
256 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
257 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
258 let service_account_authenticator =
259 crate::service_account::service_account_authenticator_from_schema(&schema);
260 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
261 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
262 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
263
264 let executor = Arc::new(Executor::with_config(schema.clone(), adapter, executor_config));
265 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
266
267 #[cfg(feature = "auth")]
269 let oidc_validator = if let Some(ref auth_config) = config.auth {
270 info!(
271 issuer = %auth_config.issuer,
272 "Initializing OIDC authentication"
273 );
274 let validator = OidcValidator::new(auth_config.clone())
275 .await
276 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
277 Some(Arc::new(validator))
278 } else {
279 None
280 };
281 #[cfg(not(feature = "auth"))]
282 let oidc_validator: Option<Arc<fraiseql_core::security::OidcValidator>> = None;
283
284 let hs256_auth = super::builder::build_hs256_auth(&config)?;
286
287 let rate_limiter = if let Some(rl) = schema_rate_limiter {
289 Some(rl)
290 } else if let Some(ref rate_config) = config.rate_limiting {
291 if rate_config.enabled {
292 info!(
293 rps_per_ip = rate_config.rps_per_ip,
294 rps_per_user = rate_config.rps_per_user,
295 "Initializing rate limiting from server config"
296 );
297 Some(Arc::new(RateLimiter::new(rate_config.clone())))
298 } else {
299 info!("Rate limiting disabled by configuration");
300 None
301 }
302 } else {
303 None
304 };
305
306 #[cfg(feature = "observers")]
308 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
309
310 #[cfg(feature = "auth")]
312 if pkce_store.is_some() && oidc_server_client.is_none() {
313 tracing::error!(
314 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
315 Auth routes will NOT be mounted."
316 );
317 }
318
319 #[cfg(feature = "auth")]
321 Self::check_redis_requirement(pkce_store.as_ref())?;
322
323 #[cfg(feature = "auth")]
325 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
326
327 let apq_enabled = config.apq_enabled;
328
329 Ok(Self {
330 config,
331 executor,
332 subscription_manager,
333 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
334 max_subscriptions_per_connection: None,
335 oidc_validator,
336 hs256_auth,
337 rate_limiter,
338 #[cfg(feature = "secrets")]
339 secrets_manager: None,
340 #[cfg(feature = "federation")]
341 circuit_breaker,
342 error_sanitizer,
343 #[cfg(feature = "auth")]
344 state_encryption,
345 #[cfg(feature = "auth")]
346 pkce_store,
347 #[cfg(feature = "auth")]
348 oidc_server_client,
349 #[cfg(feature = "auth")]
350 social_login: None,
351 #[cfg(feature = "auth")]
352 anon_signup_state: None,
353 #[cfg(feature = "auth")]
354 mfa_state: None,
355 api_key_authenticator,
356 service_account_authenticator,
357 revocation_manager,
358 apq_store: if apq_enabled {
359 Some(Arc::new(fraiseql_core::apq::InMemoryApqStorage::default())
360 as fraiseql_core::apq::ArcApqStorage)
361 } else {
362 None
363 },
364 trusted_docs,
365 #[cfg(feature = "mcp")]
366 mcp_config: None,
367 pool_tuning_config: None,
368 #[cfg(feature = "observers")]
369 observer_runtime,
370 #[cfg(feature = "auth")]
371 enrichment_pool: db_pool.clone(),
372 #[cfg(feature = "observers")]
373 db_pool,
374 storage_state: None,
375 realtime_state: None,
376 #[cfg(feature = "functions-runtime")]
377 functions_hooks: None,
378 tenant_executor_factory: None,
379 #[cfg(feature = "arrow")]
380 flight_service,
381 adapter_cache_enabled: false,
382 broadcast_manager: None,
383 presence_manager: None,
384 storage_backend: None,
385 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
387 function_store: None,
388 #[cfg(feature = "functions")]
389 function_runtime: None,
390 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
391 tasks,
392 })
393 }
394
395 #[cfg(feature = "observers")]
405 pub(super) async fn init_observer_runtime(
406 config: &ServerConfig,
407 pool: Option<&sqlx::PgPool>,
408 ) -> crate::Result<Option<Arc<RwLock<ObserverRuntime>>>> {
409 use fraiseql_observers::config::TransportKind;
410
411 let observer_config = match &config.observers {
413 Some(cfg) if cfg.enabled => cfg,
414 _ => {
415 info!("Observer runtime disabled");
416 return Ok(None);
417 },
418 };
419
420 let Some(pool) = pool else {
421 warn!("No database pool provided for observers");
422 return Ok(None);
423 };
424
425 info!("Initializing observer runtime");
426
427 let mut transport = observer_config.runtime.transport.clone().with_env_overrides();
431 let compiled_in = cfg!(feature = "observers-nats");
432 let nats_url_present = !transport.nats.url.is_empty();
433 crate::server::initialization::observer_transport_check(
434 transport.transport,
435 compiled_in,
436 nats_url_present,
437 crate::ServerConfig::is_production_mode(),
438 )?;
439
440 let usable = match transport.transport {
444 TransportKind::Postgres | TransportKind::InMemory => true,
445 TransportKind::Nats => compiled_in && nats_url_present,
446 _ => false,
447 };
448 if !usable {
449 transport.transport = TransportKind::Postgres;
450 }
451
452 transport.validate().map_err(|e| {
453 crate::ServerError::ConfigError(format!("invalid observer transport config: {e}"))
454 })?;
455
456 let runtime_config = ObserverRuntimeConfig::new(pool.clone())
457 .with_poll_interval(observer_config.runtime.poll_interval_ms)
458 .with_batch_size(observer_config.runtime.batch_size)
459 .with_channel_capacity(observer_config.runtime.channel_capacity)
460 .with_max_dlq_size(observer_config.runtime.max_dlq_size)
461 .with_transport(transport)
462 .with_email(observer_config.runtime.email.clone())
463 .with_log_payloads(observer_config.runtime.log_payloads);
464
465 let runtime = ObserverRuntime::new(runtime_config);
466 Ok(Some(Arc::new(RwLock::new(runtime))))
467 }
468}