1use 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 pub broker_client: Arc<dyn issuerd_core::BrokerClient>,
35 pub logout_notifier: Arc<dyn issuerd_core::SessionLogoutNotifier>,
40 pub signing_key_reload: Arc<dyn Fn() + Send + Sync>,
46 pub keyset_generation: Arc<std::sync::atomic::AtomicU64>,
51 pub event_listeners: HashMap<String, Arc<dyn EventListener>>,
55}
56
57impl ServerState {
58 pub fn typed_executor(&self) -> issuerd_auth_flow::typestate::TypedFlowExecutor {
60 self.typed_executor_with_flows(&[])
61 }
62
63 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 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 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 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 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 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 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 config.oauth.validate()?;
218 config.dpop.validate()?;
219 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 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 let kek = config.crypto.build_kek_provider()?;
260
261 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 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 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 pub async fn from_components(
328 config: &ServerConfig,
329 storage: Arc<dyn Storage>,
330 cache: Arc<dyn DistributedCache>,
331 ) -> Result<Self, IssuerdError> {
332 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 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 let crypto = bootstrap_crypto_provider(storage.as_ref()).await?;
358 let jwks = crypto.get_public_keys().await?;
359
360 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 token_manager.refresh_jwks().await?;
375 let token_service: Arc<dyn TokenService> = token_manager.clone();
376
377 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 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 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 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 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 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 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 #[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 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
600async 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 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
634fn 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 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 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 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 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 let admin_roles = vec![
751 "manage-realm",
752 "view-realm",
753 "manage-users",
754 "view-users",
755 "manage-clients",
756 "view-clients",
757 "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 info!("bootstrapped master realm with default admin user — change the password immediately");
792 Ok(())
793}
794
795pub 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
852pub 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 let providers = state
902 .federation_manager
903 .providers_for_realm(&issuerd_core::RealmId::new("master").unwrap())
904 .await
905 .unwrap();
906 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}