Introduce a StorageSystem for global storage system state.

This commit is contained in:
Martin Hoffmann
2025-11-17 18:16:17 +01:00
parent 38205078c8
commit 57d68435fe
39 changed files with 592 additions and 500 deletions
+5 -4
View File
@@ -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();
+37 -3
View File
@@ -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!();
+6 -3
View File
@@ -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<Self, SignerClientError> {
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 {
@@ -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<Arc<SignerMapper>>,
@@ -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<Arc<SignerMapper>>,
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<SignerProvider> {
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")]
@@ -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<SignerMapper> {
pub fn build(storage: &StorageSystem) -> KrillResult<SignerMapper> {
let store = AggregateStore::<SignerInfo>::create(
storage_uri,
SIGNERS_NS,
true,
storage, SIGNERS_NS, true,
)?;
Ok(SignerMapper { store })
}
@@ -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<Arc<SignerMapper>>,
) -> Result<Self, SignerError> {
let keys_store = Self::init_keys_store(storage_uri)?;
) -> Result<Self, OpenStoreError> {
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<KeyValueStore, SignerError> {
let store = KeyValueStore::create(storage_uri, KEYS_NS)
.map_err(|_| SignerError::InvalidStorage(storage_uri.clone()))?;
Ok(store)
storage: &StorageSystem,
conf: &OpenSslSignerConfig,
) -> Result<KeyValueStore, OpenStoreError> {
match &conf.keys_storage_uri {
Some(uri) => storage.open_uri(uri, KEYS_NS),
None => storage.open(KEYS_NS)
}
}
fn build_key(&self) -> Result<KeyIdentifier, SignerError> {
@@ -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();
+9 -8
View File
@@ -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<A: Aggregate> AggregateStore<A> {
/// 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<Self, AggregateStoreError> {
) -> Result<Self, OpenStoreError> {
Ok(Self::create_from_kv(
KeyValueStore::create(storage_uri, namespace)?, use_history_cache
storage.open(namespace)?, use_history_cache
))
}
@@ -90,12 +91,12 @@ impl<A: Aggregate> AggregateStore<A> {
/// 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<Self, AggregateStoreError> {
) -> Result<Self, OpenStoreError> {
Ok(Self::create_from_kv(
KeyValueStore::create_upgrade_store(storage_uri, namespace)?,
storage.open_upgrade(namespace)?,
use_history_cache,
))
}
+6 -5
View File
@@ -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<T: WalSupport> {
impl<T: WalSupport> WalStore<T> {
/// Creates a new store using the given storage URL and namespace.
pub fn create(
storage_uri: &Url,
storage: &StorageSystem,
namespace: &Ident,
) -> Result<Self, WalStoreError> {
) -> Result<Self, OpenStoreError> {
Ok(WalStore {
kv: KeyValueStore::create(storage_uri, namespace)?,
kv: storage.open(namespace)?,
cache: RwLock::new(HashMap::new()),
})
}
+29 -23
View File
@@ -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<Self, Error> {
) -> Result<Self, OpenStoreError> {
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") };
+4 -1
View File
@@ -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;
+104 -18
View File
@@ -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<Self, StorageConnectError> {
Ok(Self { storage_uri })
}
/// Opens the default store with the given namespace.
pub fn open(
&self, namespace: &Ident
) -> Result<KeyValueStore, OpenStoreError> {
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, OpenStoreError> {
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, OpenStoreError> {
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<Self, KeyValueError> {
@@ -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, KeyValueError> {
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<StorageConnectError> 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<OpenStoreError> for Error {
fn from(src: OpenStoreError) -> Error {
Error::custom(format_args!("{}", src))
}
}
//------------ KeyValueError -------------------------------------------------
/// This type defines possible Errors for KeyStore
+17 -13
View File
@@ -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()
}
}
+7 -6
View File
@@ -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<F>(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 {
+1 -18
View File
@@ -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> {
KeyValueStore::create(&self.storage_uri, PROPERTIES_NS)
.map_err(Error::KeyValueError)
}
pub fn key_value_store(
&self,
namespace: &Ident,
) -> KrillResult<KeyValueStore> {
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<PathBuf> {
+8 -2
View File
@@ -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<Config>,
) -> KrillResult<Self> {
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))
)
}
+3 -4
View File
@@ -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<CryptState> {
let store = config.key_value_store(CRYPT_STATE_NS)?;
pub(crate) fn crypt_init(storage: &StorageSystem) -> KrillResult<CryptState> {
let store = storage.open(CRYPT_STATE_NS)?;
if let Some(state) = store.get(None, CRYPT_STATE_KEY)? {
Ok(state)
@@ -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<Self> {
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<crypt::CryptState> {
fn init_session_key(
storage: &StorageSystem
) -> KrillResult<crypt::CryptState> {
debug!("Initializing login session encryption key");
crypt::crypt_init(config)
crypt::crypt_init(storage)
}
/// Parse HTTP Basic Authorization header
@@ -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<Config>,
) -> KrillResult<Self> {
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<CryptState> {
fn init_session_key(storage: &StorageSystem) -> KrillResult<CryptState> {
debug!("Initializing session encryption key");
crypt::crypt_init(config)
crypt::crypt_init(storage)
}
fn oidc_conf(&self) -> KrillResult<&ConfigAuthOpenIDConnect> {
+4 -3
View File
@@ -22,7 +22,7 @@ use super::response::{HyperResponse, HttpResponse};
/// The Krill HTTP server.
pub struct HttpServer {
/// The Krill “business logic.”
krill: KrillManager,
krill: Arc<KrillManager>,
/// 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<KrillManager>,
config: Arc<Config>
) -> KrillResult<Arc<Self>> {
let authorizer = Authorizer::new(krill.storage(), config.clone())?;
Ok(Self {
krill,
authorizer: Authorizer::new(config.clone())?,
authorizer,
config,
started: Timestamp::now(),
}.into())
+9 -8
View File
@@ -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() {
+7 -15
View File
@@ -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<Config>,
storage: &StorageSystem,
tasks: Arc<TaskQueue>,
signer: Arc<KrillSigner>,
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::<CertAuth>::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::<TrustAnchorProxy>::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,
+6 -5
View File
@@ -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<KrillSigner>,
) -> KrillResult<Self> {
let store = KeyValueStore::create(storage_uri, CA_OBJECTS_NS)?;
) -> Result<Self, OpenStoreError> {
let store = storage.open(CA_OBJECTS_NS)?;
Ok(CaObjectsStore {
store,
signer,
+3 -4
View File
@@ -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<Self> {
let store = KeyValueStore::create(storage_uri, namespace)?;
let store = storage.open(namespace)?;
let cache = RwLock::new(HashMap::new());
let store = Self { store, cache };
+7 -6
View File
@@ -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::<CertAuth>(config, CASERVER_NS, "CAs")?;
pub fn check_ca_objects(
storage: &StorageSystem, config: &Config
) -> UpgradeResult<()> {
let ca_store = check_agg_store::<CertAuth>(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()?;
+12 -21
View File
@@ -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<AspaMigrationConfigs> {
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::<CertAuth>::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<Self, UpgradeError> {
fn create(storage: &StorageSystem) -> Result<Self, OpenStoreError> {
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)?
})
}
+5 -12
View File
@@ -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<AspaMigrationConfigs> {
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::<CertAuth>::create_upgrade_store(
&config.storage_uri,
CASERVER_NS,
config.use_history_cache,
storage, CASERVER_NS, false,
)?;
CasMigration {
+35 -31
View File
@@ -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<RepositoryManager>,
pub(super) repo_manager: Arc<RepositoryManager>,
// Handles the internal TA and/or CAs
ca_manager: Arc<ca::CaManager>,
pub(super) ca_manager: Arc<ca::CaManager>,
// Handles the internal TA and/or CAs
bgp_analyser: Arc<BgpAnalyser>,
pub(super) bgp_analyser: Arc<BgpAnalyser>,
// Shared message queue
mq: Arc<TaskQueue>,
pub(super) tasks: Arc<TaskQueue>,
// System actor
system_actor: Actor,
pub(super) system_actor: Actor,
pub config: Arc<Config>,
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<Config>) -> KrillResult<Self> {
pub async fn build(
storage: StorageSystem, config: Arc<Config>
) -> KrillResult<Self> {
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
}
+3 -4
View File
@@ -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<Self> {
pub fn new(storage: &StorageSystem) -> KrillResult<Self> {
Ok(TaskQueue {
q: Queue::create(storage_uri, TASK_QUEUE_NS)?,
q: Queue::create(storage, TASK_QUEUE_NS)?,
})
}
}
+11 -9
View File
@@ -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<Self> {
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 {
+7 -6
View File
@@ -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<Self> {
/// Creates a new repository access proxy
pub fn create(
storage: &StorageSystem,
use_history_cache: bool
) -> KrillResult<Self> {
let store = AggregateStore::<RepositoryAccess>::create(
&config.storage_uri,
PUBSERVER_NS,
config.use_history_cache,
storage, PUBSERVER_NS, use_history_cache,
)?;
let key = MyHandle::from_str(PUBSERVER_DFLT).unwrap();
+4 -3
View File
@@ -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<Self> {
pub fn create(storage: &StorageSystem) -> KrillResult<Self> {
let store = Arc::new(WalStore::create(
&config.storage_uri, PUBSERVER_CONTENT_NS,
storage, PUBSERVER_CONTENT_NS,
)?);
store.warm()?;
+19 -14
View File
@@ -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<Config>,
tasks: Arc<TaskQueue>,
signer: Arc<KrillSigner>,
) -> Result<Self, Error> {
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 =
+11 -15
View File
@@ -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<bool> {
pub fn migrate_0_12_pubd_objects(
storage: &StorageSystem
) -> KrillResult<bool> {
let old_store: WalStore<OldRepositoryContent> = 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<bool> {
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<bool> {
/// 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::<OldRepositoryContent>(
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") },
@@ -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 {
+17 -48
View File
@@ -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<TaskQueue>,
ca_manager: Arc<CaManager>,
repo_manager: Arc<RepositoryManager>,
config: Arc<Config>,
system_actor: Actor,
started: Timestamp,
}
impl Scheduler {
pub fn build(
tasks: Arc<TaskQueue>,
ca_manager: Arc<CaManager>,
repo_manager: Arc<RepositoryManager>,
config: Arc<Config>,
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<TaskResult, FatalError> {
fn update_aggregate_store_snapshots<A: Aggregate>(
storage_uri: &Url,
storage: &StorageSystem,
namespace: &Ident,
) {
match AggregateStore::<A>::create(storage_uri, namespace, false) {
match AggregateStore::<A>::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<W: WalSupport>(
storage_uri: &Url,
storage: &StorageSystem,
namespace: &Ident,
) {
match WalStore::<W>::create(storage_uri, namespace) {
match WalStore::<W>::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::<CertAuth>(
&self.config.storage_uri,
&self.storage,
CASERVER_NS,
);
update_aggregate_store_snapshots::<SignerInfo>(
&self.config.storage_uri,
&self.storage,
SIGNERS_NS,
);
update_aggregate_store_snapshots::<Properties>(
&self.config.storage_uri,
&self.storage,
PROPERTIES_NS,
);
update_aggregate_store_snapshots::<RepositoryAccess>(
&self.config.storage_uri,
&self.storage,
PUBSERVER_NS,
);
update_wal_store_snapshots::<RepositoryContent>(
&self.config.storage_uri,
&self.storage,
PUBSERVER_CONTENT_NS,
);
+7 -4
View File
@@ -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<Arc<KrillSigner>, ConfigError> {
pub fn signer(
&self, storage: &StorageSystem
) -> Result<Arc<KrillSigner>, 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();
})
}
}
+43 -42
View File
@@ -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::<Properties>(config, PROPERTIES_NS, "Properties")?;
check_agg_store::<SignerInfo>(config, SIGNERS_NS, "Signer")?;
fn verify_target_data(
storage: &StorageSystem, config: &Config
) -> UpgradeResult<()> {
check_agg_store::<Properties>(storage, PROPERTIES_NS, "Properties")?;
check_agg_store::<SignerInfo>(storage, SIGNERS_NS, "Signer")?;
check_ca_objects(config)?;
check_ca_objects(storage, config)?;
check_agg_store::<RepositoryAccess>(
config,
storage,
PUBSERVER_NS,
"Publication Server Access",
)?;
check_wal_store::<RepositoryContent>(
config,
storage,
PUBSERVER_CONTENT_NS,
"Publication Server Objects",
)?;
check_agg_store::<TrustAnchorProxy>(
config,
storage,
TA_PROXY_SERVER_NS,
"TA Proxy",
)?;
check_agg_store::<TrustAnchorSigner>(
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<A: Aggregate>(
config: &Config,
storage: &StorageSystem,
ns: &Ident,
name: &str,
) -> UpgradeResult<AggregateStore<A>> {
info!("");
info!("Verify: {name}");
let store: AggregateStore<A> =
AggregateStore::create(&config.storage_uri, ns, false)?;
let store: AggregateStore<A> = AggregateStore::create(
storage, ns, false
)?;
if !store.list()?.is_empty() {
store.warm()?;
info!("Ok");
@@ -177,13 +180,13 @@ pub fn check_agg_store<A: Aggregate>(
}
fn check_wal_store<W: WalSupport>(
config: &Config,
storage: &StorageSystem,
ns: &Ident,
name: &str,
) -> UpgradeResult<()> {
info!("");
info!("Verify: {name}");
let store: WalStore<W> = WalStore::create(&config.storage_uri, ns)?;
let store: WalStore<W> = WalStore::create(storage, ns)?;
if !store.list()?.is_empty() {
store.warm()?;
info!("Ok");
@@ -194,8 +197,8 @@ fn check_wal_store<W: WalSupport>(
}
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();
}
}
+76 -75
View File
@@ -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<AggregateStoreError> for UpgradeError {
}
}
impl From<StorageConnectError> for UpgradeError {
fn from(e: StorageConnectError) -> Self {
UpgradeError::custom(e)
}
}
impl From<OpenStoreError> for UpgradeError {
fn from(e: OpenStoreError) -> Self {
UpgradeError::custom(e)
}
}
impl From<WalStoreError> 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<Option<UpgradeReport>> {
@@ -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<Option<UpgradeVersions>, 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::<KrillVersion>(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
+6 -11
View File
@@ -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<A: Aggregate> GenericUpgradeAggregateStore<A> {
pub fn upgrade(
name_space: &Ident,
mode: UpgradeMode,
config: &Config,
storage: &StorageSystem,
) -> UpgradeResult<AspaMigrationConfigs> {
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::<A>::create_upgrade_store(
&config.storage_uri,
storage,
name_space,
config.use_history_cache,
false,
)?;
let store_migration = GenericUpgradeAggregateStore {