1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
//! `AppState` — server state passed to all GraphQL route handlers.
use std::{path::PathBuf, sync::Arc};
use arc_swap::ArcSwap;
use fraiseql_core::{
apq::{ApqMetrics, ArcApqStorage},
runtime::Executor,
schema::CompiledSchema,
security::IntrospectionPolicy,
};
use tracing::{info, warn};
use super::{tenant_key::DomainRegistry, tenant_registry::TenantExecutorRegistry};
#[cfg(feature = "auth")]
use crate::auth::rate_limiting::{AuthRateLimitConfig, KeyedRateLimiter};
use crate::{
config::error_sanitization::ErrorSanitizer, error::GraphQLError,
metrics_server::MetricsCollector, usage::aggregator::UsageAggregator,
};
/// Server state containing executor and configuration.
#[derive(Clone)]
pub struct AppState {
/// Query executor (atomically swappable for schema hot-reload).
pub executor: Arc<ArcSwap<Executor>>,
/// Metrics collector.
pub metrics: Arc<MetricsCollector>,
/// Query result cache (optional).
#[cfg(feature = "arrow")]
pub cache: Option<Arc<fraiseql_arrow::cache::QueryCache>>,
/// Rate limiter for GraphQL validation errors (per IP).
#[cfg(feature = "auth")]
pub graphql_rate_limiter: Arc<KeyedRateLimiter>,
/// Secrets manager (optional, configured via `[fraiseql.secrets]`).
#[cfg(feature = "secrets")]
pub secrets_manager: Option<Arc<crate::secrets_manager::SecretsManager>>,
/// Field encryption service for transparent encrypt/decrypt of marked fields.
#[cfg(feature = "secrets")]
pub field_encryption: Option<Arc<crate::encryption::middleware::FieldEncryptionService>>,
/// Federation circuit breaker manager (optional, enabled via `fraiseql.toml`).
#[cfg(feature = "federation")]
pub circuit_breaker:
Option<Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>>,
/// Federation subgraph latency histogram tracker.
#[cfg(feature = "federation")]
pub federation_latency: Arc<fraiseql_core::federation::SubgraphLatencyTracker>,
/// Federation entity resolution counter metrics.
#[cfg(feature = "federation")]
pub federation_entity_metrics: Arc<fraiseql_core::federation::EntityResolutionMetrics>,
/// Federation query plan cache for plan visualization.
#[cfg(feature = "federation")]
pub federation_plan_cache: Option<Arc<fraiseql_core::federation::QueryPlanCache>>,
/// Error sanitizer — strips internal details before sending responses to clients.
pub error_sanitizer: Arc<ErrorSanitizer>,
/// State encryption service (optional, enabled via `[security.state_encryption]`).
#[cfg(feature = "auth")]
pub state_encryption: Option<Arc<crate::auth::state_encryption::StateEncryptionService>>,
/// API key authenticator (optional, enabled via `[security.api_keys]`).
pub api_key_authenticator: Option<Arc<crate::api_key::ApiKeyAuthenticator>>,
/// Service-account authenticator (optional, enabled via `[security.service_accounts]`;
/// ADR-0018).
pub service_account_authenticator:
Option<Arc<crate::service_account::ServiceAccountAuthenticator>>,
/// APQ persistent query store (optional, enabled via compiled schema config).
pub apq_store: Option<ArcApqStorage>,
/// Trusted document store (optional, enabled via `[security.trusted_documents]`).
pub trusted_docs: Option<Arc<crate::trusted_documents::TrustedDocumentStore>>,
/// APQ metrics tracker.
pub apq_metrics: Arc<ApqMetrics>,
/// Request validator (depth/complexity limits, configured from compiled schema).
pub validator: crate::validation::RequestValidator,
/// Debug configuration (optional, from `[debug]` in `fraiseql.toml`).
pub debug_config: Option<fraiseql_core::schema::DebugConfig>,
/// Maximum byte length for a query string delivered via HTTP GET.
///
/// Defaults to `100_000` (100 `KiB`). Configurable via
/// `ServerConfig::max_get_query_bytes`.
pub max_get_query_bytes: usize,
/// Whether the GraphQL-over-SSE response transport is enabled (#387).
///
/// Defaults to `false` (the `Accept: text/event-stream` header is ignored).
/// Set from `ServerConfig::enable_graphql_incremental` in `build_app_state`.
pub graphql_incremental_enabled: bool,
/// Continuation batch size for `@stream` deliveries (#387).
///
/// Set from `ServerConfig::graphql_incremental_batch_size()` in `build_app_state`.
pub graphql_incremental_batch_size: u32,
/// Introspection policy for the GraphQL request path.
///
/// Derived from `ServerConfig::introspection_enabled` /
/// `introspection_require_auth` via [`IntrospectionPolicy::from_config`] in
/// `build_app_state`. Defaults to [`IntrospectionPolicy::Disabled`]
/// (fail-closed) when no config is wired, matching the server's
/// introspection-off-by-default posture.
pub introspection_policy: IntrospectionPolicy,
/// Connection pool auto-tuner (optional, enabled via `[pool_tuning]` config).
pub pool_tuner: Option<Arc<crate::pool::PoolSizingAdvisor>>,
/// Observer runtime handle for health probes (optional, requires `observers` feature).
#[cfg(feature = "observers")]
pub observer_runtime: Option<Arc<tokio::sync::RwLock<crate::observers::ObserverRuntime>>>,
/// Multi-consumer fan-out of the entity events the observer runtime forwards
/// through the `EventBridge` (#1309). Read by the REST `/{resource}/stream` mount.
/// `Some` exactly when an observer runtime is configured, so a stream request can
/// tell "no producer" from "no events yet".
#[cfg(feature = "observers")]
pub entity_event_fanout: Option<crate::subscriptions::EntityEventFanout>,
/// Reads back what the REST `/{resource}/stream` mount already delivered, for a
/// client reconnecting with `Last-Event-ID` (#1310). `Some` exactly when an observer
/// runtime polls the local change log — the ledger whose recorded order it replays.
#[cfg(feature = "observers")]
pub stream_replay: Option<std::sync::Arc<fraiseql_observers::listener::ChangeLogReplayReader>>,
/// Schema file path for reload operations.
pub schema_path: Option<PathBuf>,
/// Database adapter reference for constructing new executors on reload.
/// How the booting constructor built its executor, so a reload rebuilds the
/// *same kind*.
///
/// Relay dispatch requires `A: RelayDatabaseAdapter`, a bound `AppState` does
/// not carry — so only the constructor that had it can say how to rebuild.
/// Reload used to call `Executor::new` unconditionally, which dropped relay
/// dispatch and made every relay query fail validation until the process
/// restarted (#750).
/// Reload mutex to serialize concurrent reload attempts.
pub(crate) reload_lock: Arc<tokio::sync::Mutex<()>>,
/// Whether the adapter-level query result cache is active.
///
/// Set to `true` when `ServerConfig::cache_enabled = true` and the server
/// was built via `Server::new` or `Server::with_relay_pagination`.
/// This reflects the adapter-level `CachedDatabaseAdapter` state, NOT the
/// Arrow flight cache (`AppState::cache`).
pub adapter_cache_enabled: bool,
/// Multi-tenant executor registry (optional).
///
/// When `Some`, the server operates in multi-tenant mode: each request's
/// tenant key selects an executor from this registry. When `None`,
/// single-tenant mode is in effect and all requests use `self.executor`.
pub tenant_registry: Option<Arc<TenantExecutorRegistry>>,
/// Factory for creating tenant executors from schema JSON + pool config.
///
/// Type-erased so that the management API handler does not need
/// `A: FromPoolConfig` on its generic bounds.
pub tenant_executor_factory: Option<crate::tenancy::TenantExecutorFactory>,
/// Domain-to-tenant mapping for Host header-based tenant resolution.
pub domain_registry: Arc<DomainRegistry>,
/// Tenant audit log (optional, for lifecycle event recording).
pub tenant_audit_log: Option<crate::tenancy::audit::AuditLogHandle>,
/// Usage aggregator — shared with the `MutationAuditLayer` tracing subscriber.
///
/// Always present (never `Option`): when audit logging is disabled the
/// aggregator simply receives no events and every query returns empty counts.
pub usage: Arc<UsageAggregator>,
/// Before-mutation hooks from the functions subsystem (optional).
///
/// When `Some`, every GraphQL mutation is checked against the trigger registry
/// before execution. The check is a single `HashMap::get` returning `None`
/// when no hooks are registered — zero overhead for mutations without hooks.
pub before_mutation_hooks: Option<Arc<crate::subsystems::BeforeMutationHooks>>,
/// Enrichment-profile identity resolver (#539). `Some` when
/// `[identity.enrichment].enabled` and an auth DB pool is present; every
/// authenticated request then resolves its DB identity and fail-closes
/// before dispatch (read-scoping via the `fraiseql.enriched.*` namespace).
#[cfg(feature = "auth")]
pub identity_resolver: Option<Arc<crate::identity::IdentityResolver>>,
/// Token revocation manager, when `[security.token_revocation]` is configured.
///
/// Reachable from `AppState` so `POST /admin/v1/users/{id}/revoke` can actually
/// revoke. It used to answer `{"success": true, "message": "All sessions
/// revoked"}` without touching any store, because the manager lived on `Server`
/// and the handler had no way to reach it (#749). `None` is now an explicit
/// `501`, never a fabricated success.
pub revocation_manager: Option<Arc<crate::token_revocation::TokenRevocationManager>>,
/// When this state was constructed — i.e. when the server started.
///
/// `GET /admin/v1/health/detailed` reported `uptime_secs` as *seconds since the
/// Unix epoch*, so a freshly-booted server claimed roughly 1.8 billion seconds
/// of uptime. A real instant is the only way to answer the question asked.
pub started_at: std::time::Instant,
/// #611 (layer 2): schema hot-reload signal for live subscriptions.
///
/// Bumped by [`swap_in_schema`](Self::reload_schema) on every successful
/// executor swap. The `/ws` mount subscribes each connection to it (via
/// [`subscribe_policy_reload`](Self::subscribe_policy_reload)); on a bump,
/// active subscriptions re-derive their row-visibility conditions against
/// the current policies — re-scoped in place, or terminated when the new
/// policy refuses (fail-closed).
pub policy_reload: Arc<tokio::sync::watch::Sender<u64>>,
/// Idempotency store for the GraphQL mutation path (#747).
///
/// A POST mutation carrying an `Idempotency-Key` header is deduplicated
/// against it: a repeat of an already-executed mutation replays the stored
/// response instead of executing again. This is the receiving half of the
/// saga at-least-once dispatch contract — a peer coordinator retries
/// ambiguous failures (timeouts, connection resets) under the same key, and
/// this store is what turns those retries into one logical effect.
/// In-memory with TTL expiry: replicas do not share it, so a load balancer
/// that re-routes a retry to another replica re-executes (document
/// `Idempotency-Key` affinity or use sticky routing for saga peers).
pub idempotency_store: Arc<dyn crate::routes::idempotency::IdempotencyStore>,
}
impl AppState {
/// Create new application state.
#[must_use]
pub fn new(executor: Arc<Executor>) -> Self {
Self {
executor: Arc::new(ArcSwap::from(executor)),
metrics: Arc::new(MetricsCollector::new()),
#[cfg(feature = "arrow")]
cache: None,
#[cfg(feature = "auth")]
graphql_rate_limiter: Arc::new(KeyedRateLimiter::new(
AuthRateLimitConfig::per_ip_standard(),
)),
#[cfg(feature = "secrets")]
secrets_manager: None,
#[cfg(feature = "secrets")]
field_encryption: None,
#[cfg(feature = "federation")]
circuit_breaker: None,
#[cfg(feature = "federation")]
federation_latency: Arc::new(fraiseql_core::federation::SubgraphLatencyTracker::new()),
#[cfg(feature = "federation")]
federation_entity_metrics: Arc::new(
fraiseql_core::federation::EntityResolutionMetrics::new(),
),
#[cfg(feature = "federation")]
federation_plan_cache: None,
error_sanitizer: Arc::new(ErrorSanitizer::disabled()),
#[cfg(feature = "auth")]
state_encryption: None,
api_key_authenticator: None,
service_account_authenticator: None,
apq_store: None,
trusted_docs: None,
apq_metrics: Arc::new(ApqMetrics::default()),
validator: crate::validation::RequestValidator::new(),
debug_config: None,
pool_tuner: None,
#[cfg(feature = "observers")]
observer_runtime: None,
#[cfg(feature = "observers")]
entity_event_fanout: None,
#[cfg(feature = "observers")]
stream_replay: None,
max_get_query_bytes: 100_000,
graphql_incremental_enabled: false,
graphql_incremental_batch_size: 100,
introspection_policy: IntrospectionPolicy::Disabled,
schema_path: None,
reload_lock: Arc::new(tokio::sync::Mutex::new(())),
adapter_cache_enabled: false,
tenant_registry: None,
tenant_executor_factory: None,
domain_registry: Arc::new(DomainRegistry::new()),
tenant_audit_log: None,
usage: Arc::clone(crate::usage::aggregator::global_aggregator()),
before_mutation_hooks: None,
#[cfg(feature = "auth")]
identity_resolver: None,
revocation_manager: None,
started_at: std::time::Instant::now(),
policy_reload: Arc::new(tokio::sync::watch::channel(0).0),
idempotency_store: crate::routes::idempotency::create_store(
crate::routes::idempotency::GRAPHQL_IDEMPOTENCY_TTL_SECS,
),
}
}
/// Subscribe to the schema hot-reload signal (#611 layer 2). Each live
/// `/ws` connection holds one receiver; a bump makes it re-derive its
/// active subscriptions' row-visibility conditions.
#[must_use]
pub fn subscribe_policy_reload(&self) -> tokio::sync::watch::Receiver<u64> {
self.policy_reload.subscribe()
}
/// Attach the token revocation manager so the admin revoke endpoint can use it.
#[must_use]
pub fn with_revocation_manager(
mut self,
manager: Arc<crate::token_revocation::TokenRevocationManager>,
) -> Self {
self.revocation_manager = Some(manager);
self
}
/// Attach the enrichment-profile identity resolver (#539).
#[cfg(feature = "auth")]
#[must_use]
pub fn with_identity_resolver(
mut self,
resolver: Arc<crate::identity::IdentityResolver>,
) -> Self {
self.identity_resolver = Some(resolver);
self
}
/// Load the current executor.
///
/// Returns a guard that keeps the executor alive for the duration of the
/// request. This is wait-free (no lock).
#[must_use]
pub fn executor(&self) -> arc_swap::Guard<Arc<Executor>> {
self.executor.load()
}
/// Atomically swap the executor.
///
/// In-flight requests that already called `executor()` continue using
/// the old executor until their guard is dropped.
pub fn swap_executor(&self, new_executor: Arc<Executor>) {
self.executor.store(new_executor);
}
/// Returns the executor for the given tenant key.
///
/// In multi-tenant mode, delegates to the `TenantExecutorRegistry`. In
/// single-tenant mode (no registry), ignores the key and returns the
/// default executor.
///
/// # Errors
///
/// Returns `FraiseQLError::Authorization` if multi-tenant mode is enabled
/// and the tenant key is explicit but not registered.
pub fn executor_for_tenant(
&self,
tenant_key: Option<&str>,
) -> fraiseql_error::Result<arc_swap::Guard<Arc<Executor>>> {
match &self.tenant_registry {
Some(registry) => registry.executor_for(tenant_key),
None => Ok(self.executor()),
}
}
/// Attach a multi-tenant executor registry.
#[must_use]
pub fn with_tenant_registry(mut self, registry: Arc<TenantExecutorRegistry>) -> Self {
self.tenant_registry = Some(registry);
self
}
/// Get the tenant registry if multi-tenant mode is enabled.
#[must_use]
pub const fn tenant_registry(&self) -> Option<&Arc<TenantExecutorRegistry>> {
self.tenant_registry.as_ref()
}
/// Attach a tenant executor factory for the management API.
#[must_use]
pub fn with_tenant_executor_factory(
mut self,
factory: crate::tenancy::TenantExecutorFactory,
) -> Self {
self.tenant_executor_factory = Some(factory);
self
}
/// Get the tenant executor factory if configured.
#[must_use]
pub const fn tenant_executor_factory(&self) -> Option<&crate::tenancy::TenantExecutorFactory> {
self.tenant_executor_factory.as_ref()
}
/// Get the domain registry for Host header-based tenant resolution.
#[must_use]
pub const fn domain_registry(&self) -> &Arc<DomainRegistry> {
&self.domain_registry
}
/// Attach a custom domain registry.
#[must_use]
pub fn with_domain_registry(mut self, registry: Arc<DomainRegistry>) -> Self {
self.domain_registry = registry;
self
}
/// Replace the usage aggregator (primarily for testing with an isolated aggregator).
#[must_use]
pub fn with_usage(mut self, usage: Arc<UsageAggregator>) -> Self {
self.usage = usage;
self
}
/// Attach a tenant audit log for lifecycle event recording.
#[must_use]
pub fn with_tenant_audit_log(mut self, log: crate::tenancy::audit::AuditLogHandle) -> Self {
self.tenant_audit_log = Some(log);
self
}
/// Get the tenant audit log if configured.
#[must_use]
pub const fn tenant_audit_log(&self) -> Option<&crate::tenancy::audit::AuditLogHandle> {
self.tenant_audit_log.as_ref()
}
/// Attach before-mutation hooks from the functions subsystem.
///
/// When set, every incoming GraphQL mutation is checked against the trigger
/// registry before execution. The check is a single `HashMap::get` returning
/// `None` when no hooks exist — zero overhead for mutations without hooks.
#[must_use]
pub fn with_functions(mut self, hooks: Arc<crate::subsystems::BeforeMutationHooks>) -> Self {
self.before_mutation_hooks = Some(hooks);
self
}
/// Configure reload support with the schema file path to reload from.
///
/// Took an adapter and the booting constructor's executor rebuilder until the
/// boundary work: a reload had to re-run the constructor that had the
/// `RelayDatabaseAdapter` bound in scope, and #750 guarded by hand against
/// recording the wrong one. [`Executor::rebuild_with`] carries relay dispatch
/// over by construction, so there is nothing left to record or to get wrong, and
/// a directly-assembled test `AppState` can now reload like any other.
#[must_use]
pub fn with_reload_config(mut self, schema_path: PathBuf) -> Self {
self.schema_path = Some(schema_path);
self
}
/// Reload the compiled schema from a file path.
///
/// Reads the schema file and hands it to the shared `swap_in_schema` seam.
///
/// # Errors
///
/// Returns an error if the file cannot be read, the JSON is invalid, schema
/// validation fails, a boot-time safety gate refuses the new schema, or the
/// new schema changes configuration that only a restart can apply. On error,
/// the current executor is unchanged.
pub async fn reload_schema(&self, path: &std::path::Path) -> Result<(), String> {
// Take the lock and check the preconditions before touching the disk: the
// lock exists to serialize *reloads*, and a second reload must be refused
// before it starts reading, not after.
let guard = self.begin_reload()?;
let json = tokio::fs::read_to_string(path)
.await
.map_err(|e| format!("Failed to read schema file {}: {e}", path.display()))?;
let schema = CompiledSchema::from_json(&json, false)
.map_err(|e| format!("Invalid schema JSON: {e}"))?;
self.swap_in_schema(schema, &guard).await
}
/// Reload the compiled schema from already-validated JSON bytes.
///
/// This avoids re-reading the schema file from disk after validation,
/// preventing TOCTOU race conditions where the file could change between
/// validation and reload.
///
/// # Errors
///
/// Returns an error if the JSON is invalid, schema validation fails, a
/// boot-time safety gate refuses the new schema, the new schema changes
/// configuration that only a restart can apply, or a reload is already in
/// progress. On error, the current executor is unchanged.
pub async fn reload_schema_from_json(&self, json: &str) -> Result<(), String> {
let guard = self.begin_reload()?;
let schema = CompiledSchema::from_json(json, false)
.map_err(|e| format!("Invalid schema JSON: {e}"))?;
self.swap_in_schema(schema, &guard).await
}
/// Acquire the reload lock and check that reload is configured at all.
///
/// Both preconditions are cheap and both are refusals rather than failures,
/// so they run before any I/O: a concurrent reload must be told "already in
/// progress" rather than racing on the file read, and an `AppState` with no
/// adapter must say so rather than reporting a read error for a reload it
/// could never have performed.
fn begin_reload(&self) -> Result<tokio::sync::MutexGuard<'_, ()>, String> {
let guard = self
.reload_lock
.try_lock()
.map_err(|_| "Reload already in progress".to_string())?;
Ok(guard)
}
/// The single hot-reload seam: validate, gate, rebuild, swap.
///
/// This is a **construction path**, and it must produce the same configured
/// runtime as boot does. It previously did not: it called `Executor::new`,
/// which uses `RuntimeConfig::default()`, so a successful reload silently
/// reverted mutation audit logging, the #421 page-size ceiling, the
/// change-log toggle and relay dispatch (#750); and it ran none of the
/// boot-time safety gates, so it could move a running server into a state
/// boot would have refused (#782).
///
/// # Errors
///
/// Returns the operator-facing message for every refusal: an incompatible
/// format version, a field marked for at-rest encryption (H12), a
/// multi-tenant schema that cannot isolate its tenants under caching (#758),
/// or boot-frozen configuration drift that requires a restart.
// Reason: one step of the awaited hot-reload sequence — `begin_reload` and the
// validation either side of it do I/O. The swap itself is a pointer store.
#[allow(unknown_lints, clippy::unused_async_trait_impl)]
async fn swap_in_schema(
&self,
schema: CompiledSchema,
// Held by the caller for the whole reload — taken in `begin_reload`,
// before any I/O. Threaded through so this seam cannot be reached without
// it.
_guard: &tokio::sync::MutexGuard<'_, ()>,
) -> Result<(), String> {
schema
.validate_producer_version()
.map_err(|msg| format!("Incompatible compiled schema: {msg}"))?;
let current = self.executor.load();
if current.schema().content_hash() == schema.content_hash() {
return Ok(()); // Same schema, no-op
}
// The boot-time safety gates, run again. A reload that skips them is a
// way to reach, at runtime, a configuration the server refused to start
// in — which for the encryption gate means writing plaintext into a field
// declared encrypted, and for the tenancy gate means serving one tenant's
// cached rows to another (#782).
crate::server::initialization::field_encryption_unsupported_check(&schema)
.map_err(|e| e.to_string())?;
crate::server::initialization::tenant_isolation_declaration_check(
&schema,
self.adapter_cache_enabled,
)
.map_err(|e| e.to_string())?;
// Refuse rather than half-apply: everything a boot-time subsystem read
// once cannot be changed by swapping the executor.
super::reload_gate::check_reloadable(current.schema(), &schema)?;
// #1390: a reloaded source must be readable on the replicas reads go to.
current
.refuse_standby_unreadable_sources(&schema)
.await
.map_err(|e| e.to_string())?;
// #611: new subscriptions pick up policy changes immediately (layer-1); warn loudly
// so operators know already-connected streams must reconnect to apply the change.
warn_on_subscription_policy_reload(current.schema(), &schema);
// Re-derive the schema-owned runtime settings on top of the *live* config,
// so programmatically-installed pieces (authorizers, RLS policy, field
// filter, query validation) survive the swap while the compiled and
// environment-derived ones are recomputed — exactly what boot does.
let config = current
.config()
.clone()
.with_compiled_schema(&schema)
.map_err(|msg| format!("Incompatible compiled schema: {msg}"))?;
// Notify the backend of the schema change (clears the query result cache if
// applicable) before the swap, while the stale derivations are still reachable.
current.on_schema_reload();
// Rebuild over the same backend. `rebuild_with` carries relay dispatch across
// by construction, so a reload cannot downgrade a relay executor (#750).
let new_executor = Arc::new(current.rebuild_with(schema, config));
// Atomic swap
self.executor.store(new_executor);
// Clear query plan caches (reference old schema)
#[cfg(feature = "arrow")]
if let Some(cache) = &self.cache {
cache.clear();
}
// #611 (layer 2): tell live subscription connections the schema changed, so
// they re-derive their row-visibility conditions against the new policies
// (re-scope in place, or terminate fail-closed). Bumped on every successful
// swap — connections with no policy-declaring subscriptions re-derive for
// free (the policy lookup short-circuits).
self.policy_reload.send_modify(|generation| *generation += 1);
info!("Schema executor swapped successfully");
Ok(())
}
/// Create new application state with custom metrics collector.
#[must_use]
pub fn with_metrics(executor: Arc<Executor>, metrics: Arc<MetricsCollector>) -> Self {
Self::new(executor).set_metrics(metrics)
}
/// Create new application state with cache.
#[cfg(feature = "arrow")]
#[must_use]
pub fn with_cache(
executor: Arc<Executor>,
cache: Arc<fraiseql_arrow::cache::QueryCache>,
) -> Self {
Self::new(executor).set_cache(cache)
}
fn set_metrics(mut self, metrics: Arc<MetricsCollector>) -> Self {
self.metrics = metrics;
self
}
#[cfg(feature = "arrow")]
fn set_cache(mut self, cache: Arc<fraiseql_arrow::cache::QueryCache>) -> Self {
self.cache = Some(cache);
self
}
/// Get query cache if configured.
#[cfg(feature = "arrow")]
#[must_use]
pub const fn cache(&self) -> Option<&Arc<fraiseql_arrow::cache::QueryCache>> {
self.cache.as_ref()
}
/// Set secrets manager (for credential and secret management).
#[cfg(feature = "secrets")]
#[must_use]
pub fn with_secrets_manager(
mut self,
secrets_manager: Arc<crate::secrets_manager::SecretsManager>,
) -> Self {
self.secrets_manager = Some(secrets_manager);
self
}
/// Get secrets manager if configured.
#[cfg(feature = "secrets")]
#[must_use]
pub const fn secrets_manager(&self) -> Option<&Arc<crate::secrets_manager::SecretsManager>> {
self.secrets_manager.as_ref()
}
/// Attach a field encryption service (derived from schema and secrets manager).
#[cfg(feature = "secrets")]
#[must_use]
pub fn with_field_encryption(
mut self,
service: Arc<crate::encryption::middleware::FieldEncryptionService>,
) -> Self {
self.field_encryption = Some(service);
self
}
/// Attach a federation circuit breaker manager.
#[cfg(feature = "federation")]
#[must_use]
pub fn with_circuit_breaker(
mut self,
circuit_breaker: Arc<crate::federation::circuit_breaker::FederationCircuitBreakerManager>,
) -> Self {
self.circuit_breaker = Some(circuit_breaker);
self
}
/// Attach an error sanitizer (loaded from `compiled.security.error_sanitization`).
#[must_use]
pub fn with_error_sanitizer(mut self, sanitizer: Arc<ErrorSanitizer>) -> Self {
self.error_sanitizer = sanitizer;
self
}
/// Attach a state encryption service (loaded from `compiled.security.state_encryption`).
#[cfg(feature = "auth")]
#[must_use]
pub fn with_state_encryption(
mut self,
svc: Arc<crate::auth::state_encryption::StateEncryptionService>,
) -> Self {
self.state_encryption = Some(svc);
self
}
/// Attach an API key authenticator (loaded from `compiled.security.api_keys`).
#[must_use]
pub fn with_api_key_authenticator(
mut self,
authenticator: Arc<crate::api_key::ApiKeyAuthenticator>,
) -> Self {
self.api_key_authenticator = Some(authenticator);
self
}
/// Attach a service-account authenticator (loaded from
/// `compiled.security.service_accounts`; ADR-0018).
#[must_use]
pub fn with_service_account_authenticator(
mut self,
authenticator: Arc<crate::service_account::ServiceAccountAuthenticator>,
) -> Self {
self.service_account_authenticator = Some(authenticator);
self
}
/// Attach an APQ store for Automatic Persisted Queries.
#[must_use]
pub fn with_apq_store(mut self, store: ArcApqStorage) -> Self {
self.apq_store = Some(store);
self
}
/// Attach a trusted document store for query allowlist enforcement.
#[must_use]
pub fn with_trusted_docs(
mut self,
store: Arc<crate::trusted_documents::TrustedDocumentStore>,
) -> Self {
self.trusted_docs = Some(store);
self
}
/// Set the request validator (query depth/complexity limits).
#[must_use]
pub const fn with_validator(mut self, validator: crate::validation::RequestValidator) -> Self {
self.validator = validator;
self
}
/// Set the introspection policy for the GraphQL request path.
///
/// Wired in `build_app_state` from the server config via
/// [`IntrospectionPolicy::from_config`]; the default is
/// [`IntrospectionPolicy::Disabled`] (fail-closed).
#[must_use]
pub const fn with_introspection_policy(mut self, policy: IntrospectionPolicy) -> Self {
self.introspection_policy = policy;
self
}
/// Attach an adaptive connection pool auto-tuner.
#[must_use]
pub fn with_pool_tuner(mut self, tuner: Arc<crate::pool::PoolSizingAdvisor>) -> Self {
self.pool_tuner = Some(tuner);
self
}
/// Set whether the adapter-level cache is active.
///
/// Called from `build_router` to thread the cache state through to admin handlers.
#[must_use]
pub const fn with_adapter_cache_enabled(mut self, enabled: bool) -> Self {
self.adapter_cache_enabled = enabled;
self
}
/// Attach observer runtime for health probes.
#[cfg(feature = "observers")]
#[must_use]
pub fn with_observer_runtime(
mut self,
runtime: Arc<tokio::sync::RwLock<crate::observers::ObserverRuntime>>,
) -> Self {
self.observer_runtime = Some(runtime);
self
}
/// Attach the entity-event fan-out the REST `/{resource}/stream` mount reads (#1309).
#[cfg(feature = "observers")]
#[must_use]
pub fn with_entity_event_fanout(
mut self,
fanout: crate::subscriptions::EntityEventFanout,
) -> Self {
self.entity_event_fanout = Some(fanout);
self
}
/// Attach the reader a resumed REST stream reads its catch-up from (#1310).
///
/// Separate from [`with_entity_event_fanout`](Self::with_entity_event_fanout)
/// because the two answer different questions: the fan-out says whether there is a
/// producer at all, this says whether what it produced was recorded anywhere. A
/// deployment can have the first without the second, and then a stream is live-only
/// and says so.
#[cfg(feature = "observers")]
#[must_use]
pub fn with_stream_replay(
mut self,
reader: std::sync::Arc<fraiseql_observers::listener::ChangeLogReplayReader>,
) -> Self {
self.stream_replay = Some(reader);
self
}
/// Sanitize a batch of errors before sending them to the client.
#[must_use]
pub fn sanitize_errors(&self, errors: Vec<GraphQLError>) -> Vec<GraphQLError> {
self.error_sanitizer.sanitize_all(errors)
}
}
/// Warn (loudly) when a schema hot-reload changes the subscription row-visibility
/// policies (#596/#611).
///
/// **New** subscriptions read the live policies (the `/ws` handler resolves them from
/// the reload-aware executor `ArcSwap`, #611 layer-1), and **already-connected**
/// subscriptions re-derive against them on the [`AppState::policy_reload`] bump (#611
/// layer-2) — re-scoped in place, or terminated fail-closed. This warning simply makes
/// a policy-changing reload visible in the operator's log next to the swap itself.
fn warn_on_subscription_policy_reload(
old_schema: &CompiledSchema,
new_schema: &CompiledSchema,
) -> bool {
let old_policies = crate::routes::subscriptions::build_subscription_policies(old_schema);
let new_policies = crate::routes::subscriptions::build_subscription_policies(new_schema);
let changed = old_policies != new_policies;
if changed {
warn!(
old = old_policies.len(),
new = new_policies.len(),
"SECURITY: schema reload changes subscription row-visibility policies. New \
subscriptions pick up the change immediately; already-connected subscriptions \
re-derive on the reload signal — re-scoped in place, or terminated fail-closed \
(#611)."
);
}
changed
}
#[cfg(test)]
mod reload_policy_warn_tests {
use fraiseql_core::schema::{
CompiledSchema, SubscriptionDefinition, SubscriptionPolicy, TypeDefinition,
};
use super::warn_on_subscription_policy_reload;
fn schema(policy: Option<SubscriptionPolicy>) -> CompiledSchema {
let mut order = TypeDefinition::new("Order", "v_order");
if let Some(p) = policy {
order = order.with_subscription_policy(p);
}
CompiledSchema {
types: vec![order],
subscriptions: vec![SubscriptionDefinition::new("orderUpdated", "Order")],
..Default::default()
}
}
fn policy() -> SubscriptionPolicy {
SubscriptionPolicy {
owner_path: "$.owner_id".to_string(),
identity_field: "user_id".to_string(),
bypass_roles: vec![],
}
}
#[test]
fn adding_a_policy_on_reload_is_flagged() {
// #611: a policy added by a hot-reload is a change → operator-visible.
assert!(warn_on_subscription_policy_reload(&schema(None), &schema(Some(policy()))));
}
#[test]
fn an_unchanged_policy_set_is_not_flagged() {
assert!(!warn_on_subscription_policy_reload(
&schema(Some(policy())),
&schema(Some(policy()))
));
assert!(!warn_on_subscription_policy_reload(&schema(None), &schema(None)));
}
}