use std::fmt::Debug;
use async_trait::async_trait;
use crate::core::{
DeadLetter, EffectKey, InboundEvent, RunId, StoreError, Subscription, Timestamp,
};
#[derive(Debug, Clone, PartialEq)]
pub struct BufferedEvent {
pub event: InboundEvent,
pub received_at: Timestamp,
}
#[derive(Debug, Clone, PartialEq)]
pub enum TargetedDelivery {
Matched(Subscription),
Duplicate,
NotWaiting,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Minter {
Nobody,
Operator(crate::core::Operator),
}
pub const ERASED_REASON: &str = "erased";
#[async_trait]
pub trait EventStore: Send + Sync + Debug {
fn tenant(&self) -> &str {
crate::core::TenantId::DEFAULT
}
async fn buffer(&self, event: &InboundEvent, at: Timestamp) -> Result<bool, StoreError>;
async fn subscribe(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError>;
async fn claim_for(
&self,
sub: &Subscription,
at: Timestamp,
) -> Result<Option<BufferedEvent>, StoreError>;
async fn match_waiter(
&self,
event: &InboundEvent,
at: Timestamp,
) -> Result<Option<Subscription>, StoreError>;
async fn deliver_to(
&self,
run: RunId,
event: &InboundEvent,
at: Timestamp,
) -> Result<TargetedDelivery, StoreError>;
async fn unsubscribe(&self, run: RunId, effect: EffectKey) -> Result<(), StoreError>;
async fn unsubscribe_run(&self, run: RunId) -> Result<usize, StoreError>;
async fn park_wait(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError>;
async fn parked_waits(&self, limit: usize) -> Result<Vec<Subscription>, StoreError>;
async fn erase_payload(&self, source: &str, id: &str) -> Result<bool, StoreError>;
async fn minter(&self, source: &str, id: &str) -> Result<Option<Minter>, StoreError>;
async fn sweep_unclaimed(
&self,
older_than: Timestamp,
reason: &str,
) -> Result<usize, StoreError>;
async fn dead_letters(&self, limit: usize) -> Result<Vec<DeadLetter>, StoreError>;
async fn waiting(&self, limit: usize) -> Result<Vec<Subscription>, StoreError>;
}