use std::sync::Arc;
use async_trait::async_trait;
use crate::case::{BufferedEvent, EventStore, TargetedDelivery};
use crate::core::{
DeadLetter, EffectKey, InboundEvent, RunId, StoreError, Subscription, TenantId, Timestamp,
};
use crate::journal::payload;
use super::KeyRing;
#[derive(Debug)]
pub struct SealedEvents {
inner: Arc<dyn EventStore>,
keys: Arc<dyn KeyRing>,
tenant: TenantId,
}
impl SealedEvents {
#[must_use]
pub fn wrap(inner: Arc<dyn EventStore>, keys: Arc<dyn KeyRing>, tenant: TenantId) -> Arc<Self> {
Arc::new(Self {
inner,
keys,
tenant,
})
}
fn scope_for(&self, event: &InboundEvent) -> String {
super::scope(
&self.tenant,
&format!("event/{}/{}", event.source, event.id),
)
}
fn aad(event: &InboundEvent) -> String {
format!("{}/{}", event.source, event.id)
}
async fn sealed(&self, event: &InboundEvent) -> Result<InboundEvent, StoreError> {
let plain = crate::core::canon::to_bytes(&event.payload).map_err(|e| {
StoreError::Backend(format!("an event payload would not serialise: {e}"))
})?;
let envelope = super::envelope::seal(
self.keys.as_ref(),
&self.scope_for(event),
Self::aad(event).as_bytes(),
&plain,
)
.await
.map_err(|e| StoreError::Backend(format!("sealing an event payload failed: {e}")))?;
let mut sealed = event.clone();
sealed.payload = payload::wrap(&envelope);
Ok(sealed)
}
async fn opened(&self, mut event: InboundEvent) -> InboundEvent {
let Some(envelope) = payload::unwrap(&event.payload) else {
return event;
};
let aad = Self::aad(&event);
if let Ok(plain) =
super::envelope::open(self.keys.as_ref(), aad.as_bytes(), &envelope).await
&& let Ok(value) = serde_json::from_slice(&plain)
{
event.payload = value;
}
event
}
}
#[async_trait]
impl EventStore for SealedEvents {
async fn buffer(&self, event: &InboundEvent, at: Timestamp) -> Result<bool, StoreError> {
self.inner.buffer(&self.sealed(event).await?, at).await
}
async fn claim_for(
&self,
sub: &Subscription,
at: Timestamp,
) -> Result<Option<BufferedEvent>, StoreError> {
let claimed = self.inner.claim_for(sub, at).await?;
Ok(match claimed {
Some(mut buffered) => {
buffered.event = self.opened(buffered.event).await;
Some(buffered)
}
None => None,
})
}
async fn match_waiter(
&self,
event: &InboundEvent,
at: Timestamp,
) -> Result<Option<Subscription>, StoreError> {
self.inner
.match_waiter(&self.sealed(event).await?, at)
.await
}
async fn deliver_to(
&self,
run: RunId,
event: &InboundEvent,
at: Timestamp,
) -> Result<TargetedDelivery, StoreError> {
self.inner
.deliver_to(run, &self.sealed(event).await?, at)
.await
}
async fn subscribe(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError> {
self.inner.subscribe(sub, at).await
}
async fn unsubscribe(&self, run: RunId, effect: EffectKey) -> Result<(), StoreError> {
self.inner.unsubscribe(run, effect).await
}
async fn sweep_unclaimed(
&self,
older_than: Timestamp,
reason: &str,
) -> Result<usize, StoreError> {
self.inner.sweep_unclaimed(older_than, reason).await
}
async fn dead_letters(&self, limit: usize) -> Result<Vec<DeadLetter>, StoreError> {
let letters = self.inner.dead_letters(limit).await?;
let mut out = Vec::with_capacity(letters.len());
for mut letter in letters {
letter.event = self.opened(letter.event).await;
out.push(letter);
}
Ok(out)
}
async fn waiting(&self, limit: usize) -> Result<Vec<Subscription>, StoreError> {
self.inner.waiting(limit).await
}
}