use std::{fmt, str::FromStr};
use rpki::{
ca::{idexchange::MyHandle},
repository::{
x509::Time,
},
};
use serde::{Deserialize, Serialize};
use crate::{
commons::{
crypto::dispatch::signerinfo::{
SignerInfo, SignerInfoEvent, SignerInfoInitEvent,
},
eventsourcing::{
Aggregate, AggregateStore, Storable,
StoredCommand, StoredCommandBuilder, WithStorableDetails,
},
storage::{Ident, KeyValueStore},
},
config::Config,
server::{
properties::Properties,
},
server::taproxy::{
TrustAnchorProxy, TrustAnchorProxyEvent, TrustAnchorProxyInitEvent,
},
tasigner::{
TrustAnchorSigner, TrustAnchorSignerEvent,
TrustAnchorSignerInitEvent,
},
};
use super::{
AspaMigrationConfigs, CommandMigrationEffect, UnconvertedEffect,
UpgradeAggregateStorePre0_14, UpgradeMode, UpgradeResult,
};
pub type OldInitSignerInfoEvent = OldStoredEvent<SignerInfoInitEvent>;
pub type OldSignerInfoEvent = OldStoredEvent<SignerInfoEvent>;
pub type OldTrustAnchorProxyInitEvent =
OldStoredEvent<TrustAnchorProxyInitEvent>;
pub type OldTrustAnchorProxyEvent = OldStoredEvent<TrustAnchorProxyEvent>;
pub type OldTrustAnchorSignerInitEvent =
OldStoredEvent<TrustAnchorSignerInitEvent>;
pub type OldTrustAnchorSignerEvent = OldStoredEvent<TrustAnchorSignerEvent>;
pub trait OldEvent:
fmt::Display + Eq + PartialEq + Storable + 'static
{
fn handle(&self) -> &MyHandle;
fn version(&self) -> u64;
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct OldStoredEvent<
E: fmt::Display + Eq + PartialEq + Storable + 'static,
> {
id: MyHandle,
version: u64,
#[serde(deserialize_with = "E::deserialize")]
details: E,
}
impl<E: fmt::Display + Eq + PartialEq + Storable + 'static>
OldStoredEvent<E>
{
pub fn new(id: &MyHandle, version: u64, event: E) -> Self {
OldStoredEvent {
id: id.clone(),
version,
details: event,
}
}
pub fn details(&self) -> &E {
&self.details
}
pub fn into_details(self) -> E {
self.details
}
pub fn unpack(self) -> (MyHandle, u64, E) {
(self.id, self.version, self.details)
}
}
impl<E: fmt::Display + Eq + PartialEq + Storable + 'static> OldEvent
for OldStoredEvent<E>
{
fn handle(&self) -> &MyHandle {
&self.id
}
fn version(&self) -> u64 {
self.version
}
}
impl<E: fmt::Display + Eq + PartialEq + Storable + 'static> fmt::Display
for OldStoredEvent<E>
{
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(
f,
"id: {} version: {} details: {}",
self.id, self.version, self.details
)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct OldStoredCommand<S: WithStorableDetails> {
actor: String,
time: Time,
handle: MyHandle,
version: u64,
sequence: u64,
#[serde(deserialize_with = "S::deserialize")]
details: S,
effect: OldStoredEffect,
}
impl<S: WithStorableDetails> OldStoredCommand<S> {
pub fn actor(&self) -> &String {
&self.actor
}
pub fn time(&self) -> Time {
self.time
}
pub fn handle(&self) -> &MyHandle {
&self.handle
}
pub fn sequence(&self) -> u64 {
self.sequence
}
pub fn effect(&self) -> &OldStoredEffect {
&self.effect
}
pub fn details(&self) -> &S {
&self.details
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case", tag = "result")]
pub enum OldStoredEffect {
Error { msg: String },
Success { events: Vec<u64> },
}
impl OldStoredEffect {
pub fn events(&self) -> Option<&Vec<u64>> {
match self {
OldStoredEffect::Error { .. } => None,
OldStoredEffect::Success { events } => Some(events),
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OldCommandKey {
pub sequence: u64,
pub timestamp_secs: i64,
pub label: Box<Ident>,
}
impl OldCommandKey {
pub fn to_storage_key(&self) -> Box<Ident> {
Ident::builder(
const { Ident::make("command--") }
).push_i64(
self.timestamp_secs
).push_ident(
const { Ident::make("--") }
).push_u64(
self.sequence
).push_ident(
const { Ident::make("--") }
).push_ident(
&self.label,
).finish_with_extension(
const { Ident::make("json") }
)
}
}
impl fmt::Display for OldCommandKey {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(
f,
"command--{}--{}--{}",
self.timestamp_secs, self.sequence, self.label
)
}
}
impl FromStr for OldCommandKey {
type Err = CommandKeyError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
let parts: Vec<&str> = s.split("--").collect();
if parts.len() != 4 || parts[0] != "command" {
Err(CommandKeyError(s.to_string()))
} else {
let timestamp_secs = i64::from_str(parts[1])
.map_err(|_| CommandKeyError(s.to_string()))?;
let sequence = u64::from_str(parts[2])
.map_err(|_| CommandKeyError(s.to_string()))?;
let label = {
let end = parts[3].to_string();
let last = if end.ends_with(".json") {
end.len() - 5
} else {
end.len()
};
Ident::boxed_from_string(
end[0..last].to_string()
).map_err(|_| CommandKeyError(s.to_string()))?
};
Ok(OldCommandKey {
sequence,
timestamp_secs,
label,
})
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct CommandKeyError(String);
impl fmt::Display for CommandKeyError {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(f, "invalid command key: {}", self.0)
}
}
pub type UpgradeAggregateStoreProperties =
GenericUpgradeAggregateStore<Properties>;
pub type UpgradeAggregateStoreSignerInfo =
GenericUpgradeAggregateStore<SignerInfo>;
pub type UpgradeAggregateStoreTrustAnchorSigner =
GenericUpgradeAggregateStore<TrustAnchorSigner>;
pub type UpgradeAggregateStoreTrustAnchorProxy =
GenericUpgradeAggregateStore<TrustAnchorProxy>;
pub struct GenericUpgradeAggregateStore<A: Aggregate> {
store_name: String,
current_kv_store: KeyValueStore,
new_kv_store: KeyValueStore,
new_agg_store: AggregateStore<A>,
}
impl<A: Aggregate> GenericUpgradeAggregateStore<A> {
pub fn upgrade(
name_space: &Ident,
mode: UpgradeMode,
config: &Config,
) -> UpgradeResult<AspaMigrationConfigs> {
let current_kv_store =
KeyValueStore::create(&config.storage_uri, name_space)?;
if current_kv_store.scopes()?.is_empty() {
Ok(AspaMigrationConfigs::default())
} else {
let new_kv_store = KeyValueStore::create_upgrade_store(
&config.storage_uri,
name_space,
)?;
let new_agg_store = AggregateStore::<A>::create_upgrade_store(
&config.storage_uri,
name_space,
config.use_history_cache,
)?;
let store_migration = GenericUpgradeAggregateStore {
store_name: name_space.to_string(),
current_kv_store,
new_kv_store,
new_agg_store,
};
store_migration.upgrade(mode)
}
}
}
impl<A: Aggregate> UpgradeAggregateStorePre0_14
for GenericUpgradeAggregateStore<A>
{
type Aggregate = A;
type OldInitEvent = A::InitEvent;
type OldEvent = A::Event;
type OldStorableDetails = A::StorableCommandDetails;
fn store_name(&self) -> &str {
&self.store_name
}
fn deployed_store(&self) -> &KeyValueStore {
&self.current_kv_store
}
fn preparation_key_value_store(&self) -> &KeyValueStore {
&self.new_kv_store
}
fn preparation_aggregate_store(
&self,
) -> &AggregateStore<Self::Aggregate> {
&self.new_agg_store
}
fn convert_init_event(
&self,
old_init: Self::OldInitEvent,
handle: MyHandle,
actor: String,
time: Time,
) -> UpgradeResult<StoredCommand<Self::Aggregate>> {
let details = A::StorableCommandDetails::make_init();
let builder =
StoredCommandBuilder::<A>::new(actor, time, handle, 0, details);
Ok(builder.finish_with_init_event(old_init))
}
fn convert_old_command(
&self,
old_command: OldStoredCommand<Self::OldStorableDetails>,
old_effect: UnconvertedEffect<Self::OldEvent>,
version: u64,
) -> UpgradeResult<CommandMigrationEffect<Self::Aggregate>> {
let new_command_builder = StoredCommandBuilder::<A>::new(
old_command.actor().clone(),
old_command.time(),
old_command.handle().clone(),
version,
old_command.details().clone(),
);
let new_command = match old_effect {
UnconvertedEffect::Error { msg } => {
new_command_builder.finish_with_error(msg)
}
UnconvertedEffect::Success { events } => {
new_command_builder.finish_with_events(events)
}
};
Ok(CommandMigrationEffect::StoredCommand(new_command))
}
}