From 57d68435fe23ebd7556715cc250bec4bab97ddbf Mon Sep 17 00:00:00 2001 From: Martin Hoffmann Date: Mon, 17 Nov 2025 18:16:17 +0100 Subject: [PATCH] Introduce a StorageSystem for global storage system state. --- src/api/ca.rs | 9 +- src/bin/krillup.rs | 40 ++++- src/cli/ta/signer.rs | 9 +- .../crypto/signing/dispatch/krillsigner.rs | 31 ++-- .../crypto/signing/dispatch/signerinfo.rs | 8 +- .../crypto/signing/signers/softsigner.rs | 46 +++--- src/commons/eventsourcing/store.rs | 17 +- src/commons/eventsourcing/wal.rs | 11 +- src/commons/queue.rs | 52 +++--- src/commons/storage/mod.rs | 5 +- src/commons/storage/store.rs | 122 +++++++++++--- src/commons/storage/test.rs | 30 ++-- src/commons/test.rs | 13 +- src/config.rs | 19 +-- src/daemon/http/auth/authorizer.rs | 10 +- src/daemon/http/auth/crypt.rs | 7 +- src/daemon/http/auth/providers/config_file.rs | 10 +- .../auth/providers/openid_connect/provider.rs | 8 +- src/daemon/http/server.rs | 7 +- src/daemon/start.rs | 17 +- src/server/ca/manager.rs | 22 +-- src/server/ca/publishing.rs | 11 +- src/server/ca/status.rs | 7 +- src/server/ca/upgrades/data_migration.rs | 13 +- .../ca/upgrades/pre_0_10_0/migration.rs | 33 ++-- .../ca/upgrades/pre_0_14_0/migration.rs | 17 +- src/server/manager.rs | 66 ++++---- src/server/mq.rs | 7 +- src/server/properties/mod.rs | 20 +-- src/server/pubd/access.rs | 13 +- src/server/pubd/content.rs | 7 +- src/server/pubd/manager.rs | 33 ++-- src/server/pubd/upgrades/mod.rs | 26 ++- .../pubd/upgrades/pre_0_10_0/migration.rs | 17 +- src/server/scheduler.rs | 65 ++------ src/tasigner/config.rs | 11 +- src/upgrades/data_migration.rs | 85 +++++----- src/upgrades/mod.rs | 151 +++++++++--------- src/upgrades/pre_0_14_0.rs | 17 +- 39 files changed, 592 insertions(+), 500 deletions(-) diff --git a/src/api/ca.rs b/src/api/ca.rs index b35b8607..447cc543 100644 --- a/src/api/ca.rs +++ b/src/api/ca.rs @@ -2230,7 +2230,7 @@ mod test { use bytes::Bytes; use rpki::crypto::PublicKeyFormat; use crate::api::ta::TrustAnchorLocator; - use crate::commons::crypto::OpenSslSigner; + use crate::commons::crypto::{OpenSslSigner, OpenSslSignerConfig}; use crate::commons::test; use super::*; @@ -2254,9 +2254,10 @@ mod test { #[test] fn mft_uri() { - test::test_in_memory(|storage_uri| { - let signer = - OpenSslSigner::build(storage_uri, "dummy", None).unwrap(); + test::test_in_memory(|storage| { + let signer = OpenSslSigner::build( + storage, &OpenSslSignerConfig::default(), "dummy", None + ).unwrap(); let key_id = signer.create_key(PublicKeyFormat::Rsa).unwrap(); let pub_key = signer.get_key_info(&key_id).unwrap(); diff --git a/src/bin/krillup.rs b/src/bin/krillup.rs index 965cf0ec..c0b4475e 100644 --- a/src/bin/krillup.rs +++ b/src/bin/krillup.rs @@ -7,6 +7,7 @@ use log::info; use log::LevelFilter; use url::Url; use krill::constants; +use krill::commons::storage::StorageSystem; use krill::config::{Config, LogType}; use krill::server::properties::PropertiesManager; use krill::upgrades::{prepare_upgrade_data_migrations, UpgradeMode}; @@ -33,9 +34,28 @@ fn main() { match options.command { Command::Prepare(_prepare) => { + let storage = match StorageSystem::new( + config.storage_uri.clone() + ) { + Ok(storage) => storage, + Err(err) => { + eprintln!("*** Error Preparing Data Migration ***"); + eprintln!("Cannot connect to storage system: {err}"); + eprintln!(); + eprintln!( + "Note that your server data has NOT been modified. \ + Do not upgrade krill" + ); + eprintln!("itself yet!"); + eprintln!( + "If you did upgrade krill, then downgrade it to \ + the previous installed version." + ); + ::std::process::exit(1); + } + }; let properties_manager = match PropertiesManager::create( - &config.storage_uri, - config.use_history_cache, + &storage, config.use_history_cache, ) { Ok(mgr) => mgr, Err(e) => { @@ -57,6 +77,7 @@ fn main() { match prepare_upgrade_data_migrations( UpgradeMode::PrepareOnly, + &storage, &config, &properties_manager, ) { @@ -93,7 +114,20 @@ fn main() { } } Command::Migrate(cmd) => { - if let Err(e) = migrate(config, cmd.target) { + let storage = match StorageSystem::new(cmd.target) { + Ok(storage) => storage, + Err(err) => { + eprintln!("*** Error Migrating DATA ***"); + eprintln!("Cannot connect to storage system: {err}"); + eprintln!( + "Note that your server data has NOT been modified." + ); + eprintln!(); + ::std::process::exit(1); + } + }; + + if let Err(e) = migrate(config, &storage) { eprintln!("*** Error Migrating DATA ***"); eprintln!("{e}"); eprintln!(); diff --git a/src/cli/ta/signer.rs b/src/cli/ta/signer.rs index bd05d76d..5c56dfa6 100644 --- a/src/cli/ta/signer.rs +++ b/src/cli/ta/signer.rs @@ -15,7 +15,7 @@ use crate::commons::crypto::KrillSigner; use crate::commons::actor::Actor; use crate::commons::error::Error as KrillError; use crate::commons::eventsourcing::{AggregateStore, AggregateStoreError}; -use crate::commons::storage::Ident; +use crate::commons::storage::{Ident, StorageSystem}; use crate::commons::httpclient; use crate::tasigner::{ Config, TrustAnchorProxySignerExchanges, @@ -122,13 +122,16 @@ pub struct TrustAnchorSignerManager { impl TrustAnchorSignerManager { pub fn create(config: Config) -> Result { + let storage = StorageSystem::new( + config.storage_uri.clone() + ).map_err(KrillError::from)?; let store = AggregateStore::create( - &config.storage_uri, + &storage, const { Ident::make("signer") }, config.use_history_cache, ).map_err(SignerClientError::other)?; let ta_handle = CaHandle::new("ta".into()); - let signer = config.signer()?; + let signer = config.signer(&storage)?; let actor = crate::constants::ACTOR_DEF_KRILLTA; Ok(TrustAnchorSignerManager { diff --git a/src/commons/crypto/signing/dispatch/krillsigner.rs b/src/commons/crypto/signing/dispatch/krillsigner.rs index 4d5a161e..cf38b8b6 100644 --- a/src/commons/crypto/signing/dispatch/krillsigner.rs +++ b/src/commons/crypto/signing/dispatch/krillsigner.rs @@ -26,7 +26,6 @@ use rpki::{ Cert, Crl, Manifest, Roa, }, }; -use url::Url; use crate::{ commons::{ @@ -40,6 +39,7 @@ use crate::{ CryptoResult, OpenSslSigner, SignSupport, }, error::Error, + storage::StorageSystem, KrillResult, }, constants::ID_CERTIFICATE_VALIDITY_YEARS, @@ -83,7 +83,7 @@ use crate::commons::crypto::{ type SignerBuilderFn = fn( &SignerType, SignerFlags, - &Url, + &StorageSystem, &str, std::time::Duration, &Option>, @@ -91,7 +91,7 @@ type SignerBuilderFn = fn( #[derive(Debug)] pub struct KrillSignerBuilder<'a> { - storage_uri: Url, + storage: &'a StorageSystem, probe_interval: Duration, signer_configs: &'a [SignerConfig], default_signer: Option<&'a SignerConfig>, @@ -100,12 +100,12 @@ pub struct KrillSignerBuilder<'a> { impl<'a> KrillSignerBuilder<'a> { pub fn new( - storage_uri: &Url, + storage: &'a StorageSystem, probe_interval: Duration, signer_configs: &'a [SignerConfig], ) -> Self { Self { - storage_uri: storage_uri.clone(), + storage, probe_interval, signer_configs, default_signer: None, @@ -174,7 +174,7 @@ impl<'a> KrillSignerBuilder<'a> { } KrillSigner::build( - &self.storage_uri, + self.storage, self.probe_interval, self.signer_configs, default_signer, @@ -190,7 +190,7 @@ pub struct KrillSigner { impl KrillSigner { fn build( - storage_uri: &Url, + storage: &StorageSystem, probe_interval: Duration, signer_configs: &[SignerConfig], default_signer: &SignerConfig, @@ -199,10 +199,10 @@ impl KrillSigner { #[cfg(not(feature = "hsm"))] let signer_mapper = None; #[cfg(feature = "hsm")] - let signer_mapper = Some(Arc::new(SignerMapper::build(storage_uri)?)); + let signer_mapper = Some(Arc::new(SignerMapper::build(storage)?)); let signers = Self::build_signers( signer_builder, - storage_uri, + storage, probe_interval, &signer_mapper, signer_configs, @@ -419,7 +419,7 @@ impl KrillSigner { impl KrillSigner { fn build_signers( signer_builder: SignerBuilderFn, - storage_uri: &Url, + storage: &StorageSystem, probe_interval: std::time::Duration, mapper: &Option>, configs: &[SignerConfig], @@ -449,7 +449,7 @@ impl KrillSigner { let signer = (signer_builder)( &config.signer_type, flags, - storage_uri, + storage, &config.name, probe_interval, mapper, @@ -465,7 +465,7 @@ impl KrillSigner { fn signer_builder( r#type: &SignerType, flags: SignerFlags, - storage_uri: &Url, + storage: &StorageSystem, name: &str, #[cfg(feature = "hsm")] probe_interval: Duration, #[cfg(not(feature = "hsm"))] _probe_interval: Duration, @@ -473,10 +473,9 @@ fn signer_builder( ) -> KrillResult { match r#type { SignerType::OpenSsl(conf) => { - let storage_uri = - conf.keys_storage_uri.as_ref().unwrap_or(storage_uri); - let signer = - OpenSslSigner::build(storage_uri, name, mapper.clone())?; + let signer = OpenSslSigner::build( + storage, conf, name, mapper.clone() + )?; Ok(SignerProvider::OpenSsl(flags, signer)) } #[cfg(feature = "hsm")] diff --git a/src/commons/crypto/signing/dispatch/signerinfo.rs b/src/commons/crypto/signing/dispatch/signerinfo.rs index fe109b3a..5c858e4b 100644 --- a/src/commons/crypto/signing/dispatch/signerinfo.rs +++ b/src/commons/crypto/signing/dispatch/signerinfo.rs @@ -8,7 +8,6 @@ use rpki::{ crypto::{KeyIdentifier, PublicKey}, }; use serde::{Deserialize, Serialize}; -use url::Url; use crate::{ commons::{ @@ -19,6 +18,7 @@ use crate::{ InitCommandDetails, InitEvent, SentCommand, SentInitCommand, WithStorableDetails, }, + storage::StorageSystem, KrillResult, }, constants::{ACTOR_DEF_KRILL, SIGNERS_NS}, @@ -442,11 +442,9 @@ impl std::fmt::Debug for SignerMapper { impl SignerMapper { /// Build a SignerMapper that will read/write its data in a subdirectory /// of the given work dir. - pub fn build(storage_uri: &Url) -> KrillResult { + pub fn build(storage: &StorageSystem) -> KrillResult { let store = AggregateStore::::create( - storage_uri, - SIGNERS_NS, - true, + storage, SIGNERS_NS, true, )?; Ok(SignerMapper { store }) } diff --git a/src/commons/crypto/signing/signers/softsigner.rs b/src/commons/crypto/signing/signers/softsigner.rs index 3a7a5d6f..c9958f62 100644 --- a/src/commons/crypto/signing/signers/softsigner.rs +++ b/src/commons/crypto/signing/signers/softsigner.rs @@ -27,7 +27,7 @@ use crate::{ dispatch::signerinfo::SignerMapper, signers::error::SignerError, SignerHandle, }, - storage::{Ident, KeyValueStore}, + storage::{Ident, KeyValueStore, StorageSystem, OpenStoreError}, }, constants::KEYS_NS, }; @@ -69,18 +69,19 @@ impl OpenSslSigner { /// SignerMapper only knows about keys created by the OpenSslSigner if /// the OpenSslSigner registers the new keys in the mapper. pub fn build( - storage_uri: &Url, + storage: &StorageSystem, + conf: &OpenSslSignerConfig, name: &str, mapper: Option>, - ) -> Result { - let keys_store = Self::init_keys_store(storage_uri)?; + ) -> Result { + let keys_store = Self::init_keys_store(storage, conf)?; let s = OpenSslSigner { name: name.to_string(), info: Some(format!( "OpenSSL Soft Signer [version: {}, keys store: {}]", openssl::version::version(), - storage_uri, + storage.default_uri(), )), handle: RwLock::new(None), // will be set later mapper, @@ -137,11 +138,13 @@ impl OpenSslSigner { impl OpenSslSigner { fn init_keys_store( - storage_uri: &Url, - ) -> Result { - let store = KeyValueStore::create(storage_uri, KEYS_NS) - .map_err(|_| SignerError::InvalidStorage(storage_uri.clone()))?; - Ok(store) + storage: &StorageSystem, + conf: &OpenSslSignerConfig, + ) -> Result { + match &conf.keys_storage_uri { + Some(uri) => storage.open_uri(uri, KEYS_NS), + None => storage.open(KEYS_NS) + } } fn build_key(&self) -> Result { @@ -385,10 +388,19 @@ pub mod tests { use super::*; + fn build_signer(storage: &StorageSystem) -> OpenSslSigner { + OpenSslSigner::build( + storage, + &OpenSslSignerConfig::default(), + "dummy", + None + ).unwrap() + } + #[test] fn should_return_subject_public_key_info() { - test::test_in_memory(|storage_uri| { - let s = OpenSslSigner::build(storage_uri, "dummy", None).unwrap(); + test::test_in_memory(|storage| { + let s = build_signer(storage); let ki = s.create_key(PublicKeyFormat::Rsa).unwrap(); s.get_key_info(&ki).unwrap(); s.destroy_key(&ki).unwrap(); @@ -410,14 +422,13 @@ pub mod tests { #[test] fn import_existing_pkcs1_openssl_key() { - test::test_in_memory(|storage_uri| { + test::test_in_memory(|storage| { // The following key was generated using OpenSSL on the command // line let pem = include_str!( "../../../../../test-resources/ta/example-pkcs1.pem" ); - let signer = - OpenSslSigner::build(storage_uri, "dummy", None).unwrap(); + let signer = build_signer(storage); let ki = signer.import_key(pem).unwrap(); signer.get_key_info(&ki).unwrap(); @@ -427,14 +438,13 @@ pub mod tests { #[test] fn import_existing_pkcs8_openssl_key() { - test::test_in_memory(|storage_uri| { + test::test_in_memory(|storage| { // The following key was generated using OpenSSL on the command // line let pem = include_str!( "../../../../../test-resources/ta/example-pkcs8.pem" ); - let signer = - OpenSslSigner::build(storage_uri, "dummy", None).unwrap(); + let signer = build_signer(storage); let ki = signer.import_key(pem).unwrap(); signer.get_key_info(&ki).unwrap(); diff --git a/src/commons/eventsourcing/store.rs b/src/commons/eventsourcing/store.rs index 0caf9317..2f80169c 100644 --- a/src/commons/eventsourcing/store.rs +++ b/src/commons/eventsourcing/store.rs @@ -12,12 +12,13 @@ use rpki::ca::idexchange::MyHandle; use rpki::repository::x509::Time; use serde::Serialize; use serde::de::DeserializeOwned; -use url::Url; use crate::api::history::{ CommandHistory, CommandHistoryCriteria, CommandHistoryRecord }; use crate::commons::error::KrillIoError; -use crate::commons::storage::{Ident, KeyValueError, KeyValueStore}; +use crate::commons::storage::{ + Ident, KeyValueError, KeyValueStore, OpenStoreError, StorageSystem +}; use super::agg::{ Aggregate, Command, InitCommand, PostSaveEventListener, PreSaveEventListener, StoredCommand @@ -76,12 +77,12 @@ impl AggregateStore { /// If `use_history_cache` is `true`, the new store will cache any /// history cache record created for any instance. pub fn create( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, use_history_cache: bool, - ) -> Result { + ) -> Result { Ok(Self::create_from_kv( - KeyValueStore::create(storage_uri, namespace)?, use_history_cache + storage.open(namespace)?, use_history_cache )) } @@ -90,12 +91,12 @@ impl AggregateStore { /// If `use_history_cache` is `true`, the new store will cache any /// history cache record created for any instance. pub fn create_upgrade_store( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, use_history_cache: bool, - ) -> Result { + ) -> Result { Ok(Self::create_from_kv( - KeyValueStore::create_upgrade_store(storage_uri, namespace)?, + storage.open_upgrade(namespace)?, use_history_cache, )) } diff --git a/src/commons/eventsourcing/wal.rs b/src/commons/eventsourcing/wal.rs index 2b951da4..322aaa1a 100644 --- a/src/commons/eventsourcing/wal.rs +++ b/src/commons/eventsourcing/wal.rs @@ -10,8 +10,9 @@ use std::sync::{Arc, RwLock}; use log::{error, warn, trace}; use rpki::ca::idexchange::MyHandle; use serde::{Deserialize, Serialize}; -use url::Url; -use crate::commons::storage::{Ident, KeyValueError, KeyValueStore}; +use crate::commons::storage::{ + Ident, KeyValueError, KeyValueStore, OpenStoreError, StorageSystem, +}; use super::store::Storable; @@ -135,11 +136,11 @@ pub struct WalStore { impl WalStore { /// Creates a new store using the given storage URL and namespace. pub fn create( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, - ) -> Result { + ) -> Result { Ok(WalStore { - kv: KeyValueStore::create(storage_uri, namespace)?, + kv: storage.open(namespace)?, cache: RwLock::new(HashMap::new()), }) } diff --git a/src/commons/queue.rs b/src/commons/queue.rs index 3a552f54..90329f23 100644 --- a/src/commons/queue.rs +++ b/src/commons/queue.rs @@ -2,9 +2,9 @@ use std::{cmp, error, fmt}; use std::time::{Duration, SystemTime, UNIX_EPOCH}; -use url::Url; use crate::commons::storage::{ - Ident, KeyValueError, KeyValueStore, Transaction + Ident, KeyValueError, KeyValueStore, OpenStoreError, StorageSystem, + Transaction, }; //------------ Configuration ------------------------------------------------- @@ -51,11 +51,11 @@ impl Queue { impl Queue { /// Creates a new queue. pub fn create( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, - ) -> Result { + ) -> Result { Ok(Queue { - store: KeyValueStore::create(storage_uri, namespace)?, + store: storage.open(namespace)?, }) } @@ -457,11 +457,12 @@ mod tests { use url::Url; use super::*; - fn queue_store(ns: &str) -> Queue { - Queue::create( - &Url::parse("memory://").unwrap(), - Ident::from_str(ns).unwrap() - ).unwrap() + fn storage_system() -> StorageSystem { + StorageSystem::new(Url::parse("memory://").unwrap()).unwrap() + } + + fn queue_store(storage: &StorageSystem, ns: &str) -> Queue { + Queue::create(storage, Ident::from_str(ns).unwrap()).unwrap() } #[test] @@ -481,12 +482,13 @@ mod tests { #[test] fn queue_thread_workers() { - let queue = queue_store("queue_thread_workers"); + let storage = storage_system(); + let queue = queue_store(&storage, "queue_thread_workers"); queue.store.wipe().unwrap(); thread::scope(|s| { - let create = s.spawn(|| { - let queue = queue_store("queue_thread_workers"); + s.spawn(|| { + let queue = queue_store(&storage, "queue_thread_workers"); for i in 1..=10 { let name = Ident::builder( @@ -503,16 +505,17 @@ mod tests { println!("> Scheduled job {}", &name); } }); - create.join().unwrap(); + }); - let keys = queue.store.execute(None, |tran| { - tran.list_keys(Queue::pending_scope()) - }).unwrap(); - assert_eq!(keys.len(), 10); + let keys = queue.store.execute(None, |tran| { + tran.list_keys(Queue::pending_scope()) + }).unwrap(); + assert_eq!(keys.len(), 10); + thread::scope(|s| { for _ in 1..=10 { - s.spawn(move || { - let queue = queue_store("queue_thread_workers"); + s.spawn(|| { + let queue = queue_store(&storage, "queue_thread_workers"); while queue.pending_tasks_remaining().unwrap() > 0 { if let Some((task_name, _)) @@ -537,7 +540,8 @@ mod tests { #[test] fn test_reschedule_long_running() { - let queue = queue_store("test_reschedule_long_running"); + let storage = storage_system(); + let queue = queue_store(&storage, "test_reschedule_long_running"); queue.store.wipe().unwrap(); let name = const { Ident::make("job") }; @@ -572,7 +576,8 @@ mod tests { #[test] fn test_reschedule_finished_task() { - let queue = queue_store("test_reschedule_finished_task"); + let storage = storage_system(); + let queue = queue_store(&storage, "test_reschedule_finished_task"); queue.store.wipe().unwrap(); let name = const { Ident::make("task") }; @@ -613,7 +618,8 @@ mod tests { #[test] fn test_schedule_with_existing_task() { - let queue = queue_store("test_schedule_with_existing_task"); + let storage = storage_system(); + let queue = queue_store(&storage, "test_schedule_with_existing_task"); queue.store.wipe().unwrap(); let name = const { Ident::make("task") }; diff --git a/src/commons/storage/mod.rs b/src/commons/storage/mod.rs index f54be028..0aa092f9 100644 --- a/src/commons/storage/mod.rs +++ b/src/commons/storage/mod.rs @@ -2,7 +2,10 @@ pub use self::backends::{Backend, Transaction, Error}; pub use self::ident::{Ident, IdentBuilder, IdentError}; -pub use self::store::{KeyValueStore, KeyValueError}; +pub use self::store::{ + KeyValueStore, KeyValueError, OpenStoreError, StorageConnectError, + StorageSystem +}; mod backends; mod ident; diff --git a/src/commons/storage/store.rs b/src/commons/storage/store.rs index 9c62af56..00cba061 100644 --- a/src/commons/storage/store.rs +++ b/src/commons/storage/store.rs @@ -1,14 +1,73 @@ //! The publicly exposed key-value store. -use std::fmt; +use std::{error, fmt}; use serde::de::DeserializeOwned; use serde::ser::Serialize; use url::Url; - use crate::commons::storage; +use crate::commons::error::Error; use crate::commons::storage::{Backend, Ident, Transaction}; +//------------ StorageSystem ------------------------------------------------- + +/// The system that provides the key-value stores. +#[derive(Debug)] +pub struct StorageSystem { + storage_uri: Url, +} + +impl StorageSystem { + /// Creates a new storage system. + /// + /// The provided URI will be used as the default storage URI. + pub fn new( + storage_uri: Url + ) -> Result { + Ok(Self { storage_uri }) + } + + /// Opens the default store with the given namespace. + pub fn open( + &self, namespace: &Ident + ) -> Result { + KeyValueStore::create(&self.storage_uri, namespace).map_err( + OpenStoreError + ) + } + + /// Opens a store for upgrades for the given namespace. + /// + /// Prefixes the namespace with `"upgrade_"`. + pub fn open_upgrade( + &self, namespace: &Ident + ) -> Result { + KeyValueStore::create( + &self.storage_uri, + &KeyValueStore::prefixed_namespace( + namespace, const { Ident::make("upgrade") } + ) + ).map_err( + OpenStoreError + ) + } + + /// Opens a store with the given storage URI and namespace. + pub fn open_uri( + &self, storage_uri: &Url, namespace: &Ident + ) -> Result { + KeyValueStore::create(storage_uri, namespace).map_err( + OpenStoreError + ) + } + + /// Returns the default URI of the storage system. + pub fn default_uri(&self) -> &Url { + &self.storage_uri + } +} + + //------------ KeyValueStore ------------------------------------------------- /// A key-value store. @@ -28,7 +87,7 @@ pub struct KeyValueStore { impl KeyValueStore { /// Creates a new store. - pub fn create( + fn create( storage_uri: &Url, namespace: &Ident, ) -> Result { @@ -171,21 +230,6 @@ impl KeyValueStore { // # Migration Support impl KeyValueStore { - /// Creates a new KeyValueStore for upgrades. - /// - /// Adds the implicit prefix "upgrade_" to the given namespace. - pub fn create_upgrade_store( - storage_uri: &Url, - namespace: &Ident, - ) -> Result { - Self::create( - storage_uri, - &Self::prefixed_namespace( - namespace, const { Ident::make("upgrade") } - ) - ) - } - fn prefixed_namespace( namespace: &Ident, prefix: &Ident, @@ -269,6 +313,48 @@ impl KeyValueStore { } +//------------ StorageConnectError ------------------------------------------- + +/// An error occured while connecting to a storage system. +#[derive(Debug)] +pub struct StorageConnectError(()); + +impl fmt::Display for StorageConnectError { + fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { + f.write_str("storage connect error") + } +} + +impl error::Error for StorageConnectError { } + +impl From for crate::commons::error::Error { + fn from(src: StorageConnectError) -> Self { + Self::custom(format_args!("{}", src)) + } +} + + +//------------ OpenStoreError ------------------------------------------------ + +/// An error occured while opening a store. +#[derive(Debug)] +pub struct OpenStoreError(KeyValueError); + +impl fmt::Display for OpenStoreError{ + fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { + self.0.fmt(f) + } +} + +impl error::Error for OpenStoreError { } + +impl From for Error { + fn from(src: OpenStoreError) -> Error { + Error::custom(format_args!("{}", src)) + } +} + + //------------ KeyValueError ------------------------------------------------- /// This type defines possible Errors for KeyStore diff --git a/src/commons/storage/test.rs b/src/commons/storage/test.rs index a5a4ada2..c68b5b47 100644 --- a/src/commons/storage/test.rs +++ b/src/commons/storage/test.rs @@ -8,7 +8,7 @@ use std::sync::{Mutex, MutexGuard}; use lazy_static::lazy_static; use tempfile::{TempDir, tempdir}; use url::Url; -use super::{Ident, KeyValueStore}; +use super::{Ident, KeyValueStore, StorageSystem}; //------------ Macro to Construct Tests -------------------------------------- @@ -73,7 +73,7 @@ const CONTENT_4: u32 = 45; // All the test functions. // -// The all need to have the same signature taking one argument as an +// They all need to have the same signature taking one argument as an // `impl Harness` and return unit. Each function will be transformed into a // test function for each of the backends (currently memory and disk). The // harness will give it access to a temporary test store atop that given @@ -317,6 +317,7 @@ trait Harness { /// it clean before the test. This is why there are a lock and a guard here. struct MemoryHarness<'a> { _guard: MutexGuard<'a, ()>, + storage: StorageSystem, } lazy_static! { @@ -338,21 +339,23 @@ impl<'a> MemoryHarness<'a> { } }; super::backends::memory::Store::wipe_all(); - return Self { _guard } + return Self { + _guard, + storage: StorageSystem::new( + Url::parse("memory:").unwrap() + ).unwrap(), + } } } } impl<'a> Harness for MemoryHarness<'a> { fn url(&self) -> Url { - Url::parse("memory:").unwrap() + self.storage.default_uri().clone() } fn store(&self, namespace: &Ident) -> KeyValueStore { - KeyValueStore::create( - &Url::parse("memory:").unwrap(), - namespace, - ).unwrap() + self.storage.open(namespace).unwrap() } } @@ -365,25 +368,26 @@ impl<'a> Harness for MemoryHarness<'a> { /// removed automatically when the harness is dropped. struct DiskHarness { _dir: TempDir, - url: Url, + storage: StorageSystem, } impl DiskHarness { fn new() -> Self { let _dir = tempdir().unwrap(); let url = format!("local://{}", _dir.path().display()); - let url = Url::parse(&url).unwrap(); - Self { _dir, url } + let storage = StorageSystem::new(Url::parse(&url).unwrap()).unwrap(); + + Self { _dir, storage } } } impl Harness for DiskHarness { fn url(&self) -> Url { - self.url.clone() + self.storage.default_uri().clone() } fn store(&self, namespace: &Ident) -> KeyValueStore { - KeyValueStore::create(&self.url, namespace).unwrap() + self.storage.open(namespace).unwrap() } } diff --git a/src/commons/test.rs b/src/commons/test.rs index 931a13d6..916387a0 100644 --- a/src/commons/test.rs +++ b/src/commons/test.rs @@ -8,17 +8,16 @@ use rpki::uri; use rpki::ca::idcert::IdCert; use url::Url; use crate::api::roa::{ConfiguredRoa, RoaConfiguration, RoaPayload}; +use crate::commons::storage::StorageSystem; /// This method returns an in-memory Key-Value store and then runs the test /// provided in the closure using it pub fn test_in_memory(op: F) where - F: FnOnce(&Url), + F: FnOnce(&StorageSystem), { - let storage_uri = mem_storage(); - - op(&storage_uri); + op(&mem_storage()); } /// This method sets up a test directory with a random name (a number) @@ -41,11 +40,13 @@ fn random_hex_string() -> String { hex::encode(bytes) } -pub fn mem_storage() -> Url { +pub fn mem_storage() -> StorageSystem { let mut bytes = [0; 8]; openssl::rand::rand_bytes(&mut bytes).unwrap(); - Url::parse(&format!("memory://{}", random_hex_string())).unwrap() + StorageSystem::new( + Url::parse(&format!("memory://{}", random_hex_string())).unwrap() + ).unwrap() } pub fn rsync(s: &str) -> uri::Rsync { diff --git a/src/config.rs b/src/config.rs index e0dd9d62..2d259c6b 100644 --- a/src/config.rs +++ b/src/config.rs @@ -26,9 +26,7 @@ use crate::{ commons::{ ext_serde, crypto::{OpenSslSignerConfig, SignSupport}, - error::{Error, KrillIoError}, - storage::{Ident, KeyValueStore}, - KrillResult, + error::KrillIoError, }, constants::*, daemon::{ @@ -897,21 +895,6 @@ pub struct Benchmark { /// # Accessors impl Config { - /// General purpose KV store, can be used to track server settings - /// etc not specific to any Aggregate or WalSupport type - pub fn general_key_value_store(&self) -> KrillResult { - KeyValueStore::create(&self.storage_uri, PROPERTIES_NS) - .map_err(Error::KeyValueError) - } - - pub fn key_value_store( - &self, - namespace: &Ident, - ) -> KrillResult { - KeyValueStore::create(&self.storage_uri, namespace) - .map_err(Error::KeyValueError) - } - /// Returns the data directory if disk was used for storage. /// This will always be true for upgrades of pre 0.14.0 versions fn data_dir(&self) -> Option { diff --git a/src/daemon/http/auth/authorizer.rs b/src/daemon/http/auth/authorizer.rs index c86ad139..5af27798 100644 --- a/src/daemon/http/auth/authorizer.rs +++ b/src/daemon/http/auth/authorizer.rs @@ -9,6 +9,7 @@ use crate::api::admin::Token; use crate::commons::KrillResult; use crate::commons::actor::Actor; use crate::commons::error::ApiAuthError; +use crate::commons::storage::StorageSystem; use crate::config::{AuthType, Config}; use crate::daemon::http::request::HyperRequest; use crate::daemon::http::response::HttpResponse; @@ -196,6 +197,7 @@ impl Authorizer { /// The authorizer will be created according to information provided via /// `config`. pub fn new( + storage: &StorageSystem, config: Arc, ) -> KrillResult { let (primary_provider, legacy_provider) = match config.auth_type { @@ -205,14 +207,18 @@ impl Authorizer { #[cfg(feature = "multi-user")] AuthType::ConfigFile => { ( - config_file::AuthProvider::new(&config)?.into(), + config_file::AuthProvider::new( + storage, &config + )?.into(), Some(admin_token::AuthProvider::new(config)) ) } #[cfg(feature = "multi-user")] AuthType::OpenIDConnect => { ( - openid_connect::AuthProvider::new(config.clone())?.into(), + openid_connect::AuthProvider::new( + storage, config.clone() + )?.into(), Some(admin_token::AuthProvider::new(config)) ) } diff --git a/src/daemon/http/auth/crypt.rs b/src/daemon/http/auth/crypt.rs index 4044d299..8f05b096 100644 --- a/src/daemon/http/auth/crypt.rs +++ b/src/daemon/http/auth/crypt.rs @@ -26,8 +26,7 @@ use serde::{Deserialize, Serialize}; use crate::commons::ext_serde; use crate::commons::KrillResult; use crate::commons::error::{ApiAuthError, Error}; -use crate::commons::storage::Ident; -use crate::config::Config; +use crate::commons::storage::{Ident, StorageSystem}; const CHACHA20_KEY_BIT_LEN: usize = 256; const CHACHA20_KEY_BYTE_LEN: usize = CHACHA20_KEY_BIT_LEN / 8; @@ -164,8 +163,8 @@ pub(crate) fn decrypt( }) } -pub(crate) fn crypt_init(config: &Config) -> KrillResult { - let store = config.key_value_store(CRYPT_STATE_NS)?; +pub(crate) fn crypt_init(storage: &StorageSystem) -> KrillResult { + let store = storage.open(CRYPT_STATE_NS)?; if let Some(state) = store.get(None, CRYPT_STATE_KEY)? { Ok(state) diff --git a/src/daemon/http/auth/providers/config_file.rs b/src/daemon/http/auth/providers/config_file.rs index 371c320c..4a224d41 100644 --- a/src/daemon/http/auth/providers/config_file.rs +++ b/src/daemon/http/auth/providers/config_file.rs @@ -12,6 +12,7 @@ use crate::api::admin::Token; use crate::commons::httpclient; use crate::commons::KrillResult; use crate::commons::error::{ApiAuthError, Error}; +use crate::commons::storage::StorageSystem; use crate::constants::{PW_HASH_LOG_N, PW_HASH_P, PW_HASH_R}; use crate::config::Config; use crate::daemon::http::auth::crypt; @@ -56,13 +57,14 @@ pub struct AuthProvider { impl AuthProvider { /// Creates an auth provider from the given config. pub fn new( + storage: &StorageSystem, config: &Config, ) -> KrillResult { let users = config.auth_users.as_ref().ok_or_else(|| { Error::ConfigError("Missing [auth_users] config section!".into()) })?.clone(); let roles = config.auth_roles.clone(); - let session_key = Self::init_session_key(config)?; + let session_key = Self::init_session_key(storage)?; Ok(Self { users, @@ -78,9 +80,11 @@ impl AuthProvider { } - fn init_session_key(config: &Config) -> KrillResult { + fn init_session_key( + storage: &StorageSystem + ) -> KrillResult { debug!("Initializing login session encryption key"); - crypt::crypt_init(config) + crypt::crypt_init(storage) } /// Parse HTTP Basic Authorization header diff --git a/src/daemon/http/auth/providers/openid_connect/provider.rs b/src/daemon/http/auth/providers/openid_connect/provider.rs index 54558305..d44f9eac 100644 --- a/src/daemon/http/auth/providers/openid_connect/provider.rs +++ b/src/daemon/http/auth/providers/openid_connect/provider.rs @@ -63,6 +63,7 @@ use crate::{ commons::{ httpclient, error::{ApiAuthError, Error}, + storage::StorageSystem, util::sha256, KrillResult, }, @@ -178,9 +179,10 @@ pub struct AuthProvider { impl AuthProvider { pub fn new( + storage: &StorageSystem, config: Arc, ) -> KrillResult { - let session_key = Self::init_session_key(&config)?; + let session_key = Self::init_session_key(storage)?; Ok(Self { config, @@ -741,9 +743,9 @@ impl AuthProvider { } } - fn init_session_key(config: &Config) -> KrillResult { + fn init_session_key(storage: &StorageSystem) -> KrillResult { debug!("Initializing session encryption key"); - crypt::crypt_init(config) + crypt::crypt_init(storage) } fn oidc_conf(&self) -> KrillResult<&ConfigAuthOpenIDConnect> { diff --git a/src/daemon/http/server.rs b/src/daemon/http/server.rs index fe5c3626..19530984 100644 --- a/src/daemon/http/server.rs +++ b/src/daemon/http/server.rs @@ -22,7 +22,7 @@ use super::response::{HyperResponse, HttpResponse}; /// The Krill HTTP server. pub struct HttpServer { /// The Krill “business logic.” - krill: KrillManager, + krill: Arc, /// The component responsible for API authorization checks authorizer: Authorizer, @@ -37,12 +37,13 @@ pub struct HttpServer { impl HttpServer { /// Creates a new server from a Krill manager and the configuration. pub fn new( - krill: KrillManager, + krill: Arc, config: Arc ) -> KrillResult> { + let authorizer = Authorizer::new(krill.storage(), config.clone())?; Ok(Self { krill, - authorizer: Authorizer::new(config.clone())?, + authorizer, config, started: Timestamp::now(), }.into()) diff --git a/src/daemon/start.rs b/src/daemon/start.rs index ae135584..8f31aeff 100644 --- a/src/daemon/start.rs +++ b/src/daemon/start.rs @@ -11,6 +11,7 @@ use tokio::sync::oneshot; use tokio_rustls::TlsAcceptor; use crate::commons::file; use crate::commons::error::Error; +use crate::commons::storage::StorageSystem; use crate::commons::version::KrillVersion; use crate::config::Config; use crate::constants::KRILL_ENV_UPGRADE_ONLY; @@ -31,16 +32,17 @@ pub async fn start_krill_daemon( write_pid_file_or_die(&config); test_data_dirs_or_die(&config); + let storage = StorageSystem::new(config.storage_uri.clone())?; + // Set up the runtime properties manager, so that we can check // the version used for the current data in storage let properties_manager = PropertiesManager::create( - &config.storage_uri, - config.use_history_cache, + &storage, config.use_history_cache, )?; // Call upgrade, this will only do actual work if needed. let upgrade_report = prepare_upgrade_data_migrations( - UpgradeMode::PrepareToFinalise, &config, &properties_manager + UpgradeMode::PrepareToFinalise, &storage, &config, &properties_manager ).map_err(|e| { match e { UpgradeError::CodeOlderThanData(_,_) => { @@ -58,7 +60,7 @@ pub async fn start_krill_daemon( if let Some(report) = &upgrade_report { finalise_data_migration( - report.versions(), &config, &properties_manager + report.versions(), &storage, &properties_manager ).map_err(|e| { Error::Custom(format!( "Finishing prepared migration failed unexpectedly. Please \ @@ -82,7 +84,7 @@ pub async fn start_krill_daemon( // Create the Krill manager, this will create the necessary data // sub-directories if needed - let krill = KrillManager::build(config.clone()).await?; + let krill = Arc::new(KrillManager::build(storage, config.clone()).await?); // Call post-start upgrades to trigger any upgrade related runtime // actions, such as re-issuing ROAs because subject name strategy has @@ -100,11 +102,10 @@ pub async fn start_krill_daemon( // Build the scheduler which will be responsible for executing // planned/triggered tasks - let scheduler = krill.build_scheduler(); - let scheduler_future = scheduler.run(); + let scheduler_future = krill.run_scheduler(); // Create the HTTP server. - let server = HttpServer::new(krill, config.clone())?; + let server = HttpServer::new(krill.clone(), config.clone())?; // Create self-signed HTTPS cert if configured and not generated earlier. if config.https_mode().is_generate_https_cert() { diff --git a/src/server/ca/manager.rs b/src/server/ca/manager.rs index 61e6243d..a0fad222 100644 --- a/src/server/ca/manager.rs +++ b/src/server/ca/manager.rs @@ -54,6 +54,7 @@ use crate::commons::cmslogger::CmsLogger; use crate::commons::crypto::KrillSigner; use crate::commons::error::{Error, Error as KrillError}; use crate::commons::eventsourcing::{Aggregate, AggregateStore, SentCommand}; +use crate::commons::storage::StorageSystem; use crate::constants::{ CASERVER_NS, STATUS_NS, TA_PROXY_SERVER_NS, TA_SIGNER_SERVER_NS, TA_NAME, ta_handle, @@ -138,6 +139,7 @@ impl CaManager { /// Return an error if any of the various stores cannot be initialized. pub async fn build( config: Arc, + storage: &StorageSystem, tasks: Arc, signer: Arc, system_actor: Actor, @@ -145,9 +147,7 @@ impl CaManager { // Create the AggregateStore for the event-sourced `CertAuth` // structures that handle most CA functions. let mut ca_store = AggregateStore::::create( - &config.storage_uri, - CASERVER_NS, - config.use_history_cache, + storage, CASERVER_NS, config.use_history_cache, )?; if let Err(e) = ca_store.warm() { @@ -176,9 +176,7 @@ impl CaManager { // and issued certificates from the `CertAuth` and is responsible // for manifests and CRL generation. let ca_objects_store = Arc::new(CaObjectsStore::create( - &config.storage_uri, - config.issuance_timing.clone(), - signer.clone(), + storage, config.issuance_timing.clone(), signer.clone(), )?); // Register the `CaObjectsStore` as a pre-save listener to the @@ -214,9 +212,7 @@ impl CaManager { // Create TA proxy store if we need it. let ta_proxy_store = if config.ta_proxy_enabled() { let mut store = AggregateStore::::create( - &config.storage_uri, - TA_PROXY_SERVER_NS, - config.use_history_cache, + storage, TA_PROXY_SERVER_NS, config.use_history_cache, )?; // We need a pre-save listener so that we can schedule: @@ -238,9 +234,7 @@ impl CaManager { let ta_signer_store = if config.ta_signer_enabled() { Some(AggregateStore::create( - &config.storage_uri, - TA_SIGNER_SERVER_NS, - config.use_history_cache, + storage, TA_SIGNER_SERVER_NS, config.use_history_cache, )?) } else { @@ -250,9 +244,7 @@ impl CaManager { // Create the status store which will maintain the last known // connection status between each CA and their parent(s) and // repository. - let status_store = CaStatusStore::create( - &config.storage_uri, STATUS_NS - )?; + let status_store = CaStatusStore::create(storage, STATUS_NS)?; Ok(CaManager { ca_store, diff --git a/src/server/ca/publishing.rs b/src/server/ca/publishing.rs index 603565e4..ee9c48b5 100644 --- a/src/server/ca/publishing.rs +++ b/src/server/ca/publishing.rs @@ -15,7 +15,6 @@ use rpki::repository::manifest::{FileAndHash, Manifest, ManifestContent}; use rpki::repository::sigobj::SignedObjectBuilder; use rpki::repository::x509::{Name, Serial, Time, Validity}; use serde::{Deserialize, Serialize}; -use url::Url; use crate::api::admin::{PublishedFile, RepositoryContact}; use crate::api::ca::{ CertInfo, IssuedCertificate, ObjectName, ReceivedCert, Revocation, @@ -26,7 +25,9 @@ use crate::commons::KrillResult; use crate::commons::crypto::KrillSigner; use crate::commons::error::Error; use crate::commons::eventsourcing::PreSaveEventListener; -use crate::commons::storage::{Ident, KeyValueStore}; +use crate::commons::storage::{ + Ident, KeyValueStore, OpenStoreError, StorageSystem +}; use crate::constants::CA_OBJECTS_NS; use crate::config::IssuanceTimingConfig; use super::aspa::{AspaInfo, AspaObjectsUpdates}; @@ -74,11 +75,11 @@ pub struct CaObjectsStore { impl CaObjectsStore { /// Creates a new CA objects store using the given configuration. pub fn create( - storage_uri: &Url, + storage: &StorageSystem, issuance_timing: IssuanceTimingConfig, signer: Arc, - ) -> KrillResult { - let store = KeyValueStore::create(storage_uri, CA_OBJECTS_NS)?; + ) -> Result { + let store = storage.open(CA_OBJECTS_NS)?; Ok(CaObjectsStore { store, signer, diff --git a/src/server/ca/status.rs b/src/server/ca/status.rs index 435662b6..d572ca73 100644 --- a/src/server/ca/status.rs +++ b/src/server/ca/status.rs @@ -9,7 +9,6 @@ use rpki::ca::idexchange::{CaHandle, ChildHandle, ParentHandle, ServiceUri}; use rpki::ca::provisioning::ResourceClassListResponse as Entitlements; use rpki::ca::publication::PublishDelta; use serde::{Deserialize, Serialize}; -use url::Url; use crate::api::ca::{ ChildConnectionStats, ChildStatus, ChildrenConnectionStats, ParentStatus, ParentStatuses, RepoStatus, @@ -18,7 +17,7 @@ use crate::api::status::ErrorResponse; use crate::commons::httpclient; use crate::commons::KrillResult; use crate::commons::error::Error; -use crate::commons::storage::{Ident, KeyValueStore}; +use crate::commons::storage::{Ident, KeyValueStore, StorageSystem}; const PARENTS_PREFIX: &Ident = Ident::make("parents-"); const CHILDREN_PREFIX: &Ident = Ident::make("children-"); @@ -64,10 +63,10 @@ pub struct CaStatusStore { impl CaStatusStore { /// Creates a new status store with the givn storage URI and namespace. pub fn create( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, ) -> KrillResult { - let store = KeyValueStore::create(storage_uri, namespace)?; + let store = storage.open(namespace)?; let cache = RwLock::new(HashMap::new()); let store = Self { store, cache }; diff --git a/src/server/ca/upgrades/data_migration.rs b/src/server/ca/upgrades/data_migration.rs index 6063f2f6..0dfb8aa8 100644 --- a/src/server/ca/upgrades/data_migration.rs +++ b/src/server/ca/upgrades/data_migration.rs @@ -1,6 +1,7 @@ use std::sync::Arc; use log::{debug, warn}; use crate::commons::crypto::KrillSignerBuilder; +use crate::commons::storage::StorageSystem; use crate::constants::CASERVER_NS; use crate::server::ca::certauth::CertAuth; use crate::server::ca::publishing::CaObjectsStore; @@ -9,8 +10,10 @@ use crate::upgrades::UpgradeResult; use crate::upgrades::data_migration::check_agg_store; -pub fn check_ca_objects(config: &Config) -> UpgradeResult<()> { - let ca_store = check_agg_store::(config, CASERVER_NS, "CAs")?; +pub fn check_ca_objects( + storage: &StorageSystem, config: &Config +) -> UpgradeResult<()> { + let ca_store = check_agg_store::(storage, CASERVER_NS, "CAs")?; // make a dummy Signer to use for the CaObjectsStore - it won't be used, // but it's needed for construction. @@ -18,7 +21,7 @@ pub fn check_ca_objects(config: &Config) -> UpgradeResult<()> { std::time::Duration::from_secs(config.signer_probe_retry_seconds); let signer = Arc::new( KrillSignerBuilder::new( - &config.storage_uri, + storage, probe_interval, &config.signers, ) @@ -28,9 +31,7 @@ pub fn check_ca_objects(config: &Config) -> UpgradeResult<()> { ); let ca_objects_store = CaObjectsStore::create( - &config.storage_uri, - config.issuance_timing.clone(), - signer, + storage, config.issuance_timing.clone(), signer, )?; let cas_with_objects = ca_objects_store.cas()?; diff --git a/src/server/ca/upgrades/pre_0_10_0/migration.rs b/src/server/ca/upgrades/pre_0_10_0/migration.rs index 59da0a48..37d9ca7f 100644 --- a/src/server/ca/upgrades/pre_0_10_0/migration.rs +++ b/src/server/ca/upgrades/pre_0_10_0/migration.rs @@ -7,13 +7,14 @@ use crate::api::aspa::ProviderAsn; use crate::commons::eventsourcing::{ AggregateStore, StoredCommand, StoredCommandBuilder }; -use crate::commons::storage::{Ident, KeyValueStore}; +use crate::commons::storage::{ + Ident, KeyValueStore, OpenStoreError, StorageSystem +}; use crate::constants::{CASERVER_NS, CA_OBJECTS_NS}; use crate::server::ca::certauth::CertAuth; use crate::server::ca::commands::CertAuthStorableCommand; use crate::server::ca::events::{CertAuthEvent, CertAuthInitEvent}; use crate::server::ca::publishing::CaObjects; -use crate::config::Config; use crate::upgrades::{ AspaMigrationConfigUpdates, AspaMigrationConfigs, CommandMigrationEffect, UpgradeAggregateStorePre0_14, UpgradeError, UpgradeMode, UpgradeResult, @@ -50,22 +51,17 @@ impl CasMigration { /// Upgrades the CAs based on the upgrade mode and config. pub fn upgrade( mode: UpgradeMode, - config: &Config, + storage: &StorageSystem, ) -> UpgradeResult { Self { - current_kv_store: KeyValueStore::create( - &config.storage_uri, CASERVER_NS - )?, - new_kv_store: KeyValueStore::create_upgrade_store( - &config.storage_uri, - CASERVER_NS, - )?, + current_kv_store: storage.open(CASERVER_NS)?, + new_kv_store: storage.open_upgrade(CASERVER_NS)?, new_agg_store: AggregateStore::::create_upgrade_store( - &config.storage_uri, + storage, CASERVER_NS, - config.use_history_cache, + false, )?, - ca_objects_migration: CaObjectsMigration::create(config)?, + ca_objects_migration: CaObjectsMigration::create(storage)?, } .upgrade(mode) } @@ -252,15 +248,10 @@ struct CaObjectsMigration { impl CaObjectsMigration { /// Creates a new migration from the configuration. - fn create(config: &Config) -> Result { + fn create(storage: &StorageSystem) -> Result { Ok(CaObjectsMigration { - current_store: KeyValueStore::create( - &config.storage_uri, CA_OBJECTS_NS - )?, - new_store: KeyValueStore::create_upgrade_store( - &config.storage_uri, - CA_OBJECTS_NS, - )? + current_store: storage.open(CA_OBJECTS_NS)?, + new_store: storage.open_upgrade(CA_OBJECTS_NS)? }) } diff --git a/src/server/ca/upgrades/pre_0_14_0/migration.rs b/src/server/ca/upgrades/pre_0_14_0/migration.rs index 3a6a4967..1d2dda44 100644 --- a/src/server/ca/upgrades/pre_0_14_0/migration.rs +++ b/src/server/ca/upgrades/pre_0_14_0/migration.rs @@ -10,10 +10,9 @@ use crate::upgrades::{ use crate::{ commons::{ eventsourcing::AggregateStore, - storage::KeyValueStore, + storage::{KeyValueStore, StorageSystem}, }, constants::CASERVER_NS, - config::Config, upgrades::UpgradeResult, }; @@ -36,19 +35,13 @@ pub struct CasMigration { impl CasMigration { pub fn upgrade( mode: UpgradeMode, - config: &Config, + storage: &StorageSystem, ) -> UpgradeResult { - let current_kv_store = - KeyValueStore::create(&config.storage_uri, CASERVER_NS)?; - let new_kv_store = KeyValueStore::create_upgrade_store( - &config.storage_uri, - CASERVER_NS, - )?; + let current_kv_store = storage.open(CASERVER_NS)?; + let new_kv_store = storage.open_upgrade(CASERVER_NS)?; let new_agg_store = AggregateStore::::create_upgrade_store( - &config.storage_uri, - CASERVER_NS, - config.use_history_cache, + storage, CASERVER_NS, false, )?; CasMigration { diff --git a/src/server/manager.rs b/src/server/manager.rs index 0000fecd..f56609da 100644 --- a/src/server/manager.rs +++ b/src/server/manager.rs @@ -34,7 +34,6 @@ use crate::{ }, mq::{now, Task, TaskQueue}, pubd::RepositoryManager, - scheduler::Scheduler, }, }; use crate::api; @@ -52,8 +51,7 @@ use crate::api::bgpsec::{BgpSecCsrInfoList, BgpSecDefinitionUpdates}; use crate::api::ca::{ CaRepoDetails, CertAuthInfo, CertAuthIssues, CertAuthList, CertAuthStats, ChildCaInfo, ChildrenConnectionStats, - IdCertInfo, RtaList, RtaName, - RtaPrepResponse, + IdCertInfo, RtaList, RtaName, RtaPrepResponse, Timestamp, }; use crate::api::history::{ CommandDetails, CommandHistory, CommandHistoryCriteria @@ -70,6 +68,7 @@ use crate::api::ta::{ ApiTrustAnchorSignedRequest, TaCertDetails, TrustAnchorSignedResponse, TrustAnchorSignerInfo, }; +use crate::commons::storage::StorageSystem; use crate::constants::{TA_NAME, ta_handle}; use crate::server::bgp::BgpAnalyser; @@ -80,31 +79,38 @@ use crate::server::bgp::BgpAnalyser; /// components. pub struct KrillManager { // The base URI for this service - service_uri: uri::Https, + pub(super) service_uri: uri::Https, + + /// The storage system used by all components. + pub(super) storage: StorageSystem, // Publication server, with configured publishers - repo_manager: Arc, + pub(super) repo_manager: Arc, // Handles the internal TA and/or CAs - ca_manager: Arc, + pub(super) ca_manager: Arc, // Handles the internal TA and/or CAs - bgp_analyser: Arc, + pub(super) bgp_analyser: Arc, // Shared message queue - mq: Arc, + pub(super) tasks: Arc, // System actor - system_actor: Actor, + pub(super) system_actor: Actor, pub config: Arc, + + pub(super) started: Timestamp, } /// # Set up and initialization impl KrillManager { /// Creates a new publication server. Note that state is preserved /// in the data storage. - pub async fn build(config: Arc) -> KrillResult { + pub async fn build( + storage: StorageSystem, config: Arc + ) -> KrillResult { let service_uri = config.service_uri(); info!("Starting {} v{}", KRILL_SERVER_APP, crate_version!()); @@ -117,7 +123,7 @@ impl KrillManager { let probe_interval = std::time::Duration::from_secs(config.signer_probe_retry_seconds); let signer = KrillSignerBuilder::new( - &config.storage_uri, + &storage, probe_interval, &config.signers, ) @@ -130,21 +136,23 @@ impl KrillManager { // Task queue Arc is shared between ca_manager, repo_manager and the // scheduler. - let mq = Arc::new(TaskQueue::new(&config.storage_uri)?); + let tasks = Arc::new(TaskQueue::new(&storage)?); // for now, support that existing embedded repositories are still // supported. this should be removed in future after people // have had a chance to separate. let repo_manager = Arc::new(RepositoryManager::build( + &storage, config.clone(), - mq.clone(), + tasks.clone(), signer.clone(), )?); let ca_manager = Arc::new( ca::CaManager::build( config.clone(), - mq.clone(), + &storage, + tasks.clone(), signer, system_actor.clone(), ) @@ -162,18 +170,20 @@ impl KrillManager { // When multi-node set ups with a shared queue are // supported then we can no longer safely reschedule // ALL running tests. See issue: #1112 - mq.reschedule_tasks_at_startup()?; + tasks.reschedule_tasks_at_startup()?; - mq.schedule(Task::QueueStartTasks, now())?; + tasks.schedule(Task::QueueStartTasks, now())?; let server = KrillManager { service_uri, + storage, repo_manager, ca_manager, bgp_analyser, - mq, + tasks, system_actor, config: config.clone(), + started: Timestamp::now(), }; // Check if we need to do any testbed or benchmarking set up. @@ -271,24 +281,18 @@ impl KrillManager { Ok(server) } - - pub fn build_scheduler(&self) -> Scheduler { - Scheduler::build( - self.mq.clone(), - self.ca_manager.clone(), - self.repo_manager.clone(), - self.config.clone(), - self.system_actor.clone(), - ) - } - - pub fn service_base_uri(&self) -> &uri::Https { - &self.service_uri - } } /// # Access to components impl KrillManager { + pub fn service_base_uri(&self) -> &uri::Https { + &self.service_uri + } + + pub fn storage(&self) -> &StorageSystem { + &self.storage + } + pub fn system_actor(&self) -> &Actor { &self.system_actor } diff --git a/src/server/mq.rs b/src/server/mq.rs index e73937b7..e1afad88 100644 --- a/src/server/mq.rs +++ b/src/server/mq.rs @@ -10,13 +10,12 @@ use rpki::ca::idexchange::{CaHandle, ParentHandle}; use rpki::ca::provisioning::{ResourceClassName, RevocationRequest}; use rpki::repository::x509::Time; use serde::{Deserialize, Serialize}; -use url::Url; use crate::api::ca::Timestamp; use crate::commons::eventsourcing; use crate::commons::{Error, KrillResult}; use crate::commons::eventsourcing::Aggregate; use crate::commons::queue::{Queue, ScheduleMode}; -use crate::commons::storage::Ident; +use crate::commons::storage::{Ident, StorageSystem}; use crate::constants::{TASK_QUEUE_NS, ta_handle}; use crate::server::ca::{CertAuth, CertAuthEvent}; use crate::server::taproxy::{TrustAnchorProxy, TrustAnchorProxyEvent}; @@ -286,9 +285,9 @@ pub struct TaskQueue { } impl TaskQueue { - pub fn new(storage_uri: &Url) -> KrillResult { + pub fn new(storage: &StorageSystem) -> KrillResult { Ok(TaskQueue { - q: Queue::create(storage_uri, TASK_QUEUE_NS)?, + q: Queue::create(storage, TASK_QUEUE_NS)?, }) } } diff --git a/src/server/properties/mod.rs b/src/server/properties/mod.rs index dd1787d1..a4f22cae 100644 --- a/src/server/properties/mod.rs +++ b/src/server/properties/mod.rs @@ -20,7 +20,6 @@ use std::{fmt, str::FromStr, sync::Arc}; use log::{log_enabled, trace}; use rpki::ca::idexchange::MyHandle; use serde::{Deserialize, Serialize}; -use url::Url; use crate::{ commons::{ @@ -30,6 +29,7 @@ use crate::{ self, Aggregate, AggregateStore, Event, InitCommandDetails, InitEvent, SentCommand, SentInitCommand, WithStorableDetails, }, + storage::StorageSystem, version::KrillVersion, KrillResult, }, @@ -296,17 +296,19 @@ pub struct PropertiesManager { impl PropertiesManager { pub fn create( - storage_uri: &Url, + storage: &StorageSystem, use_history_cache: bool, ) -> KrillResult { let main_key = MyHandle::from_str(PROPERTIES_DFLT_NAME).unwrap(); - AggregateStore::create(storage_uri, PROPERTIES_NS, use_history_cache) - .map(|store| PropertiesManager { - store, - main_key, - system_actor: ACTOR_DEF_KRILL, - }) - .map_err(Error::AggregateStoreError) + let store = AggregateStore::create( + storage, PROPERTIES_NS, use_history_cache + )?; + + Ok(PropertiesManager { + store, + main_key, + system_actor: ACTOR_DEF_KRILL, + }) } pub fn is_initialized(&self) -> bool { diff --git a/src/server/pubd/access.rs b/src/server/pubd/access.rs index 7921801e..2387c126 100644 --- a/src/server/pubd/access.rs +++ b/src/server/pubd/access.rs @@ -25,10 +25,10 @@ use crate::commons::eventsourcing::{ Aggregate, AggregateStore, CommandDetails, Event, InitCommandDetails, InitEvent, SentCommand, SentInitCommand, WithStorableDetails, }; +use crate::commons::storage::StorageSystem; use crate::constants::{ ACTOR_DEF_KRILL, PUBSERVER_DFLT, PUBSERVER_NS, TA_NAME }; -use crate::config::Config; use super::publishers::Publisher; @@ -50,12 +50,13 @@ pub struct RepositoryAccessProxy { } impl RepositoryAccessProxy { - /// Creates a new repository access proxy from the config. - pub fn create(config: &Config) -> KrillResult { + /// Creates a new repository access proxy + pub fn create( + storage: &StorageSystem, + use_history_cache: bool + ) -> KrillResult { let store = AggregateStore::::create( - &config.storage_uri, - PUBSERVER_NS, - config.use_history_cache, + storage, PUBSERVER_NS, use_history_cache, )?; let key = MyHandle::from_str(PUBSERVER_DFLT).unwrap(); diff --git a/src/server/pubd/content.rs b/src/server/pubd/content.rs index 1fffed83..9164f489 100644 --- a/src/server/pubd/content.rs +++ b/src/server/pubd/content.rs @@ -16,8 +16,9 @@ use crate::commons::error::Error; use crate::commons::eventsourcing::{ WalChange, WalCommand, WalSet, WalStore, WalSupport, }; +use crate::commons::storage::StorageSystem; use crate::constants::PUBSERVER_CONTENT_NS; -use crate::config::{Config, RrdpUpdatesConfig}; +use crate::config::RrdpUpdatesConfig; use super::rrdp::{ CurrentObjects, DeltaElements, RrdpServer, RrdpSession, RrdpSessionReset, RrdpUpdated, RrdpUpdateNeeded, @@ -43,9 +44,9 @@ pub struct RepositoryContentProxy { impl RepositoryContentProxy { /// Creates a new repository content proxy. - pub fn create(config: &Config) -> KrillResult { + pub fn create(storage: &StorageSystem) -> KrillResult { let store = Arc::new(WalStore::create( - &config.storage_uri, PUBSERVER_CONTENT_NS, + storage, PUBSERVER_CONTENT_NS, )?); store.warm()?; diff --git a/src/server/pubd/manager.rs b/src/server/pubd/manager.rs index 1d1d876f..f9802212 100644 --- a/src/server/pubd/manager.rs +++ b/src/server/pubd/manager.rs @@ -19,6 +19,7 @@ use crate::commons::actor::Actor; use crate::commons::cmslogger::CmsLogger; use crate::commons::crypto::KrillSigner; use crate::commons::error::Error; +use crate::commons::storage::StorageSystem; use crate::config::Config; use crate::server::mq::{now, Task, TaskQueue}; use super::access::RepositoryAccessProxy; @@ -53,13 +54,16 @@ pub struct RepositoryManager { impl RepositoryManager { /// Builds the repository manager. pub fn build( + storage: &StorageSystem, config: Arc, tasks: Arc, signer: Arc, ) -> Result { - let access_proxy = Arc::new(RepositoryAccessProxy::create(&config)?); + let access_proxy = Arc::new(RepositoryAccessProxy::create( + storage, config.use_history_cache, + )?); let content_proxy = Arc::new( - RepositoryContentProxy::create(&config)? + RepositoryContentProxy::create(storage)? ); Ok(RepositoryManager { @@ -354,7 +358,6 @@ mod tests { use std::time::Duration; use bytes::Bytes; use tokio::time::sleep; - use url::Url; use rpki::uri; use rpki::ca::idexchange::Handle; use rpki::ca::publication::{ListElement, PublishDelta}; @@ -362,6 +365,7 @@ mod tests { use crate::commons::file; use crate::commons::crypto::{KrillSignerBuilder, OpenSslSignerConfig}; use crate::commons::file::CurrentFile; + use crate::commons::storage::StorageSystem; use crate::commons::test::{self, https, rsync}; use crate::constants::{ ACTOR_DEF_TEST, RRDP_FIRST_SERIAL, enable_test_mode @@ -371,7 +375,7 @@ mod tests { use crate::server::pubd::rrdp::{PublicationDeltaError, RrdpServer}; use super::*; - fn publisher_alice(storage_uri: &Url) -> Publisher { + fn publisher_alice(storage: &StorageSystem) -> Publisher { // When the "hsm" feature is enabled we could be running the tests // with PKCS#11 as the default signer type. In that case, if // the backend signer is SoftHSMv2, attempting to create a second @@ -387,7 +391,7 @@ mod tests { SignerConfig::new("Alice".to_string(), signer_type); let signer_configs = &[signer_config]; KrillSignerBuilder::new( - storage_uri, + storage, Duration::from_secs(1), signer_configs, ) @@ -415,13 +419,13 @@ mod tests { } fn make_server( - storage_uri: &Url + storage: &StorageSystem ) -> (RepositoryManager, tempfile::TempDir) { let data_dir = tempfile::tempdir().unwrap(); enable_test_mode(); let mut config = Config::test( - storage_uri, + storage.default_uri(), Some(data_dir.path()), true, false, @@ -432,7 +436,7 @@ mod tests { config.process().unwrap(); let signer = KrillSignerBuilder::new( - storage_uri, + storage, Duration::from_secs(1), &config.signers, ) @@ -443,9 +447,10 @@ mod tests { let signer = Arc::new(signer); let config = Arc::new(config); - let mq = Arc::new(TaskQueue::new(&config.storage_uri).unwrap()); - let repository_manager = - RepositoryManager::build(config, mq, signer).unwrap(); + let mq = Arc::new(TaskQueue::new(storage).unwrap()); + let repository_manager = RepositoryManager::build( + storage, config, mq, signer + ).unwrap(); let uris = PublicationServerUris { rrdp_base_uri: https("https://localhost/repo/rrdp/"), @@ -460,10 +465,10 @@ mod tests { #[test] fn should_add_publisher() { // we need a disk, as repo_dir, etc. use data_dir by default - let storage_uri = test::mem_storage(); - let (server, _data_dir) = make_server(&storage_uri); + let storage = test::mem_storage(); + let (server, _data_dir) = make_server(&storage); - let alice = publisher_alice(&storage_uri); + let alice = publisher_alice(&storage); let alice_handle = Handle::from_str("alice").unwrap(); let publisher_req = diff --git a/src/server/pubd/upgrades/mod.rs b/src/server/pubd/upgrades/mod.rs index 9bf7a815..ff71b618 100644 --- a/src/server/pubd/upgrades/mod.rs +++ b/src/server/pubd/upgrades/mod.rs @@ -9,18 +9,19 @@ use log::info; use rpki::ca::idexchange::MyHandle; use crate::commons::KrillResult; use crate::commons::eventsourcing::WalStore; -use crate::commons::storage::{Ident, KeyValueStore}; +use crate::commons::storage::{Ident, StorageSystem}; use crate::constants::PUBSERVER_CONTENT_NS; -use crate::config::Config; use crate::server::pubd::content::RepositoryContent; use self::pre_0_13_0::OldRepositoryContent; /// Migrate v0.12.x RepositoryContent to the new 0.13.0+ format. /// Apply any open WAL changes to the source first. -pub fn migrate_0_12_pubd_objects(config: &Config) -> KrillResult { +pub fn migrate_0_12_pubd_objects( + storage: &StorageSystem +) -> KrillResult { let old_store: WalStore = WalStore::create( - &config.storage_uri, PUBSERVER_CONTENT_NS + storage, PUBSERVER_CONTENT_NS )?; let repo_content_handle = MyHandle::new("0".into()); @@ -29,10 +30,7 @@ pub fn migrate_0_12_pubd_objects(config: &Config) -> KrillResult { old_store.get_latest(&repo_content_handle)?.as_ref().clone(); let repo_content: RepositoryContent = old_repo_content.try_into()?; - let upgrade_store = KeyValueStore::create_upgrade_store( - &config.storage_uri, - PUBSERVER_CONTENT_NS, - )?; + let upgrade_store = storage.open_upgrade(PUBSERVER_CONTENT_NS)?; upgrade_store.store( Some(const { Ident::make("0") }), const { Ident::make("snapshot.json") }, @@ -46,10 +44,10 @@ pub fn migrate_0_12_pubd_objects(config: &Config) -> KrillResult { /// The format of the RepositoryContent did not change in 0.12, but /// the location and way of storing it did. So, migrate if present. -pub fn migrate_pre_0_12_pubd_objects(config: &Config) -> KrillResult<()> { - let old_store = KeyValueStore::create( - &config.storage_uri, PUBSERVER_CONTENT_NS - )?; +pub fn migrate_pre_0_12_pubd_objects( + storage: &StorageSystem +) -> KrillResult<()> { + let old_store = storage.open(PUBSERVER_CONTENT_NS)?; if let Ok(Some(old_repo_content)) = old_store.get::( None, const { Ident::make("0.json") } @@ -58,9 +56,7 @@ pub fn migrate_pre_0_12_pubd_objects(config: &Config) -> KrillResult<()> { info!("Found pre 0.12.0 RC2 publication server data. Migrating.."); let repo_content: RepositoryContent = old_repo_content.try_into()?; - let upgrade_store = KeyValueStore::create_upgrade_store( - &config.storage_uri, PUBSERVER_CONTENT_NS, - )?; + let upgrade_store = storage.open_upgrade(PUBSERVER_CONTENT_NS)?; upgrade_store.store( Some(const { Ident::make("0") }), const { Ident::make("snapshot.json") }, diff --git a/src/server/pubd/upgrades/pre_0_10_0/migration.rs b/src/server/pubd/upgrades/pre_0_10_0/migration.rs index cec62fa1..88e8b6a4 100644 --- a/src/server/pubd/upgrades/pre_0_10_0/migration.rs +++ b/src/server/pubd/upgrades/pre_0_10_0/migration.rs @@ -1,10 +1,9 @@ use rpki::ca::idexchange::MyHandle; use rpki::repository::x509::Time; use crate::commons::eventsourcing::{AggregateStore, StoredCommandBuilder}; -use crate::commons::storage::{Ident, KeyValueStore}; +use crate::commons::storage::{Ident, KeyValueStore, StorageSystem}; use crate::commons::version::KrillVersion; use crate::constants::PUBSERVER_NS; -use crate::config::Config; use crate::server::pubd::access::{ RepositoryAccess, RepositoryAccessEvent, RepositoryAccessInitEvent, StorableRepositoryCommand, @@ -32,19 +31,15 @@ pub struct PublicationServerRepositoryAccessMigration { impl PublicationServerRepositoryAccessMigration { pub fn upgrade( mode: UpgradeMode, - config: &Config, + storage: &StorageSystem, versions: &UpgradeVersions, ) -> UpgradeResult<()> { - let current_kv_store = - KeyValueStore::create(&config.storage_uri, PUBSERVER_NS)?; - let new_kv_store = KeyValueStore::create_upgrade_store( - &config.storage_uri, - PUBSERVER_NS, - )?; + let current_kv_store = storage.open(PUBSERVER_NS)?; + let new_kv_store = storage.open_upgrade(PUBSERVER_NS)?; let new_agg_store = AggregateStore::create_upgrade_store( - &config.storage_uri, + storage, PUBSERVER_NS, - config.use_history_cache, + false, )?; let store_migration = PublicationServerRepositoryAccessMigration { diff --git a/src/server/scheduler.rs b/src/server/scheduler.rs index 4ee22f7a..acdbde35 100644 --- a/src/server/scheduler.rs +++ b/src/server/scheduler.rs @@ -1,7 +1,7 @@ //! Deal with asynchronous scheduled processes, either triggered by an //! event that occurred, or planned (e.g. re-publishing). -use std::{collections::HashMap, sync::Arc, time::Duration}; +use std::{collections::HashMap, time::Duration}; use tokio::time::sleep; @@ -10,16 +10,13 @@ use rpki::ca::{ idexchange::{CaHandle, ParentHandle}, provisioning::{ResourceClassName, RevocationRequest}, }; -use url::Url; use crate::{ - api::ca::Timestamp, commons::{ - actor::Actor, crypto::dispatch::signerinfo::SignerInfo, error::FatalError, eventsourcing::{Aggregate, AggregateStore, WalStore, WalSupport}, - storage::Ident, + storage:: {Ident, StorageSystem}, version::KrillVersion, }, constants::{ @@ -28,49 +25,21 @@ use crate::{ SCHEDULER_RESYNC_REPO_CAS_THRESHOLD, SCHEDULER_USE_JITTER_CAS_THRESHOLD, SIGNERS_NS, }, - config::Config, server::{ - ca::{CaManager, CertAuth}, - mq::{ - in_hours, in_minutes, in_seconds, in_weeks, now, Task, TaskQueue, - }, + ca::CertAuth, + mq::{in_hours, in_minutes, in_seconds, in_weeks, now, Task}, properties::Properties, - pubd::{RepositoryAccess, RepositoryContent, RepositoryManager}, + pubd::{RepositoryAccess, RepositoryContent}, }, }; +use super::manager::KrillManager; use super::mq::TaskResult; -pub struct Scheduler { - tasks: Arc, - ca_manager: Arc, - repo_manager: Arc, - config: Arc, - system_actor: Actor, - started: Timestamp, -} - -impl Scheduler { - pub fn build( - tasks: Arc, - ca_manager: Arc, - repo_manager: Arc, - config: Arc, - system_actor: Actor, - ) -> Self { - Scheduler { - tasks, - ca_manager, - repo_manager, - config, - system_actor, - started: Timestamp::now(), - } - } - +impl KrillManager { /// Run the scheduler in the background. It will sweep the message queue /// for tasks and re-schedule new tasks as needed. - pub async fn run(&self) { + pub async fn run_scheduler(&self) { loop { while let Some((task_key, value)) = self.tasks.pop() { match serde_json::from_value(value) { @@ -504,10 +473,10 @@ impl Scheduler { // Call update_snapshots on all AggregateStores and WalStores fn update_snapshots(&self) -> Result { fn update_aggregate_store_snapshots( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, ) { - match AggregateStore::::create(storage_uri, namespace, false) { + match AggregateStore::::create(storage, namespace, false) { Err(e) => { // Note: this is highly unlikely.. probably something else // is broken and Krill would @@ -532,10 +501,10 @@ impl Scheduler { } fn update_wal_store_snapshots( - storage_uri: &Url, + storage: &StorageSystem, namespace: &Ident, ) { - match WalStore::::create(storage_uri, namespace) { + match WalStore::::create(storage, namespace) { Err(e) => { // Note: this is highly unlikely.. probably something else // is broken and Krill would @@ -558,24 +527,24 @@ impl Scheduler { } update_aggregate_store_snapshots::( - &self.config.storage_uri, + &self.storage, CASERVER_NS, ); update_aggregate_store_snapshots::( - &self.config.storage_uri, + &self.storage, SIGNERS_NS, ); update_aggregate_store_snapshots::( - &self.config.storage_uri, + &self.storage, PROPERTIES_NS, ); update_aggregate_store_snapshots::( - &self.config.storage_uri, + &self.storage, PUBSERVER_NS, ); update_wal_store_snapshots::( - &self.config.storage_uri, + &self.storage, PUBSERVER_CONTENT_NS, ); diff --git a/src/tasigner/config.rs b/src/tasigner/config.rs index ff9ba44c..094d8a49 100644 --- a/src/tasigner/config.rs +++ b/src/tasigner/config.rs @@ -11,6 +11,7 @@ use url::Url; use crate::{ commons::crypto::{KrillSigner, KrillSignerBuilder, OpenSslSignerConfig}, + commons::storage::StorageSystem, constants::OPENSSL_ONE_OFF_SIGNER_NAME, config::{LogType, SignerConfig, SignerReference, SignerType}, }; @@ -227,7 +228,9 @@ impl Config { } // Signer support - pub fn signer(&self) -> Result, ConfigError> { + pub fn signer( + &self, storage: &StorageSystem + ) -> Result, ConfigError> { // Assumes that Config::verify() has already ensured that the signer // configuration is valid and that Config::resolve() has been // used to update signer name references to resolve to the @@ -235,7 +238,7 @@ impl Config { let probe_interval = std::time::Duration::from_secs(self.signer_probe_retry_seconds); let signer = KrillSignerBuilder::new( - &self.storage_uri, + storage, probe_interval, &self.signers, ) @@ -376,11 +379,11 @@ mod tests { #[test] fn initialise_default_signers() { - test::test_in_memory(|_storage_uri| { + test::test_in_memory(|storage| { let config_string = include_str!("../../test-resources/ta/ta.conf"); let config = Config::parse_str(config_string).unwrap(); - config.signer().unwrap(); + config.signer(storage).unwrap(); }) } } diff --git a/src/upgrades/data_migration.rs b/src/upgrades/data_migration.rs index cd61dec2..f8a61864 100644 --- a/src/upgrades/data_migration.rs +++ b/src/upgrades/data_migration.rs @@ -10,12 +10,12 @@ use crate::{ commons::{ crypto::{ dispatch::signerinfo::SignerInfo, - OpenSslSigner, + OpenSslSigner, OpenSslSignerConfig, }, eventsourcing::{ Aggregate, AggregateStore, WalStore, WalSupport, }, - storage::{Ident, KeyValueStore}, + storage::{Ident, StorageSystem}, }, constants::{ KEYS_NS, PROPERTIES_NS, PUBSERVER_CONTENT_NS, @@ -37,7 +37,9 @@ use crate::{ use super::UpgradeResult; -pub fn migrate(mut config: Config, target_storage: Url) -> UpgradeResult<()> { +pub fn migrate( + config: Config, target_storage: &StorageSystem +) -> UpgradeResult<()> { // Copy the source data from config unmodified into the target_storage info!("-----------------------------------------------------------"); info!(" Krill Data Migration"); @@ -47,18 +49,17 @@ pub fn migrate(mut config: Config, target_storage: Url) -> UpgradeResult<()> { info!("STEP 1: Copy data"); info!(""); info!("From: {}", &config.storage_uri); - info!(" To: {}", &target_storage); + info!(" To: {}", &target_storage.default_uri()); info!("-----------------------------------------------------------"); info!(""); - copy_data_for_migration(&config, &target_storage)?; + copy_data_for_migration(target_storage, &config.storage_uri)?; - // Update the config file with the new target_storage - // and perform a normal data migration - the source data - // could be for an older version of Krill. - config.storage_uri = target_storage; - let properties_manager = - PropertiesManager::create(&config.storage_uri, false)?; + // Perform a normal data migration using the target storage - the source + // data could be for an older version of Krill. + let properties_manager = PropertiesManager::create( + target_storage, false + )?; info!("-----------------------------------------------------------"); info!("STEP 2: Upgrade data to current Krill version (if needed)"); @@ -66,12 +67,13 @@ pub fn migrate(mut config: Config, target_storage: Url) -> UpgradeResult<()> { info!(""); if let Some(upgrade) = prepare_upgrade_data_migrations( crate::upgrades::UpgradeMode::PrepareToFinalise, + target_storage, &config, &properties_manager, )? { finalise_data_migration( upgrade.versions(), - &config, + target_storage, &properties_manager, )?; } @@ -92,53 +94,53 @@ pub fn migrate(mut config: Config, target_storage: Url) -> UpgradeResult<()> { // That said, it's a pretty easy check to perform and it kind of makes // sense to do it to now, even if it would be to point users at deeper // source data issues. - verify_target_data(&config) + verify_target_data(target_storage, &config) } -fn verify_target_data(config: &Config) -> UpgradeResult<()> { - check_agg_store::(config, PROPERTIES_NS, "Properties")?; - check_agg_store::(config, SIGNERS_NS, "Signer")?; +fn verify_target_data( + storage: &StorageSystem, config: &Config +) -> UpgradeResult<()> { + check_agg_store::(storage, PROPERTIES_NS, "Properties")?; + check_agg_store::(storage, SIGNERS_NS, "Signer")?; - check_ca_objects(config)?; + check_ca_objects(storage, config)?; check_agg_store::( - config, + storage, PUBSERVER_NS, "Publication Server Access", )?; check_wal_store::( - config, + storage, PUBSERVER_CONTENT_NS, "Publication Server Objects", )?; check_agg_store::( - config, + storage, TA_PROXY_SERVER_NS, "TA Proxy", )?; check_agg_store::( - config, + storage, TA_SIGNER_SERVER_NS, "TA Signer", )?; - check_openssl_keys(config)?; + check_openssl_keys(storage)?; Ok(()) } -fn check_openssl_keys(config: &Config) -> UpgradeResult<()> { +fn check_openssl_keys(storage: &StorageSystem) -> UpgradeResult<()> { info!(""); info!("Verify: OpenSSL keys"); + // XXX This seems to assume that OpenSSL keys are always in the store? let open_ssl_signer = OpenSslSigner::build( - &config.storage_uri, - "test", - None, - ) - .map_err(|e| { + storage, &OpenSslSignerConfig::default(), "test", None, + ).map_err(|e| { UpgradeError::Custom(format!("Cannot create openssl signer: {e}")) })?; - let keys_key_store = KeyValueStore::create(&config.storage_uri, KEYS_NS)?; + let keys_key_store = storage.open(KEYS_NS)?; for key in keys_key_store.keys(None, "")? { let key_id = @@ -159,14 +161,15 @@ fn check_openssl_keys(config: &Config) -> UpgradeResult<()> { } pub fn check_agg_store( - config: &Config, + storage: &StorageSystem, ns: &Ident, name: &str, ) -> UpgradeResult> { info!(""); info!("Verify: {name}"); - let store: AggregateStore = - AggregateStore::create(&config.storage_uri, ns, false)?; + let store: AggregateStore = AggregateStore::create( + storage, ns, false + )?; if !store.list()?.is_empty() { store.warm()?; info!("Ok"); @@ -177,13 +180,13 @@ pub fn check_agg_store( } fn check_wal_store( - config: &Config, + storage: &StorageSystem, ns: &Ident, name: &str, ) -> UpgradeResult<()> { info!(""); info!("Verify: {name}"); - let store: WalStore = WalStore::create(&config.storage_uri, ns)?; + let store: WalStore = WalStore::create(storage, ns)?; if !store.list()?.is_empty() { store.warm()?; info!("Ok"); @@ -194,8 +197,8 @@ fn check_wal_store( } fn copy_data_for_migration( - config: &Config, - target_storage: &Url, + target_storage: &StorageSystem, + source_storage: &Url, ) -> UpgradeResult<()> { const NAMESPACES: &[&Ident] = &[ Ident::make("ca_objects"), @@ -209,13 +212,11 @@ fn copy_data_for_migration( Ident::make("ta_signer"), ]; for namespace in NAMESPACES { - let source_kv_store = KeyValueStore::create( - &config.storage_uri, namespace + let source_kv_store = target_storage.open_uri( + source_storage, namespace )?; if !source_kv_store.is_empty()? { - let target_kv_store = KeyValueStore::create( - target_storage, namespace - )?; + let target_kv_store = target_storage.open(namespace)?; target_kv_store.import(&source_kv_store)?; } } @@ -253,7 +254,7 @@ pub mod tests { // Create an in-memory target store to migrate to let target_store = test::mem_storage(); - migrate(config, target_store).unwrap(); + migrate(config, &target_store).unwrap(); } } diff --git a/src/upgrades/mod.rs b/src/upgrades/mod.rs index ec198b81..650a18a3 100644 --- a/src/upgrades/mod.rs +++ b/src/upgrades/mod.rs @@ -23,7 +23,10 @@ use crate::{ Storable, StoredCommand, WalStoreError, WithStorableDetails, }, - storage::{Ident, KeyValueError, KeyValueStore}, + storage::{ + Ident, KeyValueError, OpenStoreError, KeyValueStore, + StorageConnectError, StorageSystem, + }, version::KrillVersion, KrillResult, }, @@ -222,6 +225,18 @@ impl From for UpgradeError { } } +impl From for UpgradeError { + fn from(e: StorageConnectError) -> Self { + UpgradeError::custom(e) + } +} + +impl From for UpgradeError { + fn from(e: OpenStoreError) -> Self { + UpgradeError::custom(e) + } +} + impl From for UpgradeError { fn from(e: WalStoreError) -> Self { UpgradeError::WalStoreError(e) @@ -811,6 +826,7 @@ pub trait UpgradeAggregateStorePre0_14 { /// time. After this, the migration will be finalised. pub fn prepare_upgrade_data_migrations( mode: UpgradeMode, + storage: &StorageSystem, config: &Config, properties_manager: &PropertiesManager, ) -> UpgradeResult> { @@ -823,9 +839,9 @@ pub fn prepare_upgrade_data_migrations( // just do at startup. It is done here, because in effect it *is* a data // migration. #[cfg(feature = "hsm")] - record_preexisting_openssl_keys_in_signer_mapper(config)?; + record_preexisting_openssl_keys_in_signer_mapper(storage, config)?; - match upgrade_versions(config, properties_manager)? { + match upgrade_versions(storage, properties_manager)? { None => Ok(None), Some(versions) => { info!( @@ -840,9 +856,7 @@ pub fn prepare_upgrade_data_migrations( // easily be migrated to the new setup in 0.13.0. // Well.. it could be done, if there would be a strong use // case to put in the effort, but there really isn't. - let ca_kv_store = KeyValueStore::create( - &config.storage_uri, CASERVER_NS - )?; + let ca_kv_store = storage.open(CASERVER_NS)?; if ca_kv_store.has_scope(const { Ident::make("ta") })? { return Err(UpgradeError::OldTaMigration); } @@ -859,19 +873,20 @@ pub fn prepare_upgrade_data_migrations( } else if versions.from < KrillVersion::candidate(0, 10, 0, 1) { // Complex migrations involving command / event conversions - pubd::pre_0_10_0::PublicationServerRepositoryAccessMigration::upgrade(mode, config, &versions)?; - let aspa_configs = - ca::pre_0_10_0::CasMigration::upgrade(mode, config)?; + pubd::pre_0_10_0::PublicationServerRepositoryAccessMigration::upgrade(mode, storage, &versions)?; + let aspa_configs = ca::pre_0_10_0::CasMigration::upgrade( + mode, storage + )?; // The way that pubd objects were stored was changed as well // (since 0.13.0) - pubd::migrate_pre_0_12_pubd_objects(config)?; + pubd::migrate_pre_0_12_pubd_objects(storage)?; // Migrate remaining aggregate stores used in < 0.10.0 to the // new format in 0.14.0 where we combine // commands and events into a single key-value pair. pre_0_14_0::UpgradeAggregateStoreSignerInfo::upgrade( - SIGNERS_NS, mode, config, + SIGNERS_NS, mode, storage )?; Ok(Some(UpgradeReport::new(aspa_configs, true, versions))) @@ -887,38 +902,38 @@ pub fn prepare_upgrade_data_migrations( ); // The pubd objects storage changed in 0.13.0 - pubd::migrate_pre_0_12_pubd_objects(config)?; + pubd::migrate_pre_0_12_pubd_objects(storage)?; // Migrate aggregate stores used in < 0.12.0 to the new format // in 0.14.0 where we combine commands and // events into a single key-value pair. pre_0_14_0::UpgradeAggregateStoreSignerInfo::upgrade( - SIGNERS_NS, mode, config, + SIGNERS_NS, mode, storage + )?; + let aspa_configs = ca::pre_0_14_0::CasMigration::upgrade( + mode, storage )?; - let aspa_configs = - ca::pre_0_14_0::CasMigration::upgrade(mode, config)?; pubd::pre_0_14_0::UpgradeAggregateStoreRepositoryAccess::upgrade( PUBSERVER_NS, mode, - config, + storage, )?; Ok(Some(UpgradeReport::new(aspa_configs, true, versions))) } else if versions.from < KrillVersion::candidate(0, 13, 0, 0) { - pubd::migrate_0_12_pubd_objects(config)?; + pubd::migrate_0_12_pubd_objects(storage)?; // Migrate aggregate stores used in < 0.13.0 to the new format // in 0.14.0 where we combine commands and // events into a single key-value pair. pre_0_14_0::UpgradeAggregateStoreSignerInfo::upgrade( - SIGNERS_NS, mode, config, + SIGNERS_NS, mode, storage + )?; + let aspa_configs = ca::pre_0_14_0::CasMigration::upgrade( + mode, storage )?; - let aspa_configs = - ca::pre_0_14_0::CasMigration::upgrade(mode, config)?; pubd::pre_0_14_0::UpgradeAggregateStoreRepositoryAccess::upgrade( - PUBSERVER_NS, - mode, - config, + PUBSERVER_NS, mode, storage )?; Ok(Some(UpgradeReport::new(aspa_configs, true, versions))) @@ -926,25 +941,20 @@ pub fn prepare_upgrade_data_migrations( // Migrate aggregate stores used in < 0.14.0 to the new format // in 0.14.0 where we combine commands and // events into a single key-value pair. - let aspa_configs = - ca::pre_0_14_0::CasMigration::upgrade(mode, config)?; + let aspa_configs = ca::pre_0_14_0::CasMigration::upgrade( + mode, storage + )?; pubd::pre_0_14_0::UpgradeAggregateStoreRepositoryAccess::upgrade( - PUBSERVER_NS, - mode, - config, + PUBSERVER_NS, mode, storage )?; pre_0_14_0::UpgradeAggregateStoreSignerInfo::upgrade( - SIGNERS_NS, mode, config, + SIGNERS_NS, mode, storage )?; pre_0_14_0::UpgradeAggregateStoreTrustAnchorSigner::upgrade( - TA_SIGNER_SERVER_NS, - mode, - config, + TA_SIGNER_SERVER_NS, mode, storage )?; pre_0_14_0::UpgradeAggregateStoreTrustAnchorProxy::upgrade( - TA_PROXY_SERVER_NS, - mode, - config, + TA_PROXY_SERVER_NS, mode, storage )?; Ok(Some(UpgradeReport::new(aspa_configs, true, versions))) @@ -966,7 +976,7 @@ pub fn prepare_upgrade_data_migrations( /// - make the prepared data current pub fn finalise_data_migration( upgrade: &UpgradeVersions, - config: &Config, + storage: &StorageSystem, properties_manager: &PropertiesManager, ) -> KrillResult<()> { // For each NS @@ -997,24 +1007,20 @@ pub fn finalise_data_migration( ] { // Check if there is a non-empty upgrade store for this namespace // that would need to be migrated. - let mut upgrade_store = - KeyValueStore::create_upgrade_store(&config.storage_uri, ns)?; + let mut upgrade_store = storage.open_upgrade(ns)?; if !upgrade_store.is_empty()? { info!("Migrate new data for {ns} and archive old"); - let mut current_store = - KeyValueStore::create(&config.storage_uri, ns)?; + let mut current_store = storage.open(ns)?; if !current_store.is_empty()? { - current_store.migrate_to_archive(&config.storage_uri, ns)?; + current_store.migrate_to_archive(storage.default_uri(), ns)?; } - upgrade_store.migrate_to_current(&config.storage_uri, ns)?; + upgrade_store.migrate_to_current(storage.default_uri(), ns)?; } else { // No migration needed, but check if we have a current store // for this namespace that still includes a version file. If // so, remove it. - let current_store = KeyValueStore::create( - &config.storage_uri, ns - )?; + let current_store = storage.open(ns)?; if current_store.has(None, VERSION_KEY)? { debug!("Removing excess version key in ns: {ns}"); current_store.drop_key(None, VERSION_KEY)?; @@ -1051,10 +1057,10 @@ pub fn finalise_data_migration( /// adding the keys one by one to the mapping in the signer store, if any. #[cfg(feature = "hsm")] fn record_preexisting_openssl_keys_in_signer_mapper( + storage: &StorageSystem, config: &Config, ) -> Result<(), UpgradeError> { - let signers_key_store = - KeyValueStore::create(&config.storage_uri, SIGNERS_NS)?; + let signers_key_store = storage.open(SIGNERS_NS)?; if signers_key_store.is_empty()? { let mut num_recorded_keys = 0; // If the key value store for the "signers" namespace is empty, then @@ -1062,19 +1068,16 @@ fn record_preexisting_openssl_keys_in_signer_mapper( // from a previous krill installation (earlier version, or a custom // build that has the hsm feature disabled.) - let keys_key_store = - KeyValueStore::create(&config.storage_uri, KEYS_NS)?; + let keys_key_store = storage.open(KEYS_NS)?; info!( "Mapping OpenSSL signer keys, using uri: {}", - config.storage_uri + storage.default_uri() ); let probe_interval = std::time::Duration::from_secs(config.signer_probe_retry_seconds); let krill_signer = crate::commons::crypto::KrillSignerBuilder::new( - &config.storage_uri, - probe_interval, - &config.signers, + storage, probe_interval, &config.signers, ) .with_default_signer(config.default_signer()) .with_one_off_signer(config.one_off_signer()) @@ -1180,7 +1183,7 @@ pub async fn post_start_upgrade( /// - if the code is the same version then we do not upgrade /// - if the code is older then we need to error out fn upgrade_versions( - config: &Config, + storage: &StorageSystem, properties_manager: &PropertiesManager, ) -> Result, UpgradeError> { if properties_manager.is_initialized() { @@ -1217,7 +1220,7 @@ fn upgrade_versions( PUBSERVER_NS, PUBSERVER_CONTENT_NS, ] { - let kv_store = KeyValueStore::create(&config.storage_uri, ns)?; + let kv_store = storage.open(ns)?; if let Some(key_store_version) = kv_store.get::(None, VERSION_KEY)? { @@ -1282,13 +1285,13 @@ mod tests { copy_folder(base_dir, &temp_dir); // Copy data for the given names spaces into memory for testing. - let mem_storage_base_uri = test::mem_storage(); + let mem_storage = test::mem_storage(); // This is needed for tls_dir etc, but will be ignored here. let bogus_path = PathBuf::from("/dev/null"); let mut config = Config::test( - &mem_storage_base_uri, + mem_storage.default_uri(), Some(&bogus_path), false, false, false, false, ); @@ -1301,23 +1304,22 @@ mod tests { for ns in namespaces { let namespace = Ident::from_str(ns).unwrap(); - let source_store = KeyValueStore::create( + let source_store = mem_storage.open_uri( &source_url, namespace ).unwrap(); - let target_store = KeyValueStore::create( - &mem_storage_base_uri, namespace - ).unwrap(); + let target_store = mem_storage.open(namespace).unwrap(); target_store.import(&source_store).unwrap(); } let properties_manager = PropertiesManager::create( - &config.storage_uri, + &mem_storage, config.use_history_cache, ).unwrap(); prepare_upgrade_data_migrations( UpgradeMode::PrepareOnly, + &mem_storage, &config, &properties_manager, ).unwrap().unwrap(); @@ -1326,13 +1328,14 @@ mod tests { // again. let report = prepare_upgrade_data_migrations( UpgradeMode::PrepareToFinalise, + &mem_storage, &config, &properties_manager, ).unwrap().unwrap(); finalise_data_migration( report.versions(), - &config, + &mem_storage, &properties_manager, ).unwrap(); } @@ -1461,18 +1464,16 @@ mod tests { ).unwrap(); // Copy test data into test storage - let mem_storage_base_uri = test::mem_storage(); + let mem_storage = test::mem_storage(); let source_url = Url::parse(&format!( "local://{}", temp_dir.path().to_str().unwrap() )).unwrap(); - let source_store = KeyValueStore::create( + let source_store = mem_storage.open_uri( &source_url, KEYS_NS ).unwrap(); - let target_store = KeyValueStore::create( - &mem_storage_base_uri, KEYS_NS - ).unwrap(); + let target_store = mem_storage.open(KEYS_NS).unwrap(); target_store.import(&source_store).unwrap(); // This is needed for tls_dir etc, but will be ignored here. @@ -1558,12 +1559,13 @@ mod tests { &format!("local://{}", &temp_dir.path().to_str().unwrap())) .unwrap(); - let source_store = - KeyValueStore::create(&source_dir_url, STATUS_NS).unwrap(); + let test_storage = test::mem_storage(); - let test_storage_uri = test::mem_storage(); - let status_kv_store = - KeyValueStore::create(&test_storage_uri, STATUS_NS).unwrap(); + let source_store = test_storage.open_uri( + &source_dir_url, STATUS_NS + ).unwrap(); + + let status_kv_store = test_storage.open(STATUS_NS).unwrap(); // copy the source KV store (files) into the test KV store (in memory) status_kv_store.import(&source_store).unwrap(); @@ -1578,8 +1580,7 @@ mod tests { // Initialise the StatusStore using the new (in memory) storage, // and migrate the data. - let store = - CaStatusStore::create(&test_storage_uri, STATUS_NS).unwrap(); + let store = CaStatusStore::create(&test_storage, STATUS_NS).unwrap(); let testbed = CaHandle::from_str("testbed").unwrap(); // Get the migrated status for testbed and verify that it's equivalent diff --git a/src/upgrades/pre_0_14_0.rs b/src/upgrades/pre_0_14_0.rs index 4d3f1088..38549344 100644 --- a/src/upgrades/pre_0_14_0.rs +++ b/src/upgrades/pre_0_14_0.rs @@ -17,9 +17,8 @@ use crate::{ Aggregate, AggregateStore, Storable, StoredCommand, StoredCommandBuilder, WithStorableDetails, }, - storage::{Ident, KeyValueStore}, + storage::{Ident, KeyValueStore, StorageSystem}, }, - config::Config, server::{ properties::Properties, }, @@ -299,23 +298,19 @@ impl GenericUpgradeAggregateStore { pub fn upgrade( name_space: &Ident, mode: UpgradeMode, - config: &Config, + storage: &StorageSystem, ) -> UpgradeResult { - let current_kv_store = - KeyValueStore::create(&config.storage_uri, name_space)?; + let current_kv_store = storage.open(name_space)?; if current_kv_store.scopes()?.is_empty() { // nothing to do here Ok(AspaMigrationConfigs::default()) } else { - let new_kv_store = KeyValueStore::create_upgrade_store( - &config.storage_uri, - name_space, - )?; + let new_kv_store = storage.open_upgrade(name_space)?; let new_agg_store = AggregateStore::::create_upgrade_store( - &config.storage_uri, + storage, name_space, - config.use_history_cache, + false, )?; let store_migration = GenericUpgradeAggregateStore {