use async_trait::async_trait;
use redb::{ReadableDatabase, ReadableTable, TableDefinition};
use crate::case::{BufferedEvent, EventStore};
use crate::core::{
CaseId, CorrelationKey, DeadLetter, EffectKey, InboundEvent, RunId, StoreError, Subscription,
Timestamp,
};
use super::redb::{MAX_STR, RedbStore, be, begin_write};
type EventRow<'a> = (&'a str, &'a str, i64, &'a str, i64, u8, u8, &'a str);
const EVENTS: TableDefinition<&str, EventRow<'static>> = TableDefinition::new("inbound_events");
const EVENT_CORR: TableDefinition<(&str, &str, &str), ()> =
TableDefinition::new("inbound_correlation");
const EVENT_BY_KEY: TableDefinition<(&str, &str, i64, &str), ()> =
TableDefinition::new("inbound_by_key");
type SubRow<'a> = (&'a str, u8, u32, &'a str, &'a str, i64);
const SUBS: TableDefinition<(&str, &str, &str, &str), SubRow<'static>> =
TableDefinition::new("subscriptions");
type SubKey<'a> = (&'a str, &'a str, &'a str, i64, &'a str, &'a str);
const SUBS_BY_KEY: TableDefinition<SubKey<'static>, ()> =
TableDefinition::new("subscriptions_by_key");
const EVENTS_LIVE: TableDefinition<(i64, &str), ()> = TableDefinition::new("inbound_live");
const EVENTS_DEAD: TableDefinition<(i64, &str), ()> = TableDefinition::new("inbound_dead");
const SUBS_BY_TIME: TableDefinition<(i64, &str, &str, &str, &str), ()> =
TableDefinition::new("subscriptions_by_time");
pub(super) fn create_tables(w: &redb::WriteTransaction) -> Result<(), StoreError> {
w.open_table(EVENTS).map_err(|e| be(&e))?;
w.open_table(EVENT_CORR).map_err(|e| be(&e))?;
w.open_table(EVENT_BY_KEY).map_err(|e| be(&e))?;
w.open_table(SUBS).map_err(|e| be(&e))?;
w.open_table(SUBS_BY_KEY).map_err(|e| be(&e))?;
w.open_table(EVENTS_LIVE).map_err(|e| be(&e))?;
w.open_table(EVENTS_DEAD).map_err(|e| be(&e))?;
w.open_table(SUBS_BY_TIME).map_err(|e| be(&e))?;
Ok(())
}
fn phase_str(p: crate::core::Phase) -> &'static str {
match p {
crate::core::Phase::Forward => "forward",
crate::core::Phase::Compensating => "compensating",
}
}
fn phase_from(s: &str) -> crate::core::Phase {
match s {
"compensating" => crate::core::Phase::Compensating,
_ => crate::core::Phase::Forward,
}
}
fn ts(t: Timestamp) -> i64 {
t.unix_timestamp()
}
fn from_ts(v: i64) -> Result<Timestamp, StoreError> {
Timestamp::from_unix_timestamp(v).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("unrepresentable timestamp {v}: {e}"),
})
}
fn load_correlation(
t: &impl ReadableTable<(&'static str, &'static str, &'static str), ()>,
event_id: &str,
) -> Result<Vec<CorrelationKey>, StoreError> {
let mut out = Vec::new();
for e in t
.range((event_id, "", "")..=(event_id, MAX_STR, MAX_STR))
.map_err(|e| be(&e))?
{
let (k, _) = e.map_err(|e| be(&e))?;
let (_, ns, v) = k.value();
out.push(CorrelationKey::new(ns.to_owned(), v.to_owned()));
}
Ok(out)
}
type Waiter = (String, String, String, u8, u32, String);
fn oldest_waiter(
by_key: &impl ReadableTable<SubKey<'static>, ()>,
subs: &impl ReadableTable<(&'static str, &'static str, &'static str, &'static str), SubRow<'static>>,
kind: &str,
keys: &[CorrelationKey],
) -> Result<Option<Waiter>, StoreError> {
for k in keys {
for e in by_key
.range(
(
kind,
k.namespace.as_str(),
k.value.as_str(),
i64::MIN,
"",
"",
)
..=(
kind,
k.namespace.as_str(),
k.value.as_str(),
i64::MAX,
MAX_STR,
MAX_STR,
),
)
.map_err(|e| be(&e))?
{
let (sk, _) = e.map_err(|e| be(&e))?;
let (_, ns, val, _, run, effect) = sk.value();
if let Some(v) = subs.get((run, effect, ns, val)).map_err(|e| be(&e))? {
let (case, has_case, step, phase, _, _) = v.value();
return Ok(Some((
run.to_owned(),
effect.to_owned(),
case.to_owned(),
has_case,
step,
phase.to_owned(),
)));
}
}
}
Ok(None)
}
#[async_trait]
impl EventStore for RedbStore {
async fn buffer(&self, event: &InboundEvent, at: Timestamp) -> Result<bool, StoreError> {
let id = event.id.clone();
let kind = event.kind.clone();
let payload = serde_json::to_string(&event.payload)?;
let keys = event.correlation.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
let fresh = {
let mut ev = w.open_table(EVENTS).map_err(|e| be(&e))?;
if ev.get(id.as_str()).map_err(|e| be(&e))?.is_some() {
false
} else {
ev.insert(
id.as_str(),
(
kind.as_str(),
payload.as_str(),
ts(at),
"",
0i64,
0u8,
0u8,
"",
),
)
.map_err(|e| be(&e))?;
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.insert((ts(at), id.as_str()), ())
.map_err(|e| be(&e))?;
let mut corr = w.open_table(EVENT_CORR).map_err(|e| be(&e))?;
let mut by_key = w.open_table(EVENT_BY_KEY).map_err(|e| be(&e))?;
for k in &keys {
corr.insert((id.as_str(), k.namespace.as_str(), k.value.as_str()), ())
.map_err(|e| be(&e))?;
by_key
.insert(
(k.namespace.as_str(), k.value.as_str(), ts(at), id.as_str()),
(),
)
.map_err(|e| be(&e))?;
}
true
}
};
w.commit().map_err(|e| be(&e))?;
Ok(fresh)
})
.await
}
async fn subscribe(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError> {
let run = sub.run.to_string();
let effect = sub.effect.to_hex();
let case = sub.case.map(|c| c.to_string()).unwrap_or_default();
let has_case = u8::from(sub.case.is_some());
let step = sub.step.0;
let phase = phase_str(sub.phase);
let kind = sub.kind.clone();
let keys = sub.correlation.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
{
let mut subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut by_key = w.open_table(SUBS_BY_KEY).map_err(|e| be(&e))?;
let mut by_time = w.open_table(SUBS_BY_TIME).map_err(|e| be(&e))?;
for k in &keys {
let key = (
run.as_str(),
effect.as_str(),
k.namespace.as_str(),
k.value.as_str(),
);
if subs.get(key).map_err(|e| be(&e))?.is_none() {
subs.insert(
key,
(case.as_str(), has_case, step, phase, kind.as_str(), ts(at)),
)
.map_err(|e| be(&e))?;
by_key
.insert(
(
kind.as_str(),
k.namespace.as_str(),
k.value.as_str(),
ts(at),
run.as_str(),
effect.as_str(),
),
(),
)
.map_err(|e| be(&e))?;
by_time
.insert(
(
ts(at),
run.as_str(),
effect.as_str(),
k.namespace.as_str(),
k.value.as_str(),
),
(),
)
.map_err(|e| be(&e))?;
}
}
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
async fn claim_for(
&self,
sub: &Subscription,
at: Timestamp,
) -> Result<Option<BufferedEvent>, StoreError> {
let run = sub.run.to_string();
let kind = sub.kind.clone();
let keys = sub.correlation.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
let found = {
let by_key = w.open_table(EVENT_BY_KEY).map_err(|e| be(&e))?;
let mut events = w.open_table(EVENTS).map_err(|e| be(&e))?;
let mut hit = None;
'outer: for k in &keys {
for e in by_key
.range(
(k.namespace.as_str(), k.value.as_str(), i64::MIN, "")
..=(k.namespace.as_str(), k.value.as_str(), i64::MAX, MAX_STR),
)
.map_err(|e| be(&e))?
{
let (ek, _) = e.map_err(|e| be(&e))?;
let id = ek.value().3.to_owned();
let Some(row) = events.get(id.as_str()).map_err(|e| be(&e))?.map(|v| {
let (kd, pl, ra, _, _, hc, dead, _) = v.value();
(kd.to_owned(), pl.to_owned(), ra, hc, dead)
}) else {
continue;
};
if row.0 == kind && row.3 == 0 && row.4 == 0 {
hit = Some((id, row));
break 'outer;
}
}
}
match hit {
None => None,
Some((id, (kd, payload, received, _, _))) => {
events
.insert(
id.as_str(),
(
kd.as_str(),
payload.as_str(),
received,
run.as_str(),
ts(at),
1u8,
0u8,
"",
),
)
.map_err(|e| be(&e))?;
drop(events);
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.remove((received, id.as_str()))
.map_err(|e| be(&e))?;
let corr_t = w.open_table(EVENT_CORR).map_err(|e| be(&e))?;
let correlation = load_correlation(&corr_t, &id)?;
Some(BufferedEvent {
event: InboundEvent {
id,
kind: kd,
correlation,
payload: serde_json::from_str(&payload)?,
},
received_at: from_ts(received)?,
})
}
}
};
w.commit().map_err(|e| be(&e))?;
Ok(found)
})
.await
}
async fn match_waiter(
&self,
event: &InboundEvent,
at: Timestamp,
) -> Result<Option<Subscription>, StoreError> {
let id = event.id.clone();
let kind = event.kind.clone();
let keys = event.correlation.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
let found = {
let by_key = w.open_table(SUBS_BY_KEY).map_err(|e| be(&e))?;
let subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let hit = oldest_waiter(&by_key, &subs, &kind, &keys)?;
drop(by_key);
match hit {
None => None,
Some((run, effect, case, has_case, step, phase)) => {
let mut events = w.open_table(EVENTS).map_err(|e| be(&e))?;
let row = events.get(id.as_str()).map_err(|e| be(&e))?.map(|v| {
let (kd, pl, ra, _, _, hc, dead, _) = v.value();
(kd.to_owned(), pl.to_owned(), ra, hc, dead)
});
let claimable = row
.as_ref()
.is_some_and(|(_, _, _, hc, dead)| *hc == 0 && *dead == 0);
if claimable {
let (kd, pl, ra, _, _) = row.expect("claimable implies present");
events
.insert(
id.as_str(),
(
kd.as_str(),
pl.as_str(),
ra,
run.as_str(),
ts(at),
1u8,
0u8,
"",
),
)
.map_err(|e| be(&e))?;
drop(events);
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.remove((ra, id.as_str()))
.map_err(|e| be(&e))?;
let mut correlation = Vec::new();
for e in subs
.range(
(run.as_str(), effect.as_str(), "", "")
..=(run.as_str(), effect.as_str(), MAX_STR, MAX_STR),
)
.map_err(|e| be(&e))?
{
let (k, _) = e.map_err(|e| be(&e))?;
let (_, _, ns, v) = k.value();
correlation.push(CorrelationKey::new(ns.to_owned(), v.to_owned()));
}
Some(Subscription {
run: RunId::parse(&run).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("bad run id '{run}': {e}"),
})?,
case: if has_case == 1 {
Some(CaseId::parse(&case).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("bad case id '{case}': {e}"),
})?)
} else {
None
},
effect: EffectKey::from_hex(&effect).map_err(|e| {
StoreError::Corrupt {
seq: 0,
detail: format!("bad effect key '{effect}': {e}"),
}
})?,
step: crate::core::StepId(step),
phase: phase_from(&phase),
kind: kind.clone(),
correlation,
})
} else {
None
}
}
}
};
w.commit().map_err(|e| be(&e))?;
Ok(found)
})
.await
}
async fn unsubscribe(&self, run: RunId, effect: EffectKey) -> Result<(), StoreError> {
let (run, effect) = (run.to_string(), effect.to_hex());
self.with_db(move |db| {
let w = begin_write(db)?;
{
let mut subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut doomed = Vec::new();
for e in subs
.range(
(run.as_str(), effect.as_str(), "", "")
..=(run.as_str(), effect.as_str(), MAX_STR, MAX_STR),
)
.map_err(|e| be(&e))?
{
let (k, v) = e.map_err(|e| be(&e))?;
let (_, _, ns, val) = k.value();
let (_, _, _, _, kind, created) = v.value();
doomed.push((ns.to_owned(), val.to_owned(), kind.to_owned(), created));
}
let mut by_key = w.open_table(SUBS_BY_KEY).map_err(|e| be(&e))?;
let mut by_time = w.open_table(SUBS_BY_TIME).map_err(|e| be(&e))?;
for (ns, val, kind, created) in doomed {
subs.remove((run.as_str(), effect.as_str(), ns.as_str(), val.as_str()))
.map_err(|e| be(&e))?;
by_key
.remove((
kind.as_str(),
ns.as_str(),
val.as_str(),
created,
run.as_str(),
effect.as_str(),
))
.map_err(|e| be(&e))?;
by_time
.remove((
created,
run.as_str(),
effect.as_str(),
ns.as_str(),
val.as_str(),
))
.map_err(|e| be(&e))?;
}
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
async fn sweep_unclaimed(
&self,
older_than: Timestamp,
reason: &str,
) -> Result<usize, StoreError> {
let cutoff = ts(older_than);
let reason = reason.to_owned();
self.with_db(move |db| {
let w = begin_write(db)?;
let n = {
let live = w.open_table(EVENTS_LIVE).map_err(|e| be(&e))?;
let mut doomed = Vec::new();
for e in live
.range((i64::MIN, "")..=(cutoff, MAX_STR))
.map_err(|e| be(&e))?
{
let (k, _) = e.map_err(|e| be(&e))?;
let (at, id) = k.value();
doomed.push((at, id.to_owned()));
}
drop(live);
let mut events = w.open_table(EVENTS).map_err(|e| be(&e))?;
let mut live = w.open_table(EVENTS_LIVE).map_err(|e| be(&e))?;
let mut dead = w.open_table(EVENTS_DEAD).map_err(|e| be(&e))?;
let mut n = 0usize;
for (at, id) in doomed {
let Some(row) = events.get(id.as_str()).map_err(|e| be(&e))?.map(|v| {
let (kd, pl, ra, _, _, hc, d, _) = v.value();
(kd.to_owned(), pl.to_owned(), ra, hc, d)
}) else {
continue;
};
events
.insert(
id.as_str(),
(
row.0.as_str(),
row.1.as_str(),
row.2,
"",
0i64,
0u8,
1u8,
reason.as_str(),
),
)
.map_err(|e| be(&e))?;
live.remove((at, id.as_str())).map_err(|e| be(&e))?;
dead.insert((row.2, id.as_str()), ()).map_err(|e| be(&e))?;
n += 1;
}
n
};
w.commit().map_err(|e| be(&e))?;
Ok(n)
})
.await
}
async fn dead_letters(&self, limit: usize) -> Result<Vec<DeadLetter>, StoreError> {
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let dead = r.open_table(EVENTS_DEAD).map_err(|e| be(&e))?;
let events = r.open_table(EVENTS).map_err(|e| be(&e))?;
let corr = r.open_table(EVENT_CORR).map_err(|e| be(&e))?;
let mut out = Vec::new();
for e in dead.iter().map_err(|e| be(&e))?.rev() {
if out.len() >= limit {
break;
}
let (k, _) = e.map_err(|e| be(&e))?;
let id = k.value().1;
let Some(v) = events.get(id).map_err(|e| be(&e))? else {
continue;
};
let (kind, payload, received, _, _, _, _, reason) = v.value();
out.push(DeadLetter {
event: InboundEvent {
id: id.to_owned(),
kind: kind.to_owned(),
correlation: load_correlation(&corr, id)?,
payload: serde_json::from_str(payload)?,
},
received_at: from_ts(received)?,
reason: if reason.is_empty() {
"unclaimed".to_owned()
} else {
reason.to_owned()
},
});
}
Ok(out)
})
.await
}
async fn waiting(&self, limit: usize) -> Result<Vec<Subscription>, StoreError> {
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let by_time = r.open_table(SUBS_BY_TIME).map_err(|e| be(&e))?;
let subs = r.open_table(SUBS).map_err(|e| be(&e))?;
let mut out = Vec::new();
for e in by_time.iter().map_err(|e| be(&e))? {
if out.len() >= limit {
break;
}
let (k, _) = e.map_err(|e| be(&e))?;
let (_, run, effect, ns, val) = k.value();
let Some(v) = subs.get((run, effect, ns, val)).map_err(|e| be(&e))? else {
continue;
};
let (case, has_case, step, phase, kind, _) = v.value();
out.push(Subscription {
run: RunId::parse(run).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("bad run id '{run}': {e}"),
})?,
case: if has_case == 1 {
Some(CaseId::parse(case).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("bad case id '{case}': {e}"),
})?)
} else {
None
},
effect: EffectKey::from_hex(effect).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("bad effect key '{effect}': {e}"),
})?,
step: crate::core::StepId(step),
phase: phase_from(phase),
kind: kind.to_owned(),
correlation: vec![CorrelationKey::new(ns.to_owned(), val.to_owned())],
});
}
Ok(out)
})
.await
}
}