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, 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 #[cfg(feature = "federation")]
93 let circuit_breaker = schema.federation.as_ref().and_then(
94 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
95 );
96 #[cfg(not(feature = "federation"))]
97 let circuit_breaker: Option<()> = None;
98 #[cfg(not(feature = "federation"))]
99 let _ = &schema.federation;
100 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
101 #[cfg(feature = "auth")]
102 let state_encryption = Self::state_encryption_from_schema(&schema)?;
103 #[cfg(not(feature = "auth"))]
104 let state_encryption: Option<
105 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
106 > = None;
107 #[cfg(feature = "auth")]
108 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
109 #[cfg(not(feature = "auth"))]
110 let pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
111 #[cfg(feature = "auth")]
112 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
113 #[cfg(not(feature = "auth"))]
114 let oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
115 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
116 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
117 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
118 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
119 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
120
121 let cache_config = CacheConfig::from(config.cache_enabled);
122 let cache = QueryResultCache::new(cache_config);
123 let inner = Arc::into_inner(adapter)
125 .expect("CachedDatabaseAdapter wrapping requires exclusive Arc ownership at startup");
126 let cached = CachedDatabaseAdapter::new(inner, cache, schema.content_hash())
127 .with_ttl_overrides_from_schema(&schema);
128 let executor = Arc::new(Executor::new_with_relay(schema.clone(), Arc::new(cached)));
129 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
130
131 let mut server = Self::from_executor(
132 config,
133 executor,
134 subscription_manager,
135 circuit_breaker,
136 error_sanitizer,
137 state_encryption,
138 pkce_store,
139 oidc_server_client,
140 schema_rate_limiter,
141 api_key_authenticator,
142 revocation_manager,
143 trusted_docs,
144 db_pool,
145 tasks,
146 )
147 .await?;
148
149 server.adapter_cache_enabled = cache_config.enabled;
150
151 #[cfg(feature = "mcp")]
153 if let Some(ref cfg) = server.executor.schema().mcp_config {
154 if cfg.enabled {
155 let tool_count =
156 crate::mcp::tools::schema_to_tools(server.executor.schema(), cfg).len();
157 info!(
158 path = %cfg.path,
159 transport = %cfg.transport,
160 tools = tool_count,
161 "MCP server configured"
162 );
163 server.mcp_config = Some(cfg.clone());
164 }
165 }
166
167 if server.config.apq_enabled {
169 let apq_store: fraiseql_core::apq::ArcApqStorage =
170 Arc::new(fraiseql_core::apq::InMemoryApqStorage::default());
171 server.apq_store = Some(apq_store);
172 info!("APQ (Automatic Persisted Queries) enabled — in-memory backend");
173 }
174
175 Ok(server)
176 }
177}
178
179impl<A: DatabaseAdapter + Clone + Send + Sync + 'static> Server<A> {
180 #[cfg(feature = "arrow")]
196 pub async fn with_flight_service(
197 config: ServerConfig,
198 schema: CompiledSchema,
199 adapter: Arc<A>,
200 #[allow(unused_variables)]
201 db_pool: Option<sqlx::PgPool>,
203 flight_service: Option<FraiseQLFlightService>,
204 ) -> Result<Self> {
205 #[cfg(feature = "federation")]
207 let circuit_breaker = schema.federation.as_ref().and_then(
208 crate::federation::circuit_breaker::FederationCircuitBreakerManager::from_config,
209 );
210 #[cfg(not(feature = "federation"))]
211 let _circuit_breaker: Option<()> = None;
212 #[cfg(not(feature = "federation"))]
213 let _ = &schema.federation;
214 let error_sanitizer = Self::error_sanitizer_from_schema(&schema);
215 #[cfg(feature = "auth")]
216 let state_encryption = Self::state_encryption_from_schema(&schema)?;
217 #[cfg(not(feature = "auth"))]
218 let _state_encryption: Option<
219 std::sync::Arc<crate::auth::state_encryption::StateEncryptionService>,
220 > = None;
221 #[cfg(feature = "auth")]
222 let pkce_store = Self::pkce_store_from_schema(&schema, state_encryption.as_ref()).await?;
223 #[cfg(not(feature = "auth"))]
224 let _pkce_store: Option<std::sync::Arc<crate::auth::PkceStateStore>> = None;
225 #[cfg(feature = "auth")]
226 let oidc_server_client = Self::oidc_server_client_from_schema(&schema);
227 #[cfg(not(feature = "auth"))]
228 let _oidc_server_client: Option<std::sync::Arc<crate::auth::OidcServerClient>> = None;
229 let schema_rate_limiter = Self::rate_limiter_from_schema(&schema).await?;
230 let api_key_authenticator = crate::api_key::api_key_authenticator_from_schema(&schema);
231 let revocation_manager = crate::token_revocation::revocation_manager_from_schema(&schema)?;
232 let mut tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
233 let trusted_docs = Self::trusted_docs_from_schema(&schema, &mut tasks);
234
235 let executor = Arc::new(Executor::new(schema.clone(), adapter));
236 let subscription_manager = Arc::new(SubscriptionManager::new(Arc::new(schema)));
237
238 #[cfg(feature = "auth")]
240 let oidc_validator = if let Some(ref auth_config) = config.auth {
241 info!(
242 issuer = %auth_config.issuer,
243 "Initializing OIDC authentication"
244 );
245 let validator = OidcValidator::new(auth_config.clone())
246 .await
247 .map_err(|e| ServerError::ConfigError(format!("Failed to initialize OIDC: {e}")))?;
248 Some(Arc::new(validator))
249 } else {
250 None
251 };
252 #[cfg(not(feature = "auth"))]
253 let oidc_validator: Option<Arc<fraiseql_core::security::OidcValidator>> = None;
254
255 let hs256_auth = super::builder::build_hs256_auth(&config)?;
257
258 let rate_limiter = if let Some(rl) = schema_rate_limiter {
260 Some(rl)
261 } else if let Some(ref rate_config) = config.rate_limiting {
262 if rate_config.enabled {
263 info!(
264 rps_per_ip = rate_config.rps_per_ip,
265 rps_per_user = rate_config.rps_per_user,
266 "Initializing rate limiting from server config"
267 );
268 Some(Arc::new(RateLimiter::new(rate_config.clone())))
269 } else {
270 info!("Rate limiting disabled by configuration");
271 None
272 }
273 } else {
274 None
275 };
276
277 #[cfg(feature = "observers")]
279 let observer_runtime = Self::init_observer_runtime(&config, db_pool.as_ref()).await?;
280
281 #[cfg(feature = "auth")]
283 if pkce_store.is_some() && oidc_server_client.is_none() {
284 tracing::error!(
285 "pkce.enabled = true but [auth] is not configured or OIDC client init failed. \
286 Auth routes will NOT be mounted."
287 );
288 }
289
290 #[cfg(feature = "auth")]
292 Self::check_redis_requirement(pkce_store.as_ref())?;
293
294 #[cfg(feature = "auth")]
296 Self::spawn_pkce_cleanup(pkce_store.as_ref(), &mut tasks);
297
298 let apq_enabled = config.apq_enabled;
299
300 Ok(Self {
301 config,
302 executor,
303 subscription_manager,
304 subscription_lifecycle: Arc::new(crate::subscriptions::NoopLifecycle),
305 max_subscriptions_per_connection: None,
306 oidc_validator,
307 hs256_auth,
308 rate_limiter,
309 #[cfg(feature = "secrets")]
310 secrets_manager: None,
311 #[cfg(feature = "federation")]
312 circuit_breaker,
313 error_sanitizer,
314 #[cfg(feature = "auth")]
315 state_encryption,
316 #[cfg(feature = "auth")]
317 pkce_store,
318 #[cfg(feature = "auth")]
319 oidc_server_client,
320 #[cfg(feature = "auth")]
321 social_login: None,
322 #[cfg(feature = "auth")]
323 anon_signup_state: None,
324 #[cfg(feature = "auth")]
325 mfa_state: None,
326 api_key_authenticator,
327 revocation_manager,
328 apq_store: if apq_enabled {
329 Some(Arc::new(fraiseql_core::apq::InMemoryApqStorage::default())
330 as fraiseql_core::apq::ArcApqStorage)
331 } else {
332 None
333 },
334 trusted_docs,
335 #[cfg(feature = "mcp")]
336 mcp_config: None,
337 pool_tuning_config: None,
338 #[cfg(feature = "observers")]
339 observer_runtime,
340 #[cfg(feature = "auth")]
341 enrichment_pool: db_pool.clone(),
342 #[cfg(feature = "observers")]
343 db_pool,
344 storage_state: None,
345 realtime_state: None,
346 tenant_executor_factory: None,
347 #[cfg(feature = "arrow")]
348 flight_service,
349 adapter_cache_enabled: false,
350 broadcast_manager: None,
351 presence_manager: None,
352 storage_backend: None,
353 storage_max_upload_bytes: 100 * 1024 * 1024, #[cfg(feature = "functions")]
355 function_store: None,
356 #[cfg(feature = "functions")]
357 function_runtime: None,
358 usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
359 tasks,
360 })
361 }
362
363 #[cfg(feature = "observers")]
373 pub(super) async fn init_observer_runtime(
374 config: &ServerConfig,
375 pool: Option<&sqlx::PgPool>,
376 ) -> crate::Result<Option<Arc<RwLock<ObserverRuntime>>>> {
377 use fraiseql_observers::config::TransportKind;
378
379 let observer_config = match &config.observers {
381 Some(cfg) if cfg.enabled => cfg,
382 _ => {
383 info!("Observer runtime disabled");
384 return Ok(None);
385 },
386 };
387
388 let Some(pool) = pool else {
389 warn!("No database pool provided for observers");
390 return Ok(None);
391 };
392
393 info!("Initializing observer runtime");
394
395 let mut transport = observer_config.runtime.transport.clone().with_env_overrides();
399 let compiled_in = cfg!(feature = "observers-nats");
400 let nats_url_present = !transport.nats.url.is_empty();
401 crate::server::initialization::observer_transport_check(
402 transport.transport,
403 compiled_in,
404 nats_url_present,
405 crate::ServerConfig::is_production_mode(),
406 )?;
407
408 let usable = match transport.transport {
412 TransportKind::Postgres | TransportKind::InMemory => true,
413 TransportKind::Nats => compiled_in && nats_url_present,
414 _ => false,
415 };
416 if !usable {
417 transport.transport = TransportKind::Postgres;
418 }
419
420 transport.validate().map_err(|e| {
421 crate::ServerError::ConfigError(format!("invalid observer transport config: {e}"))
422 })?;
423
424 let runtime_config = ObserverRuntimeConfig::new(pool.clone())
425 .with_poll_interval(observer_config.runtime.poll_interval_ms)
426 .with_batch_size(observer_config.runtime.batch_size)
427 .with_channel_capacity(observer_config.runtime.channel_capacity)
428 .with_max_dlq_size(observer_config.runtime.max_dlq_size)
429 .with_transport(transport)
430 .with_email(observer_config.runtime.email.clone());
431
432 let runtime = ObserverRuntime::new(runtime_config);
433 Ok(Some(Arc::new(RwLock::new(runtime))))
434 }
435}