use std::{borrow::Cow, collections::HashMap};
use crate::{Event, Upcaster};
pub(crate) struct Alias {
pub target_key: String,
pub aggregate_type: &'static str,
pub from: &'static str,
to: &'static str,
hops: u8,
upcast: fn(&[u8]) -> Result<Vec<u8>, bitcode::Error>,
convert: bool,
}
impl Alias {
pub fn apply<'a>(&self, event: &'a Event) -> anyhow::Result<Cow<'a, Event>> {
if !self.convert {
return Ok(Cow::Borrowed(event));
}
let data = (self.upcast)(&event.data).map_err(|e| {
anyhow::anyhow!(
"failed to upcast `{}` to `{}` (event {}): {e}",
self.from,
self.to,
event.id
)
})?;
tracing::debug!(from = self.from, to = self.to, "upcast event");
let mut event = event.clone();
event.name = self.to.to_owned();
event.data = data;
Ok(Cow::Owned(event))
}
}
#[derive(Default)]
pub(crate) struct Aliases(HashMap<String, Alias>);
impl Aliases {
pub fn register(
&mut self,
aggregate_type: &'static str,
to: &'static str,
upcasters: &'static [Upcaster],
convert: bool,
) {
for upcaster in upcasters {
let key = format!("{aggregate_type}_{}", upcaster.from);
if self
.0
.get(&key)
.is_some_and(|alias| alias.hops <= upcaster.hops)
{
continue;
}
self.0.insert(
key,
Alias {
target_key: format!("{aggregate_type}_{to}"),
aggregate_type,
from: upcaster.from,
to,
hops: upcaster.hops,
upcast: upcaster.upcast,
convert,
},
);
}
}
pub fn get(&self, key: &str) -> Option<&Alias> {
self.0.get(key)
}
pub fn iter(&self) -> impl Iterator<Item = (&String, &Alias)> {
self.0.iter()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn append(suffix: u8) -> fn(&[u8]) -> Result<Vec<u8>, bitcode::Error> {
match suffix {
2 => |d| Ok([d, &[2]].concat()),
_ => |d| Ok([d, &[3]].concat()),
}
}
fn v2() -> &'static [Upcaster] {
Box::leak(Box::new([Upcaster::new("V1", 1, append(2))]))
}
fn v3() -> &'static [Upcaster] {
Box::leak(Box::new([
Upcaster::new("V1", 2, append(3)),
Upcaster::new("V2", 1, append(3)),
]))
}
#[test]
fn nearest_target_wins_whatever_the_order() {
for order in [[2u8, 3], [3, 2]] {
let mut aliases = Aliases::default();
for v in order {
match v {
2 => aliases.register("t/A", "V2", v2(), true),
_ => aliases.register("t/A", "V3", v3(), true),
}
}
assert_eq!(aliases.get("t/A_V1").unwrap().target_key, "t/A_V2");
assert_eq!(aliases.get("t/A_V2").unwrap().target_key, "t/A_V3");
}
}
#[test]
fn apply_rewrites_name_and_data_only() {
let mut aliases = Aliases::default();
aliases.register("t/A", "V2", v2(), true);
let stored = Event {
id: ulid::Ulid::generate(),
aggregate_id: "a-1".to_owned(),
aggregate_type: "t/A".to_owned(),
version: 7,
name: "V1".to_owned(),
routing_key: Some("rk".to_owned()),
data: vec![1],
metadata: Default::default(),
timestamp: 42,
timestamp_subsec: 5,
};
let upcast = aliases.get("t/A_V1").unwrap().apply(&stored).unwrap();
assert_eq!(upcast.name, "V2");
assert_eq!(upcast.data, vec![1, 2]);
assert_eq!(
(
upcast.id,
upcast.version,
upcast.timestamp,
&upcast.routing_key
),
(stored.id, 7, 42, &stored.routing_key)
);
}
#[test]
fn skip_alias_leaves_the_event_untouched() {
let mut aliases = Aliases::default();
aliases.register("t/A", "V2", v2(), false);
let stored = Event {
id: ulid::Ulid::generate(),
aggregate_id: "a-1".to_owned(),
aggregate_type: "t/A".to_owned(),
version: 1,
name: "V1".to_owned(),
routing_key: None,
data: vec![1],
metadata: Default::default(),
timestamp: 0,
timestamp_subsec: 0,
};
let upcast = aliases.get("t/A_V1").unwrap().apply(&stored).unwrap();
assert!(matches!(upcast, Cow::Borrowed(_)));
}
}