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 revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
128 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
129 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
130
131 let cache_config = CacheConfig::from(config.cache_enabled);
132 let cache = QueryResultCache::new(cache_config);
133 let inner = Arc::into_inner(adapter)
135 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
136 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
137 .with_ttl_overrides_from_schema(&schema);
138 let executor = Arc::new(Executor::with_config_and_relay(
139 schema.clone(),
140 Arc::new(cached),
141 executor_config,
142 ));
143 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
144
145 let mut server = Self::from_executor(
146 config,
147 executor,
148 subscription_manager,
149 circuit_breaker,
150 error_sanitizer,
151 state_encryption,
152 pkce_store,
153 oidc_server_client,
154 schema_rate_limiter,
155 api_key_authenticator,
156 revocation_manager,
157 trusted_docs,
158 db_pool,
159 tasks,
160 )
161 .await?;
162
163 server.adapter_cache_enabled = cache_config.enabled;
164
165 #[cfg(feature = "mcp")]
167 if let Some(ref cfg) = server.executor.schema().mcp_config {
168 if cfg.enabled {
169 let tool_count =
170 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
171 info!(
172 path = %cfg.path,
173 transport = %cfg.transport,
174 tools = tool_count,
175 "MCP server configured"
176 );
177 server.mcp_config = Some(cfg.clone());
178 }
179 }
180
181 if server.config.apq_enabled {
183 let apq_store: fraiseql_core::apq::ArcApqStorage =
184 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
185 server.apq_store = Some(apq_store);
186 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
187 }
188
189 Ok(server)
190 }
191}
192
193impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
194 #[cfg(feature = "arrow")]
210 pub async fn with_flight_service(
211 config: ServerConfig,
212 schema: CompiledSchema,
213 adapter: Arc<A>,
214 #[allow(unused_variables)]
215 db_pool: Option<sqlx::PgPool>,
217 flight_service: Option<FraiseQLFlightService>,
218 ) -> Result<Self> {
219 crate::server::initialization::field_encryption_unsupported_check(&schema)?;
222 let executor_config = RuntimeConfig::from_compiled_schema(&schema).map_err(|msg| {
225 super::ServerError::ConfigError(format!("Incompatible compiled schema: {msg}"))
226 })?;
227
228 #[cfg(feature = "federation")]
230 let circuit_breaker = schema.federation.as_ref().and_then(
231 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
232 );
233 #[cfg(not(feature = "federation"))]
237 let _ = &schema.federation;
238 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
239 #[cfg(feature = "auth")]
240 let state_encryption = Self::state_encryption_from_schema(&schema)?;
241 #[cfg(not(feature = "auth"))]
242 let _state_encryption: Option<
243 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
244 > = None;
245 #[cfg(feature = "auth")]
246 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
247 #[cfg(not(feature = "auth"))]
248 let _pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
249 #[cfg(feature = "auth")]
250 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
251 #[cfg(not(feature = "auth"))]
252 let _oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
253 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
254 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
255 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
256 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
257 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
258
259 let executor = Arc::new(Executor::with_config(schema.clone(), adapter, executor_config));
260 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
261
262 #[cfg(feature = "auth")]
264 let oidc_validator = if let Some(ref auth_config) = config.auth {
265 info!(
266 issuer = %auth_config.issuer,
267 "Initializing OIDC authentication"
268 );
269 let validator = OidcValidator::new(auth_config.clone())
270 .await
271 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
272 Some(Arc::new(validator))
273 } else {
274 None
275 };
276 #[cfg(not(feature = "auth"))]
277 let oidc_validator: Option<Arc<fraiseql_core::security::OidcValidator>> = None;
278
279 let hs256_auth = super::builder::build_hs256_auth(&config)?;
281
282 let rate_limiter = if let Some(rl) = schema_rate_limiter {
284 Some(rl)
285 } else if let Some(ref rate_config) = config.rate_limiting {
286 if rate_config.enabled {
287 info!(
288 rps_per_ip = rate_config.rps_per_ip,
289 rps_per_user = rate_config.rps_per_user,
290 "Initializing rate limiting from server config"
291 );
292 Some(Arc::new(RateLimiter::new(rate_config.clone())))
293 } else {
294 info!("Rate limiting disabled by configuration");
295 None
296 }
297 } else {
298 None
299 };
300
301 #[cfg(feature = "observers")]
303 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
304
305 #[cfg(feature = "auth")]
307 if pkce_store.is_some() && oidc_server_client.is_none() {
308 tracing::error!(
309 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
310 Auth routes will NOT be mounted."
311 );
312 }
313
314 #[cfg(feature = "auth")]
316 Self::check_redis_requirement(pkce_store.as_ref())?;
317
318 #[cfg(feature = "auth")]
320 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
321
322 let apq_enabled = config.apq_enabled;
323
324 Ok(Self {
325 config,
326 executor,
327 subscription_manager,
328 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
329 max_subscriptions_per_connection: None,
330 oidc_validator,
331 hs256_auth,
332 rate_limiter,
333 #[cfg(feature = "secrets")]
334 secrets_manager: None,
335 #[cfg(feature = "federation")]
336 circuit_breaker,
337 error_sanitizer,
338 #[cfg(feature = "auth")]
339 state_encryption,
340 #[cfg(feature = "auth")]
341 pkce_store,
342 #[cfg(feature = "auth")]
343 oidc_server_client,
344 #[cfg(feature = "auth")]
345 social_login: None,
346 #[cfg(feature = "auth")]
347 anon_signup_state: None,
348 #[cfg(feature = "auth")]
349 mfa_state: None,
350 api_key_authenticator,
351 revocation_manager,
352 apq_store: if apq_enabled {
353 Some(Arc::new(fraiseql_core::apq::InMemoryApqStorage::default())
354 as fraiseql_core::apq::ArcApqStorage)
355 } else {
356 None
357 },
358 trusted_docs,
359 #[cfg(feature = "mcp")]
360 mcp_config: None,
361 pool_tuning_config: None,
362 #[cfg(feature = "observers")]
363 observer_runtime,
364 #[cfg(feature = "auth")]
365 enrichment_pool: db_pool.clone(),
366 #[cfg(feature = "observers")]
367 db_pool,
368 storage_state: None,
369 realtime_state: None,
370 #[cfg(feature = "functions-runtime")]
371 functions_hooks: None,
372 tenant_executor_factory: None,
373 #[cfg(feature = "arrow")]
374 flight_service,
375 adapter_cache_enabled: false,
376 broadcast_manager: None,
377 presence_manager: None,
378 storage_backend: None,
379 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
381 function_store: None,
382 #[cfg(feature = "functions")]
383 function_runtime: None,
384 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
385 tasks,
386 })
387 }
388
389 #[cfg(feature = "observers")]
399 pub(super) async fn init_observer_runtime(
400 config: &ServerConfig,
401 pool: Option<&sqlx::PgPool>,
402 ) -> crate::Result<Option<Arc<RwLock<ObserverRuntime>>>> {
403 use fraiseql_observers::config::TransportKind;
404
405 let observer_config = match &config.observers {
407 Some(cfg) if cfg.enabled => cfg,
408 _ => {
409 info!("Observer runtime disabled");
410 return Ok(None);
411 },
412 };
413
414 let Some(pool) = pool else {
415 warn!("No database pool provided for observers");
416 return Ok(None);
417 };
418
419 info!("Initializing observer runtime");
420
421 let mut transport = observer_config.runtime.transport.clone().with_env_overrides();
425 let compiled_in = cfg!(feature = "observers-nats");
426 let nats_url_present = !transport.nats.url.is_empty();
427 crate::server::initialization::observer_transport_check(
428 transport.transport,
429 compiled_in,
430 nats_url_present,
431 crate::ServerConfig::is_production_mode(),
432 )?;
433
434 let usable = match transport.transport {
438 TransportKind::Postgres | TransportKind::InMemory => true,
439 TransportKind::Nats => compiled_in && nats_url_present,
440 _ => false,
441 };
442 if !usable {
443 transport.transport = TransportKind::Postgres;
444 }
445
446 transport.validate().map_err(|e| {
447 crate::ServerError::ConfigError(format!("invalid observer transport config: {e}"))
448 })?;
449
450 let runtime_config = ObserverRuntimeConfig::new(pool.clone())
451 .with_poll_interval(observer_config.runtime.poll_interval_ms)
452 .with_batch_size(observer_config.runtime.batch_size)
453 .with_channel_capacity(observer_config.runtime.channel_capacity)
454 .with_max_dlq_size(observer_config.runtime.max_dlq_size)
455 .with_transport(transport)
456 .with_email(observer_config.runtime.email.clone())
457 .with_log_payloads(observer_config.runtime.log_payloads);
458
459 let runtime = ObserverRuntime::new(runtime_config);
460 Ok(Some(Arc::new(RwLock::new(runtime))))
461 }
462}