use std::convert::TryInto;
use std::str::FromStr;
use chrono::Duration;
use rpki::{ca::idexchange::CaHandle, repository::x509::Time};
use crate::daemon::ca::CaObjects;
use crate::{
commons::{
api::StorableCaCommand,
eventsourcing::{AggregateStore, KeyStoreKey, KeyValueStore, StoredCommand, StoredValueInfo},
},
constants::{CASERVER_DIR, CA_OBJECTS_DIR, KRILL_VERSION},
daemon::{
ca::{CaEvt, CertAuth, IniDet},
config::Config,
},
upgrades::{
pre_0_10_0::{OldCaEvt, OldCaIni},
PrepareUpgradeError, UpgradeMode, UpgradeResult, UpgradeStore,
},
};
use super::OldCaObjects;
struct CaObjectsMigration {
current_store: KeyValueStore,
new_store: KeyValueStore,
}
impl CaObjectsMigration {
fn create(config: &Config) -> Result<Self, PrepareUpgradeError> {
let current_store = KeyValueStore::disk(&config.data_dir, CA_OBJECTS_DIR)?;
let new_store = KeyValueStore::disk(&config.upgrade_data_dir(), CA_OBJECTS_DIR)?;
Ok(CaObjectsMigration {
current_store,
new_store,
})
}
fn prepare_new_data_for(&self, ca: &CaHandle) -> Result<(), PrepareUpgradeError> {
let key = KeyStoreKey::simple(format!("{}.json", ca));
if let Some(old_objects) = self.current_store.get::<OldCaObjects>(&key)? {
let converted: CaObjects = old_objects.try_into()?;
self.new_store.store(&key, &converted)?;
}
Ok(())
}
}
pub struct CasMigration {
current_kv_store: KeyValueStore,
new_kv_store: KeyValueStore,
new_agg_store: AggregateStore<CertAuth>,
ca_objects_migration: CaObjectsMigration,
}
impl CasMigration {
pub fn prepare(mode: UpgradeMode, config: &Config) -> UpgradeResult<()> {
let current_kv_store = KeyValueStore::disk(&config.data_dir, CASERVER_DIR)?;
let new_kv_store = KeyValueStore::disk(&config.upgrade_data_dir(), CASERVER_DIR)?;
let new_agg_store = AggregateStore::<CertAuth>::disk(&config.upgrade_data_dir(), CASERVER_DIR)?;
let ca_objects_migration = CaObjectsMigration::create(config)?;
CasMigration {
current_kv_store,
new_kv_store,
new_agg_store,
ca_objects_migration,
}
.prepare_new_data(mode)
}
}
impl UpgradeStore for CasMigration {
fn needs_migrate(&self) -> Result<bool, PrepareUpgradeError> {
unreachable!("This is checked in upgrades/mod.rs")
}
fn prepare_new_data(&self, mode: UpgradeMode) -> Result<(), PrepareUpgradeError> {
self.preparation_store_prepare()?;
info!(
"Prepare upgrading CA command and event data to Krill version {}",
KRILL_VERSION
);
for scope in self.current_kv_store.scopes()? {
let handle = CaHandle::from_str(&scope)
.map_err(|_| PrepareUpgradeError::Custom(format!("Found invalid CA handle '{}'", scope)))?;
let mut data_upgrade_info = self.data_upgrade_info(&scope)?;
if data_upgrade_info.last_event == 0 {
let init_key = Self::event_key(&scope, 0);
let old_init: OldCaIni = self.get(&init_key)?;
let (id, _, old_ini_det) = old_init.unpack();
let ini = IniDet::new(&id, old_ini_det.into());
self.new_kv_store.store(&init_key, &ini)?;
}
let old_cmd_keys = self.command_keys(&scope, data_upgrade_info.last_command)?;
let total_commands = old_cmd_keys.len();
if data_upgrade_info.last_command == 0 {
info!("Will migrate {} commands for CA '{}'", total_commands, handle);
} else {
info!(
"Will resume migration of {} remaining commands for CA '{}'",
total_commands, handle
);
}
let info_key = KeyStoreKey::scoped(scope.to_string(), "info.json".to_string());
let old_info: StoredValueInfo = self
.current_kv_store
.get(&info_key)?
.ok_or_else(|| PrepareUpgradeError::Custom(format!("Cannot parse old info file: {}", info_key)))?;
let mut total_migrated = 0;
let time_started = Time::now();
for cmd_key in old_cmd_keys {
let cmd: StoredCommand<StorableCaCommand> = self.get(&cmd_key)?;
if let Some(event_versions) = cmd.effect().events() {
for v in event_versions {
let event_key = Self::event_key(&scope, *v);
trace!(" +- event: {}", event_key);
let evt: OldCaEvt = self.current_kv_store.get(&event_key)?.ok_or_else(|| {
PrepareUpgradeError::Custom(format!("Cannot parse old event: {}", event_key))
})?;
let evt: CaEvt = evt.try_into()?;
self.new_kv_store.store(&event_key, &evt)?;
data_upgrade_info.last_event = *v;
}
}
self.new_kv_store.store(&cmd_key, &cmd)?;
data_upgrade_info.last_command += 1;
data_upgrade_info.last_update = cmd.time();
self.update_data_upgrade_info(&scope, &data_upgrade_info)?;
total_migrated += 1;
if total_migrated % 100 == 0 {
let mut time_passed = (Time::now().timestamp() - time_started.timestamp()) as usize;
if time_passed == 0 {
time_passed = 1; }
let migrated_per_second: f64 = total_migrated as f64 / time_passed as f64;
let expected_seconds = (total_commands as f64 / migrated_per_second) as i64;
let eta = time_started + Duration::seconds(expected_seconds);
info!(
" migrated {} commands, expect to finish: {}",
total_migrated,
eta.to_rfc3339()
);
}
}
info!("Finished migrating commands for CA '{}'", scope);
{
let info = StoredValueInfo {
snapshot_version: data_upgrade_info.last_event + 1,
last_event: data_upgrade_info.last_event,
last_command: data_upgrade_info.last_command,
last_update: data_upgrade_info.last_update,
};
self.new_kv_store.store(&info_key, &info)?;
if mode.is_finalise() {
if info.last_command != old_info.last_command || info.last_event != old_info.last_event {
return Err(PrepareUpgradeError::custom(
format!("New info.json does not match old info.json when upgrading CA '{}'. Please downgrade to the previous version and provide a bug report to rpki-team@nlnetlabs.nl.", handle),
));
}
}
}
info!("Will verify the migration by rebuilding CA '{}' events", &scope);
let ca = self.new_agg_store.get_latest(&handle).map_err(|e| {
PrepareUpgradeError::Custom(format!(
"Could not rebuild state after migrating CA '{}'! Error was: {}.",
handle, e
))
})?;
self.new_agg_store.store_snapshot(&handle, ca.as_ref()).map_err(|e| {
PrepareUpgradeError::Custom(format!(
"Could not save snapshot for CA '{}' after migration! Disk full?!? Error was: {}.",
handle, e
))
})?;
info!("Will migrate the current repository objects for CA '{}'", handle);
self.ca_objects_migration.prepare_new_data_for(&handle)?;
info!("Verified migration of CA '{}'", handle);
}
match mode {
UpgradeMode::PrepareOnly => {
info!(
"Prepared migrating CAs to Krill version {}. Will save progress for final upgrade when Krill restarts.",
KRILL_VERSION
);
}
UpgradeMode::PrepareToFinalise => {
info!("Prepared migrating CAs to Krill version {}.", KRILL_VERSION);
for scope in self.current_kv_store.scopes()? {
self.remove_data_upgrade_info(&scope)?;
}
}
}
Ok(())
}
fn deployed_store(&self) -> &KeyValueStore {
&self.current_kv_store
}
fn preparation_store(&self) -> &KeyValueStore {
&self.new_kv_store
}
}