Skip to main content

issuerd_server/
state.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (C) 2026 Dmitry Andreev. <da@issuerd.org>
3//
4// ServerState: shared server components (storage, cache, crypto, token service, plugins).
5
6use std::collections::HashMap;
7use std::sync::Arc;
8
9use issuerd_auth_flow::plugin_registry::PluginRegistry;
10use issuerd_core::{
11    Authenticator, Credential, CredentialId, CredentialType, CryptoProvider, DisplayName,
12    DistributedCache, Email, EmailSender, EventListener, FederationManager, FlowConfig,
13    IssuerdError, Realm, RealmId, Storage, TokenService, User, UserId, Username,
14};
15use issuerd_token::token_manager::{TokenIssuer, TokenManager};
16
17use crate::config::ServerConfig;
18use tracing::{debug, info, instrument};
19
20#[derive(Clone)]
21pub struct ServerState {
22    pub config: ServerConfig,
23    pub storage: Arc<dyn Storage>,
24    pub cache: Arc<dyn DistributedCache>,
25    pub crypto: Arc<dyn CryptoProvider>,
26    pub token_service: Arc<dyn TokenService>,
27    pub token_manager: Arc<dyn TokenIssuer>,
28    pub plugin_registry: Arc<dyn PluginRegistry>,
29    pub federation_manager: Arc<dyn FederationManager>,
30    pub login_failure_tracker: Arc<issuerd_auth_flow::login_failures::LoginFailureTracker>,
31    pub email_sender: Arc<dyn EmailSender>,
32    /// Server-to-server client for identity brokering. A trait so
33    /// tests can stub the external IdP; production uses `ReqwestBrokerClient`.
34    pub broker_client: Arc<dyn issuerd_core::BrokerClient>,
35    /// Back-channel logout dispatcher: notifies clients with a
36    /// `backchannel_logout_uri` whenever one of their sessions is destroyed.
37    /// Fire-and-forget; `NoOpSessionLogoutNotifier` in tests that do not
38    /// exercise logout delivery.
39    pub logout_notifier: Arc<dyn issuerd_core::SessionLogoutNotifier>,
40    /// Same-node signing-key reload hook: invoked by the admin
41    /// API after key rotation/disable so the running crypto provider and
42    /// JWKS snapshot pick up the storage mutation immediately instead of
43    /// waiting for the next JWKS polling tick. Spawns the reload onto the
44    /// runtime; no-op in tests that do not wire the concrete provider.
45    pub signing_key_reload: Arc<dyn Fn() + Send + Sync>,
46    /// Monotonic counter bumped every time this node reloads the shared
47    /// signing-key set (admin reload hook, JWKS poll on change). The cached
48    /// discovery document embeds the signing-algorithm list, so rendered
49    /// discovery entries validate against it.
50    pub keyset_generation: Arc<std::sync::atomic::AtomicU64>,
51    /// Event listeners addressable by name from a realm's `events_listeners`.
52    /// Always contains `"logging"`; unknown realm listener names
53    /// are ignored at dispatch time.
54    pub event_listeners: HashMap<String, Arc<dyn EventListener>>,
55}
56
57impl ServerState {
58    /// Build a `TypedFlowExecutor` from the current plugin registry.
59    pub fn typed_executor(&self) -> issuerd_auth_flow::typestate::TypedFlowExecutor {
60        self.typed_executor_with_flows(&[])
61    }
62
63    /// Build a `TypedFlowExecutor` whose executor resolves sub-flow stages
64    /// against `flows` (pass the realm's full stored flow set).
65    pub fn typed_executor_with_flows(
66        &self,
67        flows: &[FlowConfig],
68    ) -> issuerd_auth_flow::typestate::TypedFlowExecutor {
69        issuerd_auth_flow::typestate::TypedFlowExecutor::new(
70            issuerd_auth_flow::executor::FlowExecutor::new(self.plugin_registry.clone(), flows),
71        )
72    }
73
74    /// Resolve the realm's bound top-level flow and an executor that can run
75    /// it: the flow set is loaded from storage per request — one
76    /// indexed query on Postgres, trivial on the in-memory/json backends, and
77    /// deliberately uncached so admin flow edits take effect immediately. The
78    /// realm's binding field (`bound`, `None` → `default_alias`) picks the
79    /// top-level flow; when the read fails (logged) or the alias is not in
80    /// storage, the `fallback` code constructor keeps bootstrapping and legacy
81    /// rigs working. The executor is built over the full loaded list so
82    /// `sub_flow_alias` stages resolve.
83    ///
84    /// Only the browser and registration bindings are consumed at runtime
85    /// (here and in the login/registration routes). The direct-grant,
86    /// reset-credentials, and first-broker-login bindings exist for
87    /// Keycloak-representation parity and admin delete-guardrails only — those
88    /// paths never run the flow engine.
89    pub async fn bound_flow_executor(
90        &self,
91        realm: &Realm,
92        bound: Option<&str>,
93        default_alias: &str,
94        fallback: fn(RealmId) -> FlowConfig,
95    ) -> (FlowConfig, issuerd_auth_flow::typestate::TypedFlowExecutor) {
96        let flows = match self.storage.list_flow_configs(&realm.id).await {
97            Ok(flows) => flows,
98            Err(e) => {
99                tracing::warn!(
100                    realm = %realm.id,
101                    error = %e,
102                    "cannot load flow configs from storage; falling back to code default flow"
103                );
104                Vec::new()
105            }
106        };
107        let alias = bound.unwrap_or(default_alias);
108        let flow = flows
109            .iter()
110            .find(|f| f.top_level && f.alias.as_str() == alias)
111            .cloned()
112            .unwrap_or_else(|| fallback(realm.id.clone()));
113        (flow, self.typed_executor_with_flows(&flows))
114    }
115
116    /// Resolve the realm segment from a request URL to a realm.
117    ///
118    /// URLs carry the human-readable realm name — discovery documents and
119    /// token issuers embed the realm *name* (see
120    /// `TokenManager::issuer_for_realm`). Name lookup is therefore
121    /// deterministic; the id lookup remains as a fallback so old id-spelled
122    /// bookmarks/URLs keep resolving (deterministic if an admin ever names a
123    /// realm after another realm's id: the name wins). The by-name step goes
124    /// through the shared realm cache; the id fallback stays uncached.
125    pub async fn resolve_realm(&self, segment: &str) -> Result<Option<Realm>, IssuerdError> {
126        if let Some(realm) = self.realm_by_name_cached(segment).await? {
127            return Ok(Some(realm));
128        }
129        match RealmId::new(segment) {
130            Ok(id) => self.storage.get_realm(&id).await,
131            Err(_) => Ok(None),
132        }
133    }
134
135    /// Resolve the realm an `iss` claim belongs to.
136    ///
137    /// Issuers are realm-NAME based (`{issuer_url}/realms/{name}`). The
138    /// trailing `/realms/` segment is extracted and resolved by name; callers
139    /// then use `realm.id` for storage lookups. Tokens issued before the
140    /// name-based switch (id-spelled issuers) are intentionally NOT resolved
141    /// here — they are rejected.
142    pub async fn resolve_issuer_realm(&self, issuer: &str) -> Result<Option<Realm>, IssuerdError> {
143        match issuerd_core::typestate::extract_realm_from_issuer(issuer) {
144            Some(name) => self.realm_by_name_cached(name).await,
145            None => Ok(None),
146        }
147    }
148
149    /// Realm-by-name lookup with a cache-aside layer (TTL
150    /// `[cache] read_cache_ttl_secs`, default 60 s).
151    ///
152    /// Realm rows change rarely (admin PUT / events-config save) and every
153    /// token-bearing request needs this resolution, so the full `Realm` model
154    /// is cached as JSON under `realm-by-name:{name}`. Admin mutations delete
155    /// the key synchronously (see `issuerd-admin-api` realms/events routes), so the
156    /// TTL is only the fail-safe for writes that bypass the API (direct DB
157    /// edits, provisioning at boot). Cache errors degrade to the storage path
158    /// with a WARN — a cache outage must never fail a request. Negative
159    /// results are not cached: unknown issuers fail closed on every call.
160    async fn realm_by_name_cached(&self, name: &str) -> Result<Option<Realm>, IssuerdError> {
161        let ttl_secs = self.config.cache.read_cache_ttl_secs;
162        if ttl_secs == 0 {
163            // Read-model caches disabled: pure-DB behavior.
164            return self.storage.get_realm_by_name(name).await;
165        }
166        let key = issuerd_cluster::cache_keys::realm_by_name(name);
167        match self.cache.get(&key).await {
168            Ok(Some(bytes)) => match serde_json::from_slice::<Realm>(&bytes) {
169                Ok(realm) => {
170                    debug!(realm = %realm.id, "realm-by-name cache: hit");
171                    return Ok(Some(realm));
172                }
173                Err(e) => {
174                    debug!(error = %e, "realm-by-name cache: malformed entry; re-reading");
175                }
176            },
177            Ok(None) => {}
178            Err(e) => {
179                tracing::warn!(error = %e, "realm-by-name cache read failed; falling back to storage");
180            }
181        }
182        let realm = self.storage.get_realm_by_name(name).await?;
183        if let Some(realm) = &realm {
184            match serde_json::to_vec(realm) {
185                Ok(bytes) => {
186                    if let Err(e) = self
187                        .cache
188                        .set(&key, bytes, Some(std::time::Duration::from_secs(ttl_secs)))
189                        .await
190                    {
191                        tracing::warn!(error = %e, "realm-by-name cache write failed");
192                    }
193                }
194                Err(e) => {
195                    tracing::warn!(error = %e, "realm-by-name cache serialization failed");
196                }
197            }
198        }
199        Ok(realm)
200    }
201
202    #[instrument(skip(config), err)]
203    pub async fn from_config(config: &ServerConfig) -> Result<Self, IssuerdError> {
204        // Validate up front: the master-realm bootstrap builds client redirect
205        // URIs from this value and must never panic on a bad config (P3-10).
206        match url::Url::parse(&config.issuer_url) {
207            Ok(u) if u.scheme() == "http" || u.scheme() == "https" => {}
208            _ => {
209                return Err(IssuerdError::InvalidRequest(format!(
210                    "issuer_url must be an absolute http(s) URL: {:?}",
211                    config.issuer_url
212                )));
213            }
214        }
215        // Bounded protocol values: reject out-of-range settings instead of
216        // silently clamping them (e.g. an operator typo of 6000 s).
217        config.oauth.validate()?;
218        config.dpop.validate()?;
219        // Log only the storage variant — the full config would leak the
220        // Postgres password into the logs.
221        let storage_kind = match &config.storage {
222            crate::config::StorageConfig::InMemory => "in-memory",
223            crate::config::StorageConfig::Postgres { .. } => "postgres",
224            crate::config::StorageConfig::JsonFile { .. } => "json-file",
225        };
226        debug!(
227            storage = storage_kind,
228            redis = config.redis.is_some(),
229            "initializing server state"
230        );
231
232        // Multi-node deployments must share both the persistent store
233        // (PostgreSQL: realms, users, sessions, signing keys) and the
234        // ephemeral cache (Redis: auth codes, pending auth, revocation
235        // blocklist, login failures). Anything else is a split-brain config
236        // where e.g. an auth code issued by node A is not redeemable on
237        // node B — refuse to boot. Validated before opening any connections.
238        if config.cluster.enabled {
239            let has_postgres =
240                matches!(config.storage, crate::config::StorageConfig::Postgres { .. });
241            let has_redis = config.redis.is_some() || !config.cluster.redis_nodes.is_empty();
242            if !has_postgres || !has_redis {
243                return Err(IssuerdError::ServerError(
244                    "cluster.enabled requires PostgreSQL storage and a Redis cache \
245                     (set `redis` or `cluster.redis_nodes`)"
246                        .to_string(),
247                ));
248            }
249        }
250        info!(
251            node_id = %config.cluster.resolved_node_id(),
252            cluster_enabled = config.cluster.enabled,
253            "node identity"
254        );
255
256        // Envelope encryption of signing keys at rest: validate and build the
257        // KEK provider up front so an invalid [crypto.key_encryption] section
258        // aborts boot before any connection is opened.
259        let kek = config.crypto.build_kek_provider()?;
260
261        // Storage
262        let storage: Arc<dyn Storage> = match &config.storage {
263            crate::config::StorageConfig::InMemory => {
264                if kek.is_some() {
265                    tracing::warn!(
266                        "[crypto.key_encryption] is only supported with PostgreSQL storage; \
267                         ignoring it for the in-memory backend"
268                    );
269                }
270                Arc::new(issuerd_storage::InMemoryStorage::new())
271            }
272            crate::config::StorageConfig::Postgres { url } => {
273                let pg =
274                    issuerd_storage::PostgresStorage::connect_with_key_encryption(url, kek.clone())
275                        .await?;
276                pg.run_migrations().await?;
277                if kek.is_some() {
278                    // Expand pass: legacy plaintext rows (and rows encrypted
279                    // under a previous KEK) are rewritten under the active KEK
280                    // before the keystore loads. A wrong KEK fails the daemon
281                    // here — never a plaintext fallback.
282                    let rewritten = pg.reencrypt_signing_keys_with_active_kek().await?;
283                    if rewritten > 0 {
284                        info!(count = rewritten, "re-encrypted signing keys at rest");
285                    }
286                } else {
287                    tracing::warn!(
288                        "signing keys are stored in PLAINTEXT in PostgreSQL; \
289                         configure [crypto.key_encryption] to encrypt them at rest"
290                    );
291                }
292                Arc::new(pg)
293            }
294            crate::config::StorageConfig::JsonFile { path } => {
295                if kek.is_some() {
296                    tracing::warn!(
297                        "[crypto.key_encryption] is only supported with PostgreSQL storage; \
298                         ignoring it for the JSON-file backend (protect the snapshot file at \
299                         the filesystem level)"
300                    );
301                }
302                Arc::new(issuerd_storage::JsonFileStorage::new(path)?)
303            }
304        };
305
306        // Cache: Redis Cluster node list overrides the single-node URL.
307        let cache: Arc<dyn DistributedCache> = if !config.cluster.redis_nodes.is_empty() {
308            Arc::new(
309                issuerd_cluster::RedisCache::connect_cluster(&config.cluster.redis_nodes).await?,
310            )
311        } else if let Some(redis_url) = &config.redis {
312            Arc::new(issuerd_cluster::RedisCache::connect(redis_url).await?)
313        } else {
314            Arc::new(issuerd_cluster::InMemoryCache::new())
315        };
316
317        Self::from_components(config, storage, cache).await
318    }
319
320    /// Build server state from pre-constructed shared components.
321    ///
322    /// [`ServerState::from_config`] constructs storage and cache from the
323    /// config and delegates here. Tests (and embeddings) call this directly to
324    /// share one storage/cache pair across several state instances, simulating
325    /// a multi-node cluster in one process. `config.storage` still decides the
326    /// master-realm bootstrap and JWKS refresh task eligibility.
327    pub async fn from_components(
328        config: &ServerConfig,
329        storage: Arc<dyn Storage>,
330        cache: Arc<dyn DistributedCache>,
331    ) -> Result<Self, IssuerdError> {
332        // Email: a real SMTP sender only when `[smtp] enabled = true`;
333        // otherwise the no-op sender fails loudly on use so a missing SMTP
334        // setup cannot silently swallow verification mails.
335        let email_sender: Arc<dyn EmailSender> = if config.smtp.enabled {
336            Arc::new(crate::email::SmtpEmailSender::new(config.smtp.clone()))
337        } else {
338            Arc::new(crate::email::NoOpEmailSender)
339        };
340        Self::from_components_with_email_sender(config, storage, cache, email_sender).await
341    }
342
343    /// [`ServerState::from_components`] with an explicit [`EmailSender`].
344    ///
345    /// Tests inject a recording sender here so both `state.email_sender` and
346    /// the plugin registry's email-code authenticator capture it — mutating
347    /// `state.email_sender` after construction would leave the registry's
348    /// mailer pointing at the previous sender.
349    pub async fn from_components_with_email_sender(
350        config: &ServerConfig,
351        storage: Arc<dyn Storage>,
352        cache: Arc<dyn DistributedCache>,
353        email_sender: Arc<dyn EmailSender>,
354    ) -> Result<Self, IssuerdError> {
355        // Crypto: load the shared signing-key set from storage so every node
356        // signs with the same active key and validates its peers' tokens.
357        let crypto = bootstrap_crypto_provider(storage.as_ref()).await?;
358        let jwks = crypto.get_public_keys().await?;
359
360        // TokenService. The default signing algorithm matches the bootstrapped
361        // provider config; realms may override it via their
362        // `default_signature_algorithm` attribute.
363        let token_manager = Arc::new(TokenManager::with_default_alg(
364            crypto.clone(),
365            config.issuer_url.clone(),
366            std::time::Duration::from_secs(60),
367            jwks,
368            issuerd_token::CryptoConfig::default().default_alg,
369        ));
370        // The constructor seeds both JWKS snapshots from the full published
371        // set; refresh immediately so the signing-selection snapshot becomes
372        // the true active-only set (a key disabled before this boot must
373        // never sign again, even before the first polling tick).
374        token_manager.refresh_jwks().await?;
375        let token_service: Arc<dyn TokenService> = token_manager.clone();
376
377        // Propagate keys added (or rotated) by peer nodes into this node's
378        // keystore and JWKS snapshot. Only PostgreSQL storage is shared
379        // between nodes, so polling is pointless for the other backends.
380        let keyset_generation = Arc::new(std::sync::atomic::AtomicU64::new(0));
381        if matches!(config.storage, crate::config::StorageConfig::Postgres { .. }) {
382            spawn_jwks_refresh_task(
383                storage.clone(),
384                crypto.clone(),
385                token_manager.clone(),
386                config.cluster.jwks_refresh_interval_secs,
387                keyset_generation.clone(),
388            );
389        }
390
391        let federation_manager: Arc<dyn FederationManager> = Arc::new(
392            issuerd_federation::DynamicFederationManager::new(storage.clone(), cache.clone()),
393        );
394
395        let login_failure_tracker =
396            Arc::new(issuerd_auth_flow::login_failures::LoginFailureTracker::new());
397
398        // Brokered login talks to external IdPs over HTTP.
399        let broker_client: Arc<dyn issuerd_core::BrokerClient> =
400            Arc::new(crate::broker::ReqwestBrokerClient::new()?);
401
402        let plugin_registry: Arc<dyn PluginRegistry> = {
403            let mut reg = SimplePluginRegistry::new();
404            reg.register_authenticator(Arc::new(
405                issuerd_auth_flow::built_in::CookieAuthenticator::new(
406                    token_service.clone(),
407                    storage.clone(),
408                ),
409            ));
410            reg.register_authenticator(Arc::new(
411                issuerd_auth_flow::built_in::UsernamePasswordAuthenticator::with_tracker(
412                    storage.clone(),
413                    login_failure_tracker.clone(),
414                    cache.clone(),
415                )
416                .with_federation_manager(federation_manager.clone()),
417            ));
418            reg.register_authenticator(Arc::new(
419                issuerd_auth_flow::built_in::SpnegoFlowAuthenticator::new(
420                    federation_manager.clone(),
421                ),
422            ));
423            reg.register_authenticator(Arc::new(
424                issuerd_auth_flow::built_in::OtpFormAuthenticator::new(storage.clone()),
425            ));
426            // Passwordless email one-time-code login (only factor when a
427            // realm's browser flow binds it): mails a numeric code via the
428            // realm-merged SMTP sender.
429            reg.register_authenticator(Arc::new(
430                issuerd_auth_flow::email_code::EmailCodeAuthenticator::new(
431                    storage.clone(),
432                    cache.clone(),
433                    Arc::new(crate::email::SmtpEmailCodeSender::new(email_sender.clone())),
434                ),
435            ));
436            reg.register_authenticator(Arc::new(
437                issuerd_auth_flow::built_in::ConditionalUserConfiguredAuthenticator::new(
438                    storage.clone(),
439                ),
440            ));
441            // WebAuthn second factor: the relying party is derived
442            // from the configured issuer URL (validated absolute http(s) URL
443            // before this point).
444            let (rp_id, rp_origin) =
445                issuerd_auth_flow::webauthn::relying_party_from_issuer(&config.issuer_url)
446                    .expect("issuer_url validated as absolute http(s) URL at startup");
447            reg.register_authenticator(Arc::new(
448                issuerd_auth_flow::webauthn::WebAuthnAuthenticator::new(
449                    storage.clone(),
450                    cache.clone(),
451                    rp_id,
452                    rp_origin,
453                ),
454            ));
455            reg.register_authenticator(Arc::new(
456                issuerd_auth_flow::built_in::IdentityProviderRedirectAuthenticator::new(
457                    storage.clone(),
458                ),
459            ));
460            reg.register_authenticator(Arc::new(
461                issuerd_auth_flow::built_in::RegistrationAuthenticator::new(storage.clone()),
462            ));
463            // Required actions must be registered or the flow engine cannot
464            // resolve them and VERIFY_EMAIL / UPDATE_PASSWORD are never
465            // enforced.
466            reg.register_required_action(Arc::new(
467                issuerd_auth_flow::built_in::VerifyEmailRequiredAction::new(storage.clone()),
468            ));
469            reg.register_required_action(Arc::new(
470                issuerd_auth_flow::built_in::UpdatePasswordRequiredAction::new(storage.clone())
471                    .with_federation_manager(federation_manager.clone()),
472            ));
473            reg.register_required_action(Arc::new(
474                issuerd_auth_flow::built_in::UpdateProfileRequiredAction::new(storage.clone()),
475            ));
476            reg.register_required_action(Arc::new(
477                issuerd_auth_flow::built_in::ConfigureTotpRequiredAction::new(storage.clone()),
478            ));
479            reg.register_required_action(Arc::new(
480                issuerd_auth_flow::built_in::TermsAndConditionsRequiredAction::new(storage.clone()),
481            ));
482            Arc::new(reg)
483        };
484
485        let logout_notifier: Arc<dyn issuerd_core::SessionLogoutNotifier> =
486            Arc::new(crate::routes::logout::BackchannelLogoutDispatcher::new(
487                storage.clone(),
488                token_manager.clone(),
489            )?);
490
491        // Same-node signing-key reload hook: re-read the shared
492        // key set from storage, reload the keystore, and refresh the cached
493        // JWKS snapshot — exactly what the JWKS polling task does per reload.
494        // Built here while the concrete `RingCryptoProvider` is still in
495        // reach (the `crypto` field below erases it to `dyn CryptoProvider`).
496        let signing_key_reload: Arc<dyn Fn() + Send + Sync> = {
497            let storage = storage.clone();
498            let crypto = crypto.clone();
499            let token_manager = token_manager.clone();
500            let keyset_generation = keyset_generation.clone();
501            Arc::new(move || {
502                let storage = storage.clone();
503                let crypto = crypto.clone();
504                let token_manager = token_manager.clone();
505                let keyset_generation = keyset_generation.clone();
506                tokio::spawn(async move {
507                    let keys = match storage.list_signing_keys().await {
508                        Ok(keys) => keys,
509                        Err(e) => {
510                            tracing::warn!(error = %e, "signing-key reload: cannot read signing keys");
511                            return;
512                        }
513                    };
514                    if let Err(e) = crypto.reload_keys(&keys) {
515                        tracing::warn!(error = %e, "signing-key reload: keystore reload failed");
516                        return;
517                    }
518                    if let Err(e) = token_manager.refresh_jwks().await {
519                        tracing::warn!(error = %e, "signing-key reload: JWKS snapshot reload failed");
520                        return;
521                    }
522                    keyset_generation.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
523                });
524            })
525        };
526
527        // Event listeners: the built-in "logging" listener is
528        // always available; realm `events_listeners` entries resolve against
529        // this map at dispatch time.
530        let event_listeners: HashMap<String, Arc<dyn EventListener>> = HashMap::from([(
531            "logging".to_string(),
532            Arc::new(crate::event_listeners::LoggingEventListener) as Arc<dyn EventListener>,
533        )]);
534
535        let state = Self {
536            config: config.clone(),
537            storage,
538            cache,
539            crypto,
540            token_service,
541            token_manager,
542            plugin_registry,
543            federation_manager,
544            login_failure_tracker,
545            email_sender,
546            broker_client,
547            logout_notifier,
548            signing_key_reload,
549            keyset_generation,
550            event_listeners,
551        };
552
553        // Pre-populate master realm for in-memory/json storage if no realms exist
554        #[allow(clippy::match_like_matches_macro)]
555        let storage_is_bootstrapable = match config.storage {
556            crate::config::StorageConfig::InMemory
557            | crate::config::StorageConfig::JsonFile { .. } => true,
558            _ => false,
559        };
560        let no_realms = state
561            .storage
562            .list_realms(&issuerd_core::Pagination::default())
563            .await
564            .unwrap_or_default()
565            .is_empty();
566        if storage_is_bootstrapable && no_realms {
567            #[cfg(coverage)]
568            {
569                let _ = bootstrap_master_realm(&state).await;
570            }
571            #[cfg(not(coverage))]
572            {
573                if let Err(e) = bootstrap_master_realm(&state).await {
574                    tracing::error!(error = %e, "failed to bootstrap master realm");
575                }
576            }
577        }
578
579        // Apply optional provision config exactly once.
580        if let Some(provision_path) = &config.provision {
581            match crate::provisioner::Provisioner::from_file(provision_path) {
582                Ok(provisioner) => {
583                    if let Err(e) = provisioner
584                        .apply_once(state.storage.as_ref(), &state.config.issuer_url)
585                        .await
586                    {
587                        tracing::error!(error = %e, "provision failed");
588                    }
589                }
590                Err(e) => {
591                    tracing::error!(error = %e, "failed to load provision config");
592                }
593            }
594        }
595
596        Ok(state)
597    }
598}
599
600/// Build the crypto provider from the shared signing-key set in storage.
601///
602/// On first boot (empty key set) the initial key set is generated and
603/// persisted so every node — and every restart — converges on the same keys.
604/// The initial set is a PAIR: the server-default EdDSA key (the signing
605/// default for realms without a `default_signature_algorithm` attribute) and
606/// an active RS256 key — OIDC Core §15.1 makes RS256 mandatory-to-implement,
607/// so discovery must advertise it from the start and realms explicitly pinned
608/// to RS256 work without a rotation. A concurrent first boot may persist a
609/// second pair; that is benign: all keys are published in JWKS, all validate,
610/// and all nodes sign with the newest active key of the resolved algorithm.
611/// Deployments upgrading with an existing key set keep their stored keys
612/// untouched.
613async fn bootstrap_crypto_provider(
614    storage: &dyn Storage,
615) -> Result<Arc<issuerd_token::RingCryptoProvider>, IssuerdError> {
616    let crypto_config = issuerd_token::CryptoConfig::default();
617    let mut keys = storage.list_signing_keys().await?;
618    if keys.is_empty() {
619        for generated in issuerd_token::KeyStore::generate_initial_key_set(
620            crypto_config.default_alg,
621            crypto_config.rsa_key_size,
622        )? {
623            let stored = generated.to_stored(true);
624            storage.create_signing_key(&stored).await?;
625            info!(kid = %stored.kid, alg = %stored.alg, "generated and persisted initial signing key");
626        }
627        // Re-list so keys persisted by a concurrently booting node are picked up.
628        keys = storage.list_signing_keys().await?;
629    }
630    let provider = issuerd_token::RingCryptoProvider::from_signing_keys(crypto_config, &keys)?;
631    Ok(Arc::new(provider))
632}
633
634/// Periodically re-read the shared signing-key set from storage and, when it
635/// changed, reload the keystore and the JWKS snapshot so keys added or rotated
636/// by peer nodes propagate to this node without a restart.
637fn spawn_jwks_refresh_task(
638    storage: Arc<dyn Storage>,
639    crypto: Arc<issuerd_token::RingCryptoProvider>,
640    token_manager: Arc<TokenManager<issuerd_token::RingCryptoProvider>>,
641    interval_secs: u64,
642    keyset_generation: Arc<std::sync::atomic::AtomicU64>,
643) {
644    tokio::spawn(async move {
645        let mut interval =
646            tokio::time::interval(std::time::Duration::from_secs(interval_secs.max(1)));
647        // Skip the immediate first tick — keys were just loaded at boot.
648        interval.tick().await;
649        loop {
650            interval.tick().await;
651            let keys = match storage.list_signing_keys().await {
652                Ok(keys) => keys,
653                Err(e) => {
654                    tracing::warn!(error = %e, "JWKS refresh: cannot read signing keys");
655                    continue;
656                }
657            };
658            let current_kids: Vec<String> = match crypto.get_public_keys().await {
659                Ok(set) => {
660                    let mut kids: Vec<String> =
661                        set.keys.iter().map(|k| k.kid.to_string()).collect();
662                    kids.sort();
663                    kids
664                }
665                Err(e) => {
666                    tracing::warn!(error = %e, "JWKS refresh: cannot read current JWKS");
667                    continue;
668                }
669            };
670            let mut stored_kids: Vec<String> = keys.iter().map(|k| k.kid.to_string()).collect();
671            stored_kids.sort();
672            if stored_kids == current_kids {
673                continue;
674            }
675            if let Err(e) = crypto.reload_keys(&keys) {
676                tracing::warn!(error = %e, "JWKS refresh: keystore reload failed");
677                continue;
678            }
679            if let Err(e) = token_manager.refresh_jwks().await {
680                tracing::warn!(error = %e, "JWKS refresh: JWKS snapshot reload failed");
681                continue;
682            }
683            keyset_generation.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
684            info!(keys = keys.len(), "reloaded signing keys from shared storage");
685        }
686    });
687}
688
689async fn bootstrap_master_realm(state: &ServerState) -> Result<(), IssuerdError> {
690    // Use the configured issuer_url as the base for built-in client redirect URIs
691    // so in-memory/json bootstraps match the actual server URL (e.g. https).
692    let base_url = state.config.issuer_url.trim_end_matches('/').to_string();
693
694    let realm = Realm {
695        id: RealmId::new("master").unwrap(),
696        name: issuerd_core::RealmName::new("master").unwrap(),
697        display_name: Some(DisplayName::new("Master").unwrap()),
698        enabled: true,
699        ..Default::default()
700    };
701    state.storage.create_realm(&realm).await?;
702
703    let admin_user = User {
704        id: UserId::new("admin").unwrap(),
705        realm_id: RealmId::new("master").unwrap(),
706        username: Username::new("admin").expect("admin username must be valid"),
707        email: Some(Email::new("admin@localhost.local").expect("admin email must be valid")),
708        email_verified: true,
709        first_name: Some(DisplayName::new("Admin").unwrap()),
710        last_name: Some(DisplayName::new("User").unwrap()),
711        enabled: true,
712        federation_link: None,
713        attributes: HashMap::new(),
714        required_actions: Vec::new(),
715        created_at: chrono::Utc::now(),
716        updated_at: chrono::Utc::now(),
717    };
718    state.storage.create_user(&realm.id, &admin_user).await?;
719
720    // Hash password with argon2
721    use argon2::{password_hash::SaltString, Argon2, PasswordHasher};
722    use rand::rngs::OsRng;
723    let salt = SaltString::generate(&mut OsRng);
724    let argon2 = Argon2::default();
725    let hash = argon2
726        .hash_password("admin".as_bytes(), &salt)
727        .map_err(|e| IssuerdError::ServerError(format!("hash failed: {e}")))?
728        .to_string();
729
730    let cred = Credential {
731        id: CredentialId::new(issuerd_core::utils::generate_id()).unwrap(),
732        credential_type: CredentialType::Password,
733        user_label: Some("Password".to_string()),
734        created_date: chrono::Utc::now(),
735        secret_data: hash.into_bytes(),
736        credential_data: serde_json::json!({"hash_algorithm": "argon2id"}),
737        priority: 1,
738    };
739    state.storage.create_credential(&realm.id, &admin_user.id, &cred).await?;
740
741    // Create the built-in clients every realm gets (Keycloak parity):
742    // admin-cli for admin-style logins and account-console for the Account SPA.
743    let admin_cli = issuerd_admin_api::realms::build_admin_cli_client(&realm.id, &base_url)?;
744    state.storage.create_client(&realm.id, &admin_cli).await?;
745    let account_client =
746        issuerd_admin_api::realms::build_account_console_client(&realm.id, "master", &base_url)?;
747    state.storage.create_client(&realm.id, &account_client).await?;
748
749    // Create standard realm-management roles for the admin API
750    let admin_roles = vec![
751        "manage-realm",
752        "view-realm",
753        "manage-users",
754        "view-users",
755        "manage-clients",
756        "view-clients",
757        // Required by the admin impersonation endpoint.
758        "impersonation",
759    ];
760    for role_name in &admin_roles {
761        let role = issuerd_core::Role {
762            id: issuerd_core::RoleId::new(issuerd_core::utils::generate_id()).unwrap(),
763            name: issuerd_core::RoleName::new(role_name.to_string()).unwrap(),
764            description: Some(format!("{role_name} role")),
765            realm_id: realm.id.clone(),
766            client_role: false,
767            client_id: None,
768            composite: false,
769            composites: vec![],
770            attributes: HashMap::new(),
771        };
772        if let Err(e) = state.storage.create_role(&realm.id, &role).await {
773            tracing::warn!(error = %e, role = %role.name, "master realm bootstrap: failed to create role");
774        }
775        if let Err(e) = state.storage.add_user_realm_role(&realm.id, &admin_user.id, &role.id).await
776        {
777            tracing::warn!(error = %e, role = %role.name, "master realm bootstrap: failed to grant role to admin user");
778        }
779    }
780
781    // The built-in browser and registration flows are seeded by the storage
782    // layer's `create_realm` hook
783    // (`issuerd_storage::seed::seed_builtin_flows`), which also makes them appear in
784    // the admin API.
785    // TODO: Keycloak minimal parity — also create on realm bootstrap:
786    //   - direct grant
787    //   - reset credentials
788    //   - first broker login
789    //   - first login flow
790
791    info!("bootstrapped master realm with default admin user — change the password immediately");
792    Ok(())
793}
794
795// ---------------------------------------------------------------------------
796// SimplePluginRegistry
797// ---------------------------------------------------------------------------
798
799pub struct SimplePluginRegistry {
800    authenticators: HashMap<String, Arc<dyn Authenticator>>,
801    required_actions: HashMap<String, Arc<dyn issuerd_core::RequiredAction>>,
802}
803
804impl Default for SimplePluginRegistry {
805    fn default() -> Self {
806        Self::new()
807    }
808}
809
810impl SimplePluginRegistry {
811    pub fn new() -> Self {
812        Self {
813            authenticators: HashMap::new(),
814            required_actions: HashMap::new(),
815        }
816    }
817
818    pub fn register_authenticator(&mut self, auth: Arc<dyn Authenticator>) {
819        self.authenticators.insert(auth.id().to_string(), auth);
820    }
821
822    pub fn register_required_action(&mut self, action: Arc<dyn issuerd_core::RequiredAction>) {
823        self.required_actions.insert(action.id().to_string(), action);
824    }
825}
826
827#[async_trait::async_trait]
828impl PluginRegistry for SimplePluginRegistry {
829    async fn get_authenticator(
830        &self,
831        id: &str,
832    ) -> Result<Option<Arc<dyn Authenticator>>, IssuerdError> {
833        Ok(self.authenticators.get(id).cloned())
834    }
835
836    async fn get_required_action(
837        &self,
838        id: &str,
839    ) -> Result<Option<Arc<dyn issuerd_core::RequiredAction>>, IssuerdError> {
840        Ok(self.required_actions.get(id).cloned())
841    }
842
843    fn list_authenticator_ids(&self) -> Vec<String> {
844        self.authenticators.keys().cloned().collect()
845    }
846
847    fn list_required_action_ids(&self) -> Vec<String> {
848        self.required_actions.keys().cloned().collect()
849    }
850}
851
852// ---------------------------------------------------------------------------
853// Default browser flow for MVP
854// ---------------------------------------------------------------------------
855
856// The built-in flow constructors moved to `issuerd_core::flows` so the
857// storage layer can seed them on realm creation. Re-exported here to keep the
858// existing `crate::state::{default_browser_flow, registration_flow}` imports
859// working.
860pub use issuerd_core::flows::{default_browser_flow, registration_flow};
861
862#[cfg(test)]
863mod tests {
864    use super::*;
865    use issuerd_core::{FlowStageId, Requirement};
866
867    #[tokio::test]
868    async fn server_state_from_config_inmemory() {
869        let cfg = ServerConfig::default();
870        let state = ServerState::from_config(&cfg).await.unwrap();
871        let realms = state.storage.list_realms(&issuerd_core::Pagination::default()).await.unwrap();
872        assert!(!realms.is_empty());
873        assert_eq!(realms[0].name, "master");
874    }
875
876    #[tokio::test]
877    async fn from_config_rejects_out_of_range_auth_code_ttl() {
878        for bad in [0, 9, 601] {
879            let cfg = ServerConfig {
880                oauth: crate::config::OAuthConfig {
881                    auth_code_ttl_secs: bad,
882                },
883                ..Default::default()
884            };
885            let err = match ServerState::from_config(&cfg).await {
886                Ok(_) => panic!("out-of-range auth_code_ttl_secs ({bad}) must abort boot"),
887                Err(e) => e,
888            };
889            assert!(
890                err.to_string().contains("auth_code_ttl_secs"),
891                "error must name the offending key: {err}"
892            );
893        }
894    }
895
896    #[tokio::test]
897    async fn server_state_builds_federation_manager() {
898        let cfg = ServerConfig::default();
899        let state = ServerState::from_config(&cfg).await.unwrap();
900        // Federation manager should be non-null and functional
901        let providers = state
902            .federation_manager
903            .providers_for_realm(&issuerd_core::RealmId::new("master").unwrap())
904            .await
905            .unwrap();
906        // Master realm has no federation providers configured, so list should be empty
907        assert!(providers.is_empty());
908    }
909
910    #[tokio::test]
911    async fn simple_plugin_registry() {
912        let mut reg = SimplePluginRegistry::new();
913        assert!(reg.list_authenticator_ids().is_empty());
914        assert!(reg.list_required_action_ids().is_empty());
915
916        let auth = Arc::new(issuerd_auth_flow::built_in::UsernamePasswordAuthenticator::new(
917            Arc::new(issuerd_storage::InMemoryStorage::new()),
918        ));
919        reg.register_authenticator(auth.clone());
920        assert_eq!(reg.list_authenticator_ids(), vec!["auth-username-password"]);
921        let fetched = reg.get_authenticator("auth-username-password").await.unwrap();
922        assert!(fetched.is_some());
923
924        let action = Arc::new(issuerd_auth_flow::built_in::VerifyEmailRequiredAction::new(
925            Arc::new(issuerd_storage::InMemoryStorage::new()),
926        ));
927        reg.register_required_action(action.clone());
928        assert_eq!(reg.list_required_action_ids(), vec!["VERIFY_EMAIL"]);
929        let fetched = reg.get_required_action("VERIFY_EMAIL").await.unwrap();
930        assert!(fetched.is_some());
931    }
932
933    #[test]
934    fn default_browser_flow_structure() {
935        let flow = default_browser_flow(RealmId::new("test").unwrap());
936        assert_eq!(flow.alias, issuerd_core::Alias::new("browser").unwrap());
937        assert_eq!(flow.stages.len(), 7);
938        assert_eq!(flow.stages[0].id, FlowStageId::new("cookie-auth").unwrap());
939        assert_eq!(flow.stages[1].id, FlowStageId::new("auth-spnego").unwrap());
940        // The kc_idp_hint broker redirect runs after cookie/SPNEGO so an
941        // existing SSO session still wins over the hint.
942        assert_eq!(flow.stages[2].id, FlowStageId::new("idp-redirect").unwrap());
943        assert_eq!(flow.stages[2].requirement, Requirement::Alternative);
944        assert_eq!(flow.stages[3].id, FlowStageId::new("username-password").unwrap());
945        // Conditional TOTP/WebAuthn second factor at the flow tail.
946        assert_eq!(flow.stages[4].id, FlowStageId::new("conditional-user-configured").unwrap());
947        assert_eq!(flow.stages[4].requirement, Requirement::Conditional);
948        assert_eq!(flow.stages[5].id, FlowStageId::new("auth-otp-form").unwrap());
949        assert_eq!(flow.stages[5].requirement, Requirement::Conditional);
950        assert_eq!(flow.stages[6].id, FlowStageId::new("auth-webauthn").unwrap());
951        // Optional, NOT Conditional — see the stage comment: a Conditional
952        // WebAuthn stage would be scope-skipped after an Attempted OTP stage.
953        assert_eq!(flow.stages[6].requirement, Requirement::Optional);
954    }
955
956    #[test]
957    fn registration_flow_structure() {
958        let flow = registration_flow(RealmId::new("test").unwrap());
959        assert_eq!(flow.alias, issuerd_core::Alias::new("registration").unwrap());
960        assert_eq!(flow.stages.len(), 1);
961        assert_eq!(flow.stages[0].id, FlowStageId::new("registration").unwrap());
962        assert_eq!(
963            flow.stages[0].authenticator,
964            issuerd_core::Alias::new("auth-registration").unwrap()
965        );
966        assert_eq!(flow.stages[0].requirement, Requirement::Required);
967    }
968
969    #[tokio::test]
970    async fn server_state_from_config_json_file() {
971        let dir = std::env::temp_dir()
972            .join(format!("issuerd_server_json_test_{}", issuerd_core::utils::generate_id()));
973        let _ = std::fs::create_dir_all(&dir);
974        let path = dir.join("db.json");
975        let cfg = ServerConfig {
976            storage: crate::config::StorageConfig::JsonFile { path: path.clone() },
977            ..Default::default()
978        };
979        let state = ServerState::from_config(&cfg).await.unwrap();
980        let realms = state.storage.list_realms(&issuerd_core::Pagination::default()).await.unwrap();
981        assert!(!realms.is_empty());
982        let _ = std::fs::remove_dir_all(&dir);
983    }
984
985    #[test]
986    fn simple_plugin_registry_default() {
987        let reg = SimplePluginRegistry::default();
988        assert!(reg.list_authenticator_ids().is_empty());
989        assert!(reg.list_required_action_ids().is_empty());
990    }
991
992    #[tokio::test]
993    async fn bootstrap_master_realm_create_realm_fails() {
994        let mut mock_storage = issuerd_core::MockStorage::new();
995        mock_storage
996            .expect_create_realm()
997            .returning(|_| Err(issuerd_core::IssuerdError::ServerError("db fail".to_string())));
998        mock_storage.expect_list_realms().returning(|_| Ok(vec![]));
999
1000        let state = ServerState {
1001            config: ServerConfig::default(),
1002            storage: Arc::new(mock_storage),
1003            cache: Arc::new(issuerd_cluster::InMemoryCache::new()),
1004            crypto: Arc::new(
1005                issuerd_token::RingCryptoProvider::new(issuerd_token::CryptoConfig::default())
1006                    .unwrap(),
1007            ),
1008            token_service: Arc::new(issuerd_token::token_manager::TokenManager::new(
1009                Arc::new(
1010                    issuerd_token::RingCryptoProvider::new(issuerd_token::CryptoConfig::default())
1011                        .unwrap(),
1012                ),
1013                "http://localhost:8080".to_string(),
1014                std::time::Duration::from_secs(60),
1015                issuerd_core::JwkSet { keys: vec![] },
1016            )),
1017            token_manager: Arc::new(issuerd_token::token_manager::TokenManager::new(
1018                Arc::new(
1019                    issuerd_token::RingCryptoProvider::new(issuerd_token::CryptoConfig::default())
1020                        .unwrap(),
1021                ),
1022                "http://localhost:8080".to_string(),
1023                std::time::Duration::from_secs(60),
1024                issuerd_core::JwkSet { keys: vec![] },
1025            )),
1026            plugin_registry: Arc::new(SimplePluginRegistry::new()),
1027            federation_manager: Arc::new(issuerd_federation::NoOpFederationManager),
1028            login_failure_tracker: Arc::new(
1029                issuerd_auth_flow::login_failures::LoginFailureTracker::new(),
1030            ),
1031            email_sender: Arc::new(crate::email::NoOpEmailSender),
1032            broker_client: Arc::new(issuerd_core::MockBrokerClient::new()),
1033            logout_notifier: Arc::new(issuerd_core::NoOpSessionLogoutNotifier),
1034            signing_key_reload: Arc::new(|| {}),
1035            keyset_generation: Arc::new(std::sync::atomic::AtomicU64::new(0)),
1036            event_listeners: HashMap::new(),
1037        };
1038
1039        let result = bootstrap_master_realm(&state).await;
1040        assert!(result.is_err());
1041    }
1042
1043    #[tokio::test]
1044    async fn server_state_from_config_postgres_invalid_url_fails() {
1045        let cfg = ServerConfig {
1046            storage: crate::config::StorageConfig::Postgres {
1047                url: "postgres://invalid_host:5432/db".to_string(),
1048            },
1049            ..Default::default()
1050        };
1051        let result = ServerState::from_config(&cfg).await;
1052        assert!(result.is_err());
1053    }
1054
1055    #[tokio::test]
1056    async fn server_state_from_config_redis_invalid_url_fails() {
1057        let cfg = ServerConfig {
1058            redis: Some("redis://invalid_host:6379".to_string()),
1059            ..Default::default()
1060        };
1061        let result = ServerState::from_config(&cfg).await;
1062        assert!(result.is_err());
1063    }
1064
1065    #[tokio::test]
1066    async fn cluster_enabled_requires_postgres_storage() {
1067        // In-memory storage + no Redis → boot must fail fast.
1068        let cfg = ServerConfig {
1069            cluster: crate::config::ClusterConfig {
1070                enabled: true,
1071                ..Default::default()
1072            },
1073            ..Default::default()
1074        };
1075        let result = ServerState::from_config(&cfg).await;
1076        let msg = format!("{}", result.err().unwrap());
1077        assert!(msg.contains("cluster.enabled"), "unexpected error: {msg}");
1078    }
1079
1080    #[tokio::test]
1081    async fn cluster_enabled_requires_redis_cache() {
1082        // Postgres configured but no Redis → boot must fail before connecting.
1083        let cfg = ServerConfig {
1084            storage: crate::config::StorageConfig::Postgres {
1085                url: "postgres://invalid_host:5432/db".to_string(),
1086            },
1087            cluster: crate::config::ClusterConfig {
1088                enabled: true,
1089                ..Default::default()
1090            },
1091            ..Default::default()
1092        };
1093        let result = ServerState::from_config(&cfg).await;
1094        let msg = format!("{}", result.err().unwrap());
1095        assert!(msg.contains("cluster.enabled"), "unexpected error: {msg}");
1096    }
1097
1098    #[tokio::test]
1099    async fn cluster_enabled_with_redis_fails_on_connect_not_validation() {
1100        // Valid cluster shape (Postgres + Redis): validation passes, boot then
1101        // fails on the unreachable Postgres — a different error.
1102        let cfg = ServerConfig {
1103            storage: crate::config::StorageConfig::Postgres {
1104                url: "postgres://invalid_host:5432/db".to_string(),
1105            },
1106            redis: Some("redis://invalid_host:6379".to_string()),
1107            cluster: crate::config::ClusterConfig {
1108                enabled: true,
1109                ..Default::default()
1110            },
1111            ..Default::default()
1112        };
1113        let result = ServerState::from_config(&cfg).await;
1114        let msg = format!("{}", result.err().unwrap());
1115        assert!(!msg.contains("cluster.enabled"), "validation should have passed: {msg}");
1116    }
1117
1118    #[tokio::test]
1119    async fn inmemory_boot_persists_signing_key_in_storage() {
1120        let cfg = ServerConfig::default();
1121        let state = ServerState::from_config(&cfg).await.unwrap();
1122        let keys = state.storage.list_signing_keys().await.unwrap();
1123        // Fresh boot persists the initial PAIR: the EdDSA default signing key
1124        // and the OIDC Core MTI RS256 key, both active.
1125        assert_eq!(keys.len(), 2);
1126        assert!(keys.iter().all(|k| k.active));
1127        let algs: Vec<_> = keys.iter().map(|k| k.alg).collect();
1128        assert!(algs.contains(&issuerd_core::Algorithm::EdDsa));
1129        assert!(algs.contains(&issuerd_core::Algorithm::Rs256));
1130        // The published JWKS must contain the same keys.
1131        let jwks = state.crypto.get_public_keys().await.unwrap();
1132        assert!(keys.iter().all(|k| jwks.keys.iter().any(|j| j.kid == k.kid)));
1133    }
1134
1135    #[tokio::test]
1136    async fn server_state_from_config_skips_bootstrap_when_realms_exist() {
1137        let mut mock_storage = issuerd_core::MockStorage::new();
1138        mock_storage.expect_list_realms().returning(|_| {
1139            Ok(vec![issuerd_core::Realm {
1140                id: issuerd_core::RealmId::new("existing").unwrap(),
1141                name: issuerd_core::RealmName::new("existing").unwrap(),
1142                display_name: None,
1143                enabled: true,
1144                ..Default::default()
1145            }])
1146        });
1147
1148        let state = ServerState {
1149            config: ServerConfig::default(),
1150            storage: Arc::new(mock_storage),
1151            cache: Arc::new(issuerd_cluster::InMemoryCache::new()),
1152            crypto: Arc::new(
1153                issuerd_token::RingCryptoProvider::new(issuerd_token::CryptoConfig::default())
1154                    .unwrap(),
1155            ),
1156            token_service: Arc::new(issuerd_token::token_manager::TokenManager::new(
1157                Arc::new(
1158                    issuerd_token::RingCryptoProvider::new(issuerd_token::CryptoConfig::default())
1159                        .unwrap(),
1160                ),
1161                "http://localhost:8080".to_string(),
1162                std::time::Duration::from_secs(60),
1163                issuerd_core::JwkSet { keys: vec![] },
1164            )),
1165            token_manager: Arc::new(issuerd_token::token_manager::TokenManager::new(
1166                Arc::new(
1167                    issuerd_token::RingCryptoProvider::new(issuerd_token::CryptoConfig::default())
1168                        .unwrap(),
1169                ),
1170                "http://localhost:8080".to_string(),
1171                std::time::Duration::from_secs(60),
1172                issuerd_core::JwkSet { keys: vec![] },
1173            )),
1174            plugin_registry: Arc::new(SimplePluginRegistry::new()),
1175            federation_manager: Arc::new(issuerd_federation::NoOpFederationManager),
1176            login_failure_tracker: Arc::new(
1177                issuerd_auth_flow::login_failures::LoginFailureTracker::new(),
1178            ),
1179            email_sender: Arc::new(crate::email::NoOpEmailSender),
1180            broker_client: Arc::new(issuerd_core::MockBrokerClient::new()),
1181            logout_notifier: Arc::new(issuerd_core::NoOpSessionLogoutNotifier),
1182            signing_key_reload: Arc::new(|| {}),
1183            keyset_generation: Arc::new(std::sync::atomic::AtomicU64::new(0)),
1184            event_listeners: HashMap::new(),
1185        };
1186
1187        // bootstrap_master_realm should not be called because realms exist
1188        let realms = state.storage.list_realms(&issuerd_core::Pagination::default()).await.unwrap();
1189        assert_eq!(realms.len(), 1);
1190    }
1191
1192    #[tokio::test]
1193    async fn bound_flow_executor_loads_storage_flows_with_code_fallback() {
1194        let cfg = ServerConfig::default();
1195        let state = ServerState::from_config(&cfg).await.unwrap();
1196        let realm = state
1197            .storage
1198            .get_realm(&RealmId::new("master").unwrap())
1199            .await
1200            .unwrap()
1201            .unwrap();
1202
1203        // Realm creation seeded the builtin flows into storage.
1204        let stored = state.storage.list_flow_configs(&realm.id).await.unwrap();
1205        assert!(stored.iter().any(|f| f.top_level && f.alias.as_str() == "browser"));
1206
1207        // A realm binding to a custom stored flow wins over the code default.
1208        let mut custom = registration_flow(realm.id.clone());
1209        custom.alias = issuerd_core::Alias::new("custom-top").unwrap();
1210        custom.built_in = false;
1211        state.storage.create_flow_config(&realm.id, &custom).await.unwrap();
1212        let (flow, _executor) = state
1213            .bound_flow_executor(&realm, Some("custom-top"), "browser", default_browser_flow)
1214            .await;
1215        assert_eq!(flow.alias.as_str(), "custom-top");
1216        assert_eq!(flow.stages.len(), 1, "stored flow, not the 7-stage code default");
1217
1218        // An unbound realm (None) resolves the default alias from storage; an
1219        // unknown binding falls back to the code constructor.
1220        let (flow, _executor) =
1221            state.bound_flow_executor(&realm, None, "browser", default_browser_flow).await;
1222        assert_eq!(flow.alias.as_str(), "browser");
1223        let (flow, _executor) = state
1224            .bound_flow_executor(&realm, Some("missing"), "browser", default_browser_flow)
1225            .await;
1226        assert_eq!(flow.alias.as_str(), "browser");
1227        assert_eq!(flow.stages.len(), 7, "code default browser flow as fallback");
1228    }
1229
1230    #[tokio::test]
1231    async fn bootstrap_seeds_and_assigns_impersonation_role() {
1232        let cfg = ServerConfig::default();
1233        let state = ServerState::from_config(&cfg).await.unwrap();
1234        let realm_id = RealmId::new("master").unwrap();
1235        let role = state
1236            .storage
1237            .get_role_by_name(&realm_id, "impersonation")
1238            .await
1239            .unwrap()
1240            .expect("impersonation role must be seeded at bootstrap");
1241        assert!(!role.client_role);
1242
1243        let assigned = state
1244            .storage
1245            .list_user_realm_roles(&realm_id, &UserId::new("admin").unwrap())
1246            .await
1247            .unwrap();
1248        assert!(assigned.contains(&role.id), "bootstrap admin holds the impersonation role");
1249    }
1250
1251    #[tokio::test]
1252    async fn resolve_issuer_realm_caches_and_reflects_invalidation() {
1253        let cfg = ServerConfig::default();
1254        let state = ServerState::from_config(&cfg).await.unwrap();
1255        let issuer = "http://localhost:8080/realms/master";
1256        let key = issuerd_cluster::cache_keys::realm_by_name("master");
1257
1258        let realm = state
1259            .resolve_issuer_realm(issuer)
1260            .await
1261            .unwrap()
1262            .expect("master realm resolves");
1263        assert!(
1264            state.cache.get(&key).await.unwrap().is_some(),
1265            "first resolve populates the realm-by-name cache"
1266        );
1267
1268        let mut updated = realm.clone();
1269        updated.display_name = Some(DisplayName::new("Mutated").unwrap());
1270        state.storage.update_realm(&updated).await.unwrap();
1271
1272        let cached = state.resolve_issuer_realm(issuer).await.unwrap().unwrap();
1273        assert_eq!(
1274            cached.display_name, realm.display_name,
1275            "second resolve is served from cache and still shows the old value"
1276        );
1277
1278        state.cache.delete(&key).await.unwrap();
1279        let fresh = state.resolve_issuer_realm(issuer).await.unwrap().unwrap();
1280        assert_eq!(
1281            fresh.display_name, updated.display_name,
1282            "after invalidation the resolve reflects the storage change"
1283        );
1284    }
1285
1286    #[tokio::test]
1287    async fn jwks_refresh_task_reloads_only_on_key_set_changes() {
1288        use std::sync::atomic::Ordering;
1289
1290        let storage: Arc<dyn issuerd_core::Storage> =
1291            Arc::new(issuerd_storage::InMemoryStorage::new());
1292        let crypto_config = issuerd_token::CryptoConfig::default();
1293        let mut keys = Vec::new();
1294        for generated in issuerd_token::KeyStore::generate_initial_key_set(
1295            crypto_config.default_alg,
1296            crypto_config.rsa_key_size,
1297        )
1298        .unwrap()
1299        {
1300            let stored = generated.to_stored(true);
1301            storage.create_signing_key(&stored).await.unwrap();
1302            keys.push(stored);
1303        }
1304        let crypto = Arc::new(
1305            issuerd_token::RingCryptoProvider::from_signing_keys(crypto_config, &keys).unwrap(),
1306        );
1307        let jwks = crypto.get_public_keys().await.unwrap();
1308        let token_manager = Arc::new(TokenManager::with_default_alg(
1309            crypto.clone(),
1310            "http://localhost:8080".to_string(),
1311            std::time::Duration::from_secs(60),
1312            jwks,
1313            issuerd_token::CryptoConfig::default().default_alg,
1314        ));
1315        let generation = Arc::new(std::sync::atomic::AtomicU64::new(0));
1316        spawn_jwks_refresh_task(
1317            storage.clone(),
1318            crypto.clone(),
1319            token_manager,
1320            1,
1321            generation.clone(),
1322        );
1323
1324        // Multiple polling intervals with an unchanged key set: no reload may
1325        // fire (a reload every tick would also invalidate discovery caches
1326        // cluster-wide for no reason).
1327        tokio::time::sleep(std::time::Duration::from_millis(2300)).await;
1328        assert_eq!(
1329            generation.load(Ordering::Relaxed),
1330            0,
1331            "an unchanged key set must never trigger a reload"
1332        );
1333
1334        // A key persisted by a peer node propagates on a subsequent tick.
1335        let extra =
1336            issuerd_token::KeyStore::generate_key(issuerd_core::Algorithm::EdDsa, 2048).unwrap();
1337        let extra_kid = extra.kid.to_string();
1338        storage.create_signing_key(&extra.to_stored(true)).await.unwrap();
1339        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(4);
1340        while generation.load(Ordering::Relaxed) == 0 && std::time::Instant::now() < deadline {
1341            tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1342        }
1343        assert!(
1344            generation.load(Ordering::Relaxed) >= 1,
1345            "a changed key set must trigger a reload"
1346        );
1347        let kids: Vec<String> = crypto
1348            .get_public_keys()
1349            .await
1350            .unwrap()
1351            .keys
1352            .iter()
1353            .map(|k| k.kid.to_string())
1354            .collect();
1355        assert!(kids.contains(&extra_kid), "keystore must reload with the new key");
1356    }
1357}