use std::ops::Deref;
use std::sync::{Arc, Mutex, PoisonError};
use async_trait::async_trait;
use serde_json::{Map, Value};
use super::event::{FIELD_KIND, FIELD_SEQ};
#[cfg(test)]
use super::event::{kind_of, seq_of, validate_event, FIELD_EPOCH_MS};
#[cfg(test)]
use super::History;
use super::{KnlError, KnlResult};
pub const SCHEMA_VERSION_FIELD: &str = "_schema_version";
pub const CURRENT_SCHEMA_VERSION: u64 = 1;
pub fn kernel_upcasters() -> Vec<Arc<dyn Upcaster>> {
Vec::new()
}
pub(super) fn stamp_schema_version(event: &mut Map<String, Value>) {
event.insert(
SCHEMA_VERSION_FIELD.to_string(),
Value::from(CURRENT_SCHEMA_VERSION),
);
}
pub trait Upcaster: Send + Sync {
fn upcast(&self, event: Value) -> Value;
}
pub fn apply_upcasters(chain: &[Arc<dyn Upcaster>], events: Vec<Value>) -> Vec<Value> {
events
.into_iter()
.map(|event| chain.iter().fold(event, |event, up| up.upcast(event)))
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Current(Map<String, Value>);
impl Current {
fn from_upcasted(event: Value) -> KnlResult<Self> {
let Value::Object(map) = event else {
return Err(KnlError::Corruption(format!(
"stored event is not an object: {event}"
)));
};
debug_assert_eq!(
map.get(SCHEMA_VERSION_FIELD).and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION),
"the upcaster chain must bring every event to the current schema \
version; a step is missing from kernel_upcasters(): {map:?}"
);
Ok(Self(map))
}
#[cfg(test)]
pub fn assume_current(event: Value) -> Self {
match event {
Value::Object(map) => Self(map),
other => panic!("a test fixture event must be an object, got {other}"),
}
}
pub fn kind(&self) -> &str {
self.0.get(FIELD_KIND).and_then(Value::as_str).unwrap_or("")
}
pub fn seq(&self) -> u64 {
self.0.get(FIELD_SEQ).and_then(Value::as_u64).unwrap_or(0)
}
pub fn into_inner(self) -> Map<String, Value> {
self.0
}
}
impl Deref for Current {
type Target = Map<String, Value>;
fn deref(&self) -> &Self::Target {
&self.0
}
}
impl std::fmt::Display for Current {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match serde_json::to_string(&self.0) {
Ok(json) => f.write_str(&json),
Err(_) => write!(f, "{:?}", self.0),
}
}
}
pub type Decision = Box<dyn FnOnce(Vec<Value>) -> Option<Map<String, Value>> + Send + 'static>;
pub type CurrentDecision = Box<dyn FnOnce(Vec<Current>) -> Option<Map<String, Value>> + Send>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Split<T> {
pub own: Vec<T>,
pub other: Vec<T>,
}
impl<T> Split<T> {
pub fn own(own: Vec<T>) -> Self {
Self {
own,
other: Vec::new(),
}
}
}
pub type SplitDecision =
Box<dyn FnOnce(Split<Value>) -> Option<Split<Map<String, Value>>> + Send + 'static>;
pub type CurrentSplitDecision =
Box<dyn FnOnce(Split<Current>) -> Option<Split<Map<String, Value>>> + Send>;
pub type ChildrenDecision = Box<dyn FnOnce(Vec<String>) -> Map<String, Value> + Send + 'static>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ChildScan {
pub opened: String,
pub closed: String,
pub parent_field: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Committed {
pub seq: u64,
pub epoch_ms: u64,
}
#[async_trait]
pub trait EventStore: Send + Sync {
async fn append(&mut self, event: Map<String, Value>) -> KnlResult<Committed>;
async fn append_many(&mut self, events: Vec<Map<String, Value>>) -> KnlResult<Vec<Committed>> {
let mut committed = Vec::with_capacity(events.len());
for event in events {
committed.push(self.append(event).await?);
}
Ok(committed)
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: Decision,
) -> KnlResult<Option<Committed>>;
async fn append_if_many(
&mut self,
other: &str,
kinds: Option<&[&str]>,
decide: SplitDecision,
) -> KnlResult<Option<Split<Committed>>> {
let _ = (other, kinds, decide);
Err(KnlError::Unsupported(
"this store keeps one stream, so it cannot write two in one transaction".to_string(),
))
}
async fn append_with_open_children(
&mut self,
scan: &ChildScan,
decide: ChildrenDecision,
) -> KnlResult<Committed> {
let _ = scan;
self.append(decide(Vec::new())).await
}
fn database(&self) -> Option<&str> {
None
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>>;
async fn read(&self, from_seq: u64, limit: usize) -> KnlResult<Vec<Value>> {
self.read_kinds(None, from_seq, limit).await
}
async fn read_last(&self, n: usize) -> KnlResult<Vec<Value>> {
let mut events = self.read(0, usize::MAX).await?;
let start = events.len().saturating_sub(n);
Ok(events.split_off(start))
}
async fn head(&self) -> KnlResult<Option<u64>>;
async fn len(&self) -> KnlResult<usize>;
async fn is_empty(&self) -> KnlResult<bool> {
Ok(self.len().await? == 0)
}
async fn query(&self, plan: &super::query::QueryPlan) -> KnlResult<super::query::QueryRows> {
let _ = plan;
Err(KnlError::Unsupported(
"this store keeps no queryable table".to_string(),
))
}
fn detach_append(&self, event: Map<String, Value>) {
let kind = event
.get(FIELD_KIND)
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
tracing::warn!(
%kind,
"knl: this store cannot record a detached append; the event was not written"
);
}
}
#[cfg(test)]
#[derive(Debug, Clone)]
pub struct MemEventStore {
history: History,
}
#[cfg(test)]
impl MemEventStore {
pub fn new() -> Self {
Self {
history: History::new(),
}
}
pub fn history(&self) -> &History {
&self.history
}
}
#[cfg(test)]
impl Default for MemEventStore {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
#[async_trait]
impl EventStore for MemEventStore {
async fn append(&mut self, mut event: Map<String, Value>) -> KnlResult<Committed> {
stamp_schema_version(&mut event);
let seq = self.history.append(event)?;
let epoch_ms = self
.history
.events()
.last()
.and_then(|event| event.get(FIELD_EPOCH_MS))
.and_then(Value::as_u64)
.unwrap_or(0);
Ok(Committed { seq, epoch_ms })
}
async fn append_many(&mut self, events: Vec<Map<String, Value>>) -> KnlResult<Vec<Committed>> {
for event in &events {
validate_event(event)?;
}
let mut committed = Vec::with_capacity(events.len());
for event in events {
committed.push(self.append(event).await?);
}
Ok(committed)
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: Decision,
) -> KnlResult<Option<Committed>> {
let events = self.read_kinds(kinds, 0, usize::MAX).await?;
match decide(events) {
Some(event) => self.append(event).await.map(Some),
None => Ok(None),
}
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
let mut events = self.history.since(from_seq);
if let Some(kinds) = kinds {
events.retain(|event| kinds.contains(&kind_of(event)));
}
events.truncate(limit);
Ok(events)
}
async fn read_last(&self, n: usize) -> KnlResult<Vec<Value>> {
let mut events = self.history.since(0);
let start = events.len().saturating_sub(n);
Ok(events.split_off(start))
}
async fn head(&self) -> KnlResult<Option<u64>> {
Ok(self.history.events().last().map(seq_of))
}
async fn len(&self) -> KnlResult<usize> {
Ok(self.history.len())
}
}
pub struct CurrentStore {
inner: Box<dyn EventStore>,
chain: Vec<Arc<dyn Upcaster>>,
}
impl CurrentStore {
pub fn new(inner: Box<dyn EventStore>, chain: Vec<Arc<dyn Upcaster>>) -> Self {
Self { inner, chain }
}
pub fn detach_append(&self, event: Map<String, Value>) {
self.inner.detach_append(event);
}
fn project(chain: &[Arc<dyn Upcaster>], events: Vec<Value>) -> KnlResult<Vec<Current>> {
apply_upcasters(chain, events)
.into_iter()
.map(Current::from_upcasted)
.collect()
}
pub async fn append(&mut self, event: Map<String, Value>) -> KnlResult<Committed> {
self.inner.append(event).await
}
pub async fn append_many(
&mut self,
events: Vec<Map<String, Value>>,
) -> KnlResult<Vec<Committed>> {
self.inner.append_many(events).await
}
pub async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: CurrentDecision,
) -> KnlResult<Option<Committed>> {
let chain = self.chain.clone();
let failure: Arc<Mutex<Option<KnlError>>> = Arc::new(Mutex::new(None));
let parked = Arc::clone(&failure);
let upcasted: Decision =
Box::new(
move |events: Vec<Value>| match Self::project(&chain, events) {
Ok(current) => decide(current),
Err(fault) => {
*parked.lock().unwrap_or_else(PoisonError::into_inner) = Some(fault);
None
}
},
);
let committed = self.inner.append_if(kinds, upcasted).await;
let parked = failure
.lock()
.unwrap_or_else(PoisonError::into_inner)
.take();
match parked {
Some(fault) => Err(fault),
None => committed,
}
}
pub async fn append_if_many(
&mut self,
other: &str,
kinds: Option<&[&str]>,
decide: CurrentSplitDecision,
) -> KnlResult<Option<Split<Committed>>> {
let chain = self.chain.clone();
let failure: Arc<Mutex<Option<KnlError>>> = Arc::new(Mutex::new(None));
let parked = Arc::clone(&failure);
let upcasted: SplitDecision = Box::new(move |events: Split<Value>| {
let projected = Self::project(&chain, events.own).and_then(|own| {
let other = Self::project(&chain, events.other)?;
Ok(Split { own, other })
});
match projected {
Ok(current) => decide(current),
Err(fault) => {
*parked.lock().unwrap_or_else(PoisonError::into_inner) = Some(fault);
None
}
}
});
let committed = self.inner.append_if_many(other, kinds, upcasted).await;
let parked = failure
.lock()
.unwrap_or_else(PoisonError::into_inner)
.take();
match parked {
Some(fault) => Err(fault),
None => committed,
}
}
pub async fn append_with_open_children(
&mut self,
scan: &ChildScan,
decide: ChildrenDecision,
) -> KnlResult<Committed> {
self.inner.append_with_open_children(scan, decide).await
}
pub fn database(&self) -> Option<&str> {
self.inner.database()
}
pub async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Current>> {
let events = self.inner.read_kinds(kinds, from_seq, limit).await?;
Self::project(&self.chain, events)
}
pub async fn read(&self, from_seq: u64, limit: usize) -> KnlResult<Vec<Current>> {
self.read_kinds(None, from_seq, limit).await
}
pub async fn read_last(&self, n: usize) -> KnlResult<Vec<Current>> {
let events = self.inner.read_last(n).await?;
Self::project(&self.chain, events)
}
pub async fn head(&self) -> KnlResult<Option<u64>> {
self.inner.head().await
}
pub async fn len(&self) -> KnlResult<usize> {
self.inner.len().await
}
pub async fn is_empty(&self) -> KnlResult<bool> {
self.inner.is_empty().await
}
pub async fn query(
&self,
plan: &super::query::QueryPlan,
) -> KnlResult<super::query::QueryRows> {
self.inner.query(plan).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::knl::event::{KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT};
use serde_json::json;
fn obj(value: Value) -> Map<String, Value> {
match value {
Value::Object(map) => map,
other => panic!("test fixture must be an object, got {other}"),
}
}
fn ev(i: usize) -> Map<String, Value> {
obj(json!({ "kind": format!("e{i}") }))
}
fn decide(
f: impl FnOnce(Vec<Value>) -> Option<Map<String, Value>> + Send + 'static,
) -> Decision {
Box::new(f)
}
fn decide_current(
f: impl FnOnce(Vec<Current>) -> Option<Map<String, Value>> + Send + 'static,
) -> CurrentDecision {
Box::new(f)
}
#[tokio::test]
async fn append_assigns_gap_free_monotonic_seq_from_one() {
let mut store = MemEventStore::new();
assert!(store.is_empty().await.expect("is_empty"));
assert_eq!(store.len().await.expect("len"), 0);
let a = store.append(ev(1)).await.expect("append e1");
let b = store.append(ev(2)).await.expect("append e2");
let c = store.append(ev(3)).await.expect("append e3");
assert_eq!((a.seq, b.seq, c.seq), (1, 2, 3));
assert!(a.epoch_ms >= 1 || a.epoch_ms == 0, "epoch is stamped");
assert_eq!(store.len().await.expect("len"), 3);
assert!(!store.is_empty().await.expect("is_empty"));
let stored = store.read(0, usize::MAX).await.expect("read");
let stored_epoch = stored[0]
.get(FIELD_EPOCH_MS)
.and_then(Value::as_u64)
.expect("epoch is on the stored event");
assert_eq!(stored_epoch, a.epoch_ms);
}
#[tokio::test]
async fn a_rejected_append_records_nothing_and_burns_no_seq() {
let mut store = MemEventStore::new();
store
.append(obj(json!({ "text": "no kind" })))
.await
.expect_err("kind is required");
assert_eq!(store.len().await.expect("len"), 0);
assert_eq!(store.append(ev(1)).await.expect("append").seq, 1);
}
#[tokio::test]
async fn append_if_decides_on_the_stream_and_writes_only_a_some() {
let mut store = MemEventStore::new();
store.append(ev(1)).await.expect("seed");
let seen = Arc::new(Mutex::new(0_usize));
let counted = Arc::clone(&seen);
let committed = store
.append_if(
None,
decide(move |events| {
*counted.lock().expect("not poisoned") = events.len();
Some(ev(2))
}),
)
.await
.expect("append_if");
assert_eq!(
*seen.lock().expect("not poisoned"),
1,
"decide was handed the whole stream"
);
assert_eq!(committed.map(|c| c.seq), Some(2));
let nothing = store
.append_if(None, decide(|_| None))
.await
.expect("append_if");
assert_eq!(nothing, None);
assert_eq!(store.len().await.expect("len"), 2, "a None writes nothing");
assert_eq!(store.append(ev(3)).await.expect("append").seq, 3);
}
#[tokio::test]
async fn append_if_validates_the_event_the_decision_returns() {
let mut store = MemEventStore::new();
store
.append_if(None, decide(|_| Some(obj(json!({ "text": "no kind" })))))
.await
.expect_err("kind is required");
assert_eq!(store.len().await.expect("len"), 0);
}
#[tokio::test]
async fn append_if_shows_the_decision_only_the_kinds_it_asked_for() {
let mut store = MemEventStore::new();
store
.append(obj(
json!({ "kind": KIND_BUDGET_GRANTED, "data": { "amount": 100 } }),
))
.await
.expect("the grant");
store.append(ev(1)).await.expect("noise");
store.append(ev(2)).await.expect("more noise");
let seen: Arc<Mutex<Vec<String>>> = Arc::default();
let recorded = Arc::clone(&seen);
let committed = store
.append_if(
Some(&[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]),
decide(move |events| {
*recorded.lock().expect("not poisoned") =
events.iter().map(|e| kind_of(e).to_string()).collect();
Some(obj(
json!({ "kind": KIND_BUDGET_SPENT, "data": { "amount": 10 } }),
))
}),
)
.await
.expect("append_if");
assert_eq!(
*seen.lock().expect("not poisoned"),
[KIND_BUDGET_GRANTED],
"only the kinds asked for"
);
assert_eq!(
committed.map(|c| c.seq),
Some(4),
"the write is not filtered"
);
let recorded = Arc::clone(&seen);
store
.append_if(
Some(&[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]),
decide(move |events| {
*recorded.lock().expect("not poisoned") =
events.iter().map(|e| kind_of(e).to_string()).collect();
None
}),
)
.await
.expect("append_if");
assert_eq!(
*seen.lock().expect("not poisoned"),
[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]
);
}
#[tokio::test]
async fn append_many_records_the_batch_in_order_or_records_nothing() {
let mut store = MemEventStore::new();
let committed = store
.append_many(vec![ev(1), ev(2)])
.await
.expect("the batch");
assert_eq!(
committed.iter().map(|c| c.seq).collect::<Vec<_>>(),
[1, 2],
"the batch lands in the order it was given"
);
let stored = store.read(0, usize::MAX).await.expect("read");
let kinds: Vec<&str> = stored.iter().map(kind_of).collect();
assert_eq!(kinds, ["e1", "e2"]);
store
.append_many(vec![ev(3), obj(json!({ "text": "no kind" }))])
.await
.expect_err("kind is required");
assert_eq!(
store.len().await.expect("len"),
2,
"a failed batch wrote nothing"
);
assert_eq!(
store.append(ev(4)).await.expect("append").seq,
3,
"no seq burnt"
);
}
#[tokio::test]
async fn read_kinds_selects_by_kind_and_still_pages() {
let mut store = MemEventStore::new();
store
.append(obj(
json!({ "kind": KIND_BUDGET_GRANTED, "data": { "amount": 100 } }),
))
.await
.expect("grant");
store.append(ev(1)).await.expect("noise");
store
.append(obj(
json!({ "kind": KIND_BUDGET_SPENT, "data": { "amount": 10 } }),
))
.await
.expect("spend");
store.append(ev(2)).await.expect("more noise");
let ledger = store
.read_kinds(
Some(&[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]),
0,
usize::MAX,
)
.await
.expect("read_kinds");
let kinds: Vec<&str> = ledger.iter().map(kind_of).collect();
assert_eq!(kinds, [KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]);
assert_eq!(
seq_of(&ledger[1]),
3,
"the seq is the stream's, not the fold's"
);
assert_eq!(
store
.read_kinds(Some(&[KIND_BUDGET_GRANTED]), 2, usize::MAX)
.await
.expect("read_kinds")
.len(),
0,
"the grant is before from_seq"
);
assert_eq!(
store
.read_kinds(Some(&[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]), 0, 1)
.await
.expect("read_kinds")
.len(),
1
);
assert!(store
.read_kinds(Some(&[]), 0, usize::MAX)
.await
.expect("read_kinds")
.is_empty());
assert_eq!(
store
.read_kinds(None, 0, usize::MAX)
.await
.expect("read_kinds")
.len(),
4,
"read() is read_kinds(None, ..)"
);
assert_eq!(store.read(0, usize::MAX).await.expect("read").len(), 4);
}
#[tokio::test]
async fn read_pages_by_from_seq_and_limit() {
let mut store = MemEventStore::new();
for i in 1..=5 {
store.append(ev(i)).await.expect("append");
}
assert_eq!(store.read(0, usize::MAX).await.expect("read").len(), 5);
assert_eq!(store.read(1, usize::MAX).await.expect("read").len(), 5);
assert_eq!(store.read(3, usize::MAX).await.expect("read").len(), 3);
assert_eq!(store.read(6, usize::MAX).await.expect("read").len(), 0);
let page = store.read(2, 2).await.expect("read");
assert_eq!(page.len(), 2);
assert_eq!(kind_of(&page[0]), "e2");
assert_eq!(kind_of(&page[1]), "e3");
assert!(store.read(0, 0).await.expect("read").is_empty());
}
#[tokio::test]
async fn head_is_none_when_empty_then_tracks_the_max() {
let mut store = MemEventStore::new();
assert_eq!(store.head().await.expect("head"), None);
store.append(ev(1)).await.expect("append");
assert_eq!(store.head().await.expect("head"), Some(1));
store.append(ev(2)).await.expect("append");
assert_eq!(store.head().await.expect("head"), Some(2));
store
.append(obj(json!({ "text": "no kind" })))
.await
.expect_err("kind is required");
assert_eq!(store.head().await.expect("head"), Some(2));
}
#[tokio::test]
async fn default_matches_new_and_starts_seq_at_one() {
let mut store = MemEventStore::default();
assert!(store.is_empty().await.expect("is_empty"));
assert_eq!(store.append(ev(1)).await.expect("append").seq, 1);
}
#[tokio::test]
async fn reads_are_copies_so_the_store_cannot_be_reached_through_them() {
let mut store = MemEventStore::new();
store.append(ev(1)).await.expect("append");
let mut copy = store.read(0, usize::MAX).await.expect("read");
copy[0][FIELD_KIND] = Value::String("TAMPERED".into());
let again = store.read(0, usize::MAX).await.expect("read");
assert_eq!(kind_of(&again[0]), "e1");
}
#[tokio::test]
async fn append_stamps_the_current_schema_version() {
let mut store = MemEventStore::new();
store.append(ev(1)).await.expect("append");
let stored = store.read(0, usize::MAX).await.expect("read");
assert_eq!(
stored[0].get(SCHEMA_VERSION_FIELD).and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION),
"a stored event carries the version it was written under: {}",
stored[0]
);
}
#[tokio::test]
async fn a_single_stream_store_cannot_write_two_streams() {
let mut store = MemEventStore::new();
assert_eq!(store.database(), None, "one stream, no database to share");
let err = store
.append_if_many(
"another-stream",
None,
Box::new(|_| {
panic!("the decision must not run: there is nowhere for its other half to go")
}),
)
.await
.expect_err("two streams are a durable backend's");
assert_eq!(err.kind(), KnlError::UNSUPPORTED, "{err}");
assert!(!err.is_retryable(), "asking again changes nothing: {err}");
assert_eq!(
store.len().await.expect("len"),
0,
"and nothing was written"
);
}
#[tokio::test]
async fn a_single_stream_store_finds_no_children_and_appends_anyway() {
let mut store = MemEventStore::new();
let scan = ChildScan {
opened: "session_opened".to_string(),
closed: "session_closed".to_string(),
parent_field: "parent".to_string(),
};
let seen: Arc<Mutex<Option<usize>>> = Arc::default();
let counted = Arc::clone(&seen);
let committed = store
.append_with_open_children(
&scan,
Box::new(move |children| {
*counted.lock().expect("not poisoned") = Some(children.len());
ev(1)
}),
)
.await
.expect("the close lands");
assert_eq!(*seen.lock().expect("not poisoned"), Some(0));
assert_eq!(committed.seq, 1);
assert_eq!(store.len().await.expect("len"), 1);
}
#[tokio::test]
async fn a_store_with_no_table_refuses_a_query() {
use crate::knl::query::{plan, QueryOpts, QueryParams};
let store = MemEventStore::new();
let asked = plan(
"SELECT 1",
QueryParams::None,
&QueryOpts::default(),
"a-stream",
)
.expect("a plan");
let err = store
.query(&asked)
.await
.expect_err("there is no table to query");
assert_eq!(err.kind(), KnlError::UNSUPPORTED, "{err}");
assert!(!err.is_retryable(), "asking again changes nothing: {err}");
}
#[test]
fn an_empty_upcaster_chain_is_the_identity() {
let events = vec![
json!({ "kind": "a", "seq": 1 }),
json!({ "kind": "b", "seq": 2 }),
];
assert_eq!(apply_upcasters(&[], events.clone()), events);
}
#[test]
fn upcasters_compose_per_event_in_registration_order() {
use std::sync::Arc;
struct Tag(&'static str);
impl Upcaster for Tag {
fn upcast(&self, mut event: Value) -> Value {
let map = event.as_object_mut().expect("event is an object");
let trace = map
.entry("trace")
.or_insert_with(|| Value::Array(Vec::new()));
trace
.as_array_mut()
.expect("trace is an array")
.push(Value::from(self.0));
event
}
}
let chain: Vec<Arc<dyn Upcaster>> = vec![Arc::new(Tag("first")), Arc::new(Tag("second"))];
let out = apply_upcasters(&chain, vec![json!({ "kind": "x" }), json!({ "kind": "y" })]);
assert_eq!(out[0]["trace"], json!(["first", "second"]));
assert_eq!(out[1]["trace"], json!(["first", "second"]));
}
#[tokio::test]
async fn the_seam_projects_on_read_and_leaves_writes_untouched() {
struct Mark;
impl Upcaster for Mark {
fn upcast(&self, mut event: Value) -> Value {
let map = event.as_object_mut().expect("event is an object");
map.entry("trace")
.or_insert_with(|| Value::Array(Vec::new()))
.as_array_mut()
.expect("trace is an array")
.push(Value::from("mark"));
event
}
}
let chain: Vec<Arc<dyn Upcaster>> = vec![Arc::new(Mark)];
let mut store = CurrentStore::new(Box::new(MemEventStore::new()), chain);
let a = store.append(ev(1)).await.expect("append e1");
assert_eq!(a.seq, 1);
assert_eq!(store.head().await.expect("head"), Some(1));
assert_eq!(store.len().await.expect("len"), 1);
assert!(!store.is_empty().await.expect("is_empty"));
let first = store.read(0, usize::MAX).await.expect("read");
assert_eq!(first[0]["trace"], json!(["mark"]), "read applies the chain");
let second = store.read(0, usize::MAX).await.expect("read again");
assert_eq!(
second[0]["trace"],
json!(["mark"]),
"stored bytes carry no marker; read adds exactly one"
);
let b = store.append(ev(2)).await.expect("append e2");
assert_eq!(b.seq, 2);
assert_eq!(store.head().await.expect("head"), Some(2));
assert_eq!(store.len().await.expect("len"), 2);
let both = store.read(0, usize::MAX).await.expect("read both");
assert_eq!(both.len(), 2);
assert_eq!(both[1].kind(), "e2");
assert_eq!(both[1].seq(), 2, "a Current keeps its coordinates");
assert_eq!(both[1]["trace"], json!(["mark"]));
}
#[tokio::test]
async fn the_seam_with_an_empty_chain_returns_events_unchanged() {
let mut store = CurrentStore::new(Box::new(MemEventStore::new()), Vec::new());
assert!(store.is_empty().await.expect("is_empty"));
store.append(ev(1)).await.expect("append");
let read = store.read(0, usize::MAX).await.expect("read");
assert_eq!(read.len(), 1);
assert_eq!(read[0].kind(), "e1", "the event passes through unchanged");
assert!(
read[0].get("trace").is_none(),
"an empty chain adds nothing: {:?}",
read[0]
);
assert_eq!(store.head().await.expect("head"), Some(1));
assert_eq!(store.len().await.expect("len"), 1);
}
#[tokio::test]
async fn the_kind_filter_selects_on_the_stored_kind_and_the_chain_runs_after() {
struct RenameSpent;
impl Upcaster for RenameSpent {
fn upcast(&self, mut event: Value) -> Value {
let Some(map) = event.as_object_mut() else {
return event;
};
if map.get(FIELD_KIND).and_then(Value::as_str) == Some("old_spent") {
map.insert(FIELD_KIND.to_string(), Value::from(KIND_BUDGET_SPENT));
}
event
}
}
let chain: Vec<Arc<dyn Upcaster>> = vec![Arc::new(RenameSpent)];
let mut store = CurrentStore::new(Box::new(MemEventStore::new()), chain);
store
.append(obj(
json!({ "kind": KIND_BUDGET_GRANTED, "data": { "amount": 100 } }),
))
.await
.expect("the grant");
store
.append(obj(
json!({ "kind": "old_spent", "data": { "amount": 10 } }),
))
.await
.expect("a settlement under the older name");
let renamed = store
.read_kinds(Some(&["old_spent"]), 0, usize::MAX)
.await
.expect("read_kinds");
assert_eq!(renamed.len(), 1);
assert_eq!(
renamed[0].kind(),
KIND_BUDGET_SPENT,
"the chain ran after the selection"
);
assert!(
store
.read_kinds(Some(&[KIND_BUDGET_SPENT]), 0, usize::MAX)
.await
.expect("read_kinds")
.is_empty(),
"the filter cannot see a kind the chain has not produced yet"
);
let seen: Arc<Mutex<Vec<String>>> = Arc::default();
let recorded = Arc::clone(&seen);
store
.append_if(
None,
decide_current(move |events| {
*recorded.lock().expect("not poisoned") =
events.iter().map(|e| e.kind().to_string()).collect();
None
}),
)
.await
.expect("append_if");
assert_eq!(
*seen.lock().expect("not poisoned"),
[KIND_BUDGET_GRANTED, KIND_BUDGET_SPENT]
);
}
struct RenameOldKind;
impl Upcaster for RenameOldKind {
fn upcast(&self, mut event: Value) -> Value {
let version = event
.get(SCHEMA_VERSION_FIELD)
.and_then(Value::as_u64)
.unwrap_or(1);
if version >= 2 {
return event;
}
let Some(map) = event.as_object_mut() else {
return event;
};
if map.get(FIELD_KIND).and_then(Value::as_str) == Some("old_kind") {
map.insert(FIELD_KIND.to_string(), Value::from("new_kind"));
}
map.insert(SCHEMA_VERSION_FIELD.to_string(), Value::from(2_u64));
event
}
}
#[test]
fn the_kernel_chain_is_empty_and_the_current_version_is_one() {
assert_eq!(CURRENT_SCHEMA_VERSION, 1);
assert!(
kernel_upcasters().is_empty(),
"no shape has been released, so no step is owed"
);
let stored = json!({ "kind": "note", "seq": 1, SCHEMA_VERSION_FIELD: 1 });
assert_eq!(
apply_upcasters(&kernel_upcasters(), vec![stored.clone()]),
vec![stored],
"an empty chain reads the log back verbatim"
);
}
#[tokio::test]
async fn new_events_are_stamped_with_the_current_version() {
let mut store = MemEventStore::new();
store.append(ev(1)).await.expect("append");
let stored = store.read(0, usize::MAX).await.expect("read");
assert_eq!(
stored[0].get(SCHEMA_VERSION_FIELD).and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION),
"{}",
stored[0]
);
}
#[test]
fn an_upcaster_leaves_a_current_or_unrecognised_event_unchanged() {
let chain: Vec<Arc<dyn Upcaster>> = vec![Arc::new(RenameOldKind)];
let current = json!({ "kind": "old_kind", "seq": 1, SCHEMA_VERSION_FIELD: 2 });
assert_eq!(
apply_upcasters(&chain, vec![current.clone()]),
vec![current],
"an event at the version the step produces is not stepped again"
);
let out = apply_upcasters(&chain, vec![json!(42), json!({ "kind": "note", "seq": 1 })]);
assert_eq!(out[0], json!(42), "a non-object passes straight through");
assert_eq!(kind_of(&out[1]), "note", "an unknown kind keeps its name");
}
}