use async_trait::async_trait;
use redb::{ReadableDatabase, ReadableTable, TableDefinition};
use crate::case::{BufferedEvent, EventStore, TargetedDelivery};
use crate::core::{
CaseId, CorrelationKey, DeadLetter, EffectKey, InboundEvent, RunId, StoreError, Subscription,
Timestamp,
};
use super::redb::{MAX_STR, RedbStore, be, begin_write, decoded, is_sealed};
fn phase_from(s: &str) -> Result<crate::core::Phase, StoreError> {
decoded("step phase", s, crate::core::Phase::parse(s))
}
type EventRow<'a> = (
&'a str,
&'a str,
&'a str,
&'a str,
i64,
&'a str,
i64,
u8,
u8,
&'a str,
&'a str,
&'a str,
);
fn encode_minter(by: Option<&crate::core::Operator>) -> (String, String) {
by.map_or_else(
|| (String::new(), String::new()),
|o| (o.actor().to_owned(), o.basis().as_str().to_owned()),
)
}
fn decode_minter(actor: &str, basis: &str) -> Result<Option<crate::core::Operator>, StoreError> {
match (actor.is_empty(), basis.is_empty()) {
(true, true) => Ok(None),
(false, false) => super::decode_operator(actor, basis, "inbound_events").map(Some),
_ => Err(StoreError::Corrupt {
seq: 0,
detail: "inbound_events holds half an operator: a minted event carries a \
name and what established it, or neither"
.to_owned(),
}),
}
}
const EVENTS: TableDefinition<(&str, &str), EventRow<'static>> =
TableDefinition::new("inbound_events");
const EVENT_CORR: TableDefinition<(&str, &str, &str, &str), ()> =
TableDefinition::new("inbound_correlation");
const EVENT_BY_KEY: TableDefinition<(&str, &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, &str), SubRow<'static>> =
TableDefinition::new("subscriptions");
type SubKey<'a> = (&'a str, &'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<(&str, i64, &str), ()> = TableDefinition::new("inbound_live");
const EVENTS_DEAD: TableDefinition<(&str, i64, &str), ()> = TableDefinition::new("inbound_dead");
const EVENTS_CLAIMED: TableDefinition<(&str, &str, &str), ()> =
TableDefinition::new("inbound_claimed");
const SUBS_BY_TIME: TableDefinition<(&str, i64, &str, &str, &str, &str), ()> =
TableDefinition::new("subscriptions_by_time");
const PARKED: TableDefinition<(&str, &str, &str), i64> =
TableDefinition::new("subscriptions_parked");
const EVENTS_ERASED: TableDefinition<(&str, &str), ()> = TableDefinition::new("inbound_erased");
pub(super) fn create_tables(w: &redb::WriteTransaction) -> Result<(), StoreError> {
w.open_table(EVENTS_ERASED).map_err(|e| be(&e))?;
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(EVENTS_CLAIMED).map_err(|e| be(&e))?;
w.open_table(SUBS_BY_TIME).map_err(|e| be(&e))?;
w.open_table(PARKED).map_err(|e| be(&e))?;
Ok(())
}
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, &'static str), ()>,
tenant: &str,
event_id: &str,
) -> Result<Vec<CorrelationKey>, StoreError> {
let mut out = Vec::new();
for e in t
.range((tenant, event_id, "", "")..=(tenant, 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 claim_row(
events: &mut redb::Table<'_, (&'static str, &'static str), EventRow<'static>>,
key: (&str, &str),
row: (&str, &str, &str, &str, i64),
claim: (&str, i64),
by: (&str, &str),
) -> Result<(), StoreError> {
let (tenant, id) = key;
let (src, bare, kind, payload, received) = row;
let (run, at) = claim;
events
.insert(
(tenant, id),
(
src, bare, kind, payload, received, run, at, 1u8, 0u8, "", by.0, by.1,
),
)
.map_err(|e| be(&e))?;
Ok(())
}
fn strip_payload(
events: &mut redb::Table<'_, (&'static str, &'static str), EventRow<'static>>,
tenant: &str,
key: &str,
) -> Result<bool, StoreError> {
let Some(row) = events.get((tenant, key)).map_err(|e| be(&e))?.map(|v| {
let (src, bid, kd, _, ra, cb, ca, hc, dead, reason, ba, bb) = v.value();
(
src.to_owned(),
bid.to_owned(),
kd.to_owned(),
ra,
cb.to_owned(),
ca,
hc,
dead,
reason.to_owned(),
ba.to_owned(),
bb.to_owned(),
)
}) else {
return Ok(false);
};
events
.insert(
(tenant, key),
(
row.0.as_str(),
row.1.as_str(),
row.2.as_str(),
"null",
row.3,
row.4.as_str(),
row.5,
row.6,
row.7,
row.8.as_str(),
row.9.as_str(),
row.10.as_str(),
),
)
.map_err(|e| be(&e))?;
Ok(true)
}
fn dead_letter_as_erased(
w: &redb::WriteTransaction,
events: &mut redb::Table<'_, (&'static str, &'static str), EventRow<'static>>,
tenant: &str,
key: &str,
received: i64,
) -> Result<(), StoreError> {
let Some(row) = events.get((tenant, key)).map_err(|e| be(&e))?.map(|v| {
let (src, bid, kd, pl, ra, _, _, _, _, _, ba, bb) = v.value();
(
src.to_owned(),
bid.to_owned(),
kd.to_owned(),
pl.to_owned(),
ra,
ba.to_owned(),
bb.to_owned(),
)
}) else {
return Ok(());
};
events
.insert(
(tenant, key),
(
row.0.as_str(),
row.1.as_str(),
row.2.as_str(),
row.3.as_str(),
row.4,
"",
0i64,
0u8,
1u8,
crate::case::ERASED_REASON,
row.5.as_str(),
row.6.as_str(),
),
)
.map_err(|e| be(&e))?;
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.remove((tenant, received, key))
.map_err(|e| be(&e))?;
w.open_table(EVENTS_DEAD)
.map_err(|e| be(&e))?
.insert((tenant, received, key), ())
.map_err(|e| be(&e))?;
Ok(())
}
fn unpark_for(
w: &redb::WriteTransaction,
tenant: &str,
run: &str,
event_id: &str,
) -> Result<(), StoreError> {
let keys = load_correlation(
&w.open_table(EVENT_CORR).map_err(|e| be(&e))?,
tenant,
event_id,
)?;
let effects: Vec<String> = w
.open_table(PARKED)
.map_err(|e| be(&e))?
.range((tenant, run, "")..=(tenant, run, MAX_STR))
.map_err(|e| be(&e))?
.map(|e| e.map(|(k, _)| k.value().2.to_owned()).map_err(|e| be(&e)))
.collect::<Result<_, _>>()?;
let subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut matching = Vec::new();
for effect in effects {
for k in &keys {
if subs
.get((
tenant,
run,
effect.as_str(),
k.namespace.as_str(),
k.value.as_str(),
))
.map_err(|e| be(&e))?
.is_some()
{
matching.push(effect.clone());
break;
}
}
}
drop(subs);
let mut parked = w.open_table(PARKED).map_err(|e| be(&e))?;
for effect in matching {
parked
.remove((tenant, run, effect.as_str()))
.map_err(|e| be(&e))?;
}
Ok(())
}
fn index_correlation(
w: &redb::WriteTransaction,
tenant: &str,
id: &str,
keys: &[CorrelationKey],
at: i64,
) -> Result<(), StoreError> {
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((tenant, id, k.namespace.as_str(), k.value.as_str()), ())
.map_err(|e| be(&e))?;
by_key
.insert((tenant, k.namespace.as_str(), k.value.as_str(), at, id), ())
.map_err(|e| be(&e))?;
}
Ok(())
}
fn oldest_waiter(
w: &redb::WriteTransaction,
tenant: &str,
by_key: &impl ReadableTable<SubKey<'static>, ()>,
subs: &impl ReadableTable<
(
&'static str,
&'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(
(
tenant,
kind,
k.namespace.as_str(),
k.value.as_str(),
i64::MIN,
"",
"",
)
..=(
tenant,
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 is_sealed(w, tenant, run)? {
continue;
}
if is_parked(w, tenant, run, effect)? {
continue;
}
if let Some(v) = subs
.get((tenant, 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)
}
fn is_parked(
w: &redb::WriteTransaction,
tenant: &str,
run: &str,
effect: &str,
) -> Result<bool, StoreError> {
Ok(w.open_table(PARKED)
.map_err(|e| be(&e))?
.get((tenant, run, effect))
.map_err(|e| be(&e))?
.is_some())
}
fn register_wait(
w: &redb::WriteTransaction,
tenant: &str,
sub: &Subscription,
at: i64,
) -> 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 phase = sub.phase.as_str();
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 &sub.correlation {
let key = (
tenant,
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,
sub.step.0,
phase,
sub.kind.as_str(),
at,
),
)
.map_err(|e| be(&e))?;
by_key
.insert(
(
tenant,
sub.kind.as_str(),
k.namespace.as_str(),
k.value.as_str(),
at,
run.as_str(),
effect.as_str(),
),
(),
)
.map_err(|e| be(&e))?;
by_time
.insert(
(
tenant,
at,
run.as_str(),
effect.as_str(),
k.namespace.as_str(),
k.value.as_str(),
),
(),
)
.map_err(|e| be(&e))?;
}
}
Ok(())
}
fn drop_wait(
w: &redb::WriteTransaction,
tenant: &str,
run: &str,
effect: &str,
) -> Result<(), StoreError> {
let mut subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut doomed = Vec::new();
for e in subs
.range((tenant, run, effect, "", "")..=(tenant, run, effect, 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((tenant, run, effect, ns.as_str(), val.as_str()))
.map_err(|e| be(&e))?;
by_key
.remove((
tenant,
kind.as_str(),
ns.as_str(),
val.as_str(),
created,
run,
effect,
))
.map_err(|e| be(&e))?;
by_time
.remove((tenant, created, run, effect, ns.as_str(), val.as_str()))
.map_err(|e| be(&e))?;
}
w.open_table(PARKED)
.map_err(|e| be(&e))?
.remove((tenant, run, effect))
.map_err(|e| be(&e))?;
Ok(())
}
fn shed_claimed(w: &redb::WriteTransaction, tenant: &str, run: &str) -> Result<(), StoreError> {
let mut claimed = w.open_table(EVENTS_CLAIMED).map_err(|e| be(&e))?;
let held: Vec<String> = claimed
.range((tenant, run, "")..=(tenant, run, MAX_STR))
.map_err(|e| be(&e))?
.map(|entry| {
entry
.map(|(key, _)| key.value().2.to_owned())
.map_err(|error| be(&error))
})
.collect::<Result<_, _>>()?;
let mut events = w.open_table(EVENTS).map_err(|e| be(&e))?;
for id in held {
strip_payload(&mut events, tenant, &id)?;
claimed
.remove((tenant, run, id.as_str()))
.map_err(|e| be(&e))?;
}
Ok(())
}
#[async_trait]
#[allow(clippy::too_many_lines)]
impl EventStore for RedbStore {
fn tenant(&self) -> &str {
self.tenant_str()
}
async fn buffer(&self, event: &InboundEvent, at: Timestamp) -> Result<bool, StoreError> {
let tenant = self.tenant_name();
let id = event.dedup_key();
let bare_id = event.id.clone();
let source = event.source.clone();
let (by_actor, by_basis) = encode_minter(event.by.as_ref());
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((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.is_some()
{
false
} else {
ev.insert(
(tenant.as_str(), id.as_str()),
(
source.as_str(),
bare_id.as_str(),
kind.as_str(),
payload.as_str(),
ts(at),
"",
0i64,
0u8,
0u8,
"",
by_actor.as_str(),
by_basis.as_str(),
),
)
.map_err(|e| be(&e))?;
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.insert((tenant.as_str(), ts(at), id.as_str()), ())
.map_err(|e| be(&e))?;
index_correlation(&w, &tenant, &id, &keys, ts(at))?;
true
}
};
w.commit().map_err(|e| be(&e))?;
Ok(fresh)
})
.await
}
async fn subscribe(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError> {
let tenant = self.tenant_name();
let sub = sub.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
register_wait(&w, &tenant, &sub, ts(at))?;
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
async fn claim_for(
&self,
sub: &Subscription,
at: Timestamp,
) -> Result<Option<BufferedEvent>, StoreError> {
let tenant = self.tenant_name();
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(
(
tenant.as_str(),
k.namespace.as_str(),
k.value.as_str(),
i64::MIN,
"",
)
..=(
tenant.as_str(),
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().4.to_owned();
let Some(row) = events
.get((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (src, bid, kd, pl, ra, claimed_by, _, hc, dead, _, ba, bb) =
v.value();
(
kd.to_owned(),
pl.to_owned(),
ra,
hc,
dead,
src.to_owned(),
bid.to_owned(),
claimed_by.to_owned(),
ba.to_owned(),
bb.to_owned(),
)
})
else {
continue;
};
let erased = w
.open_table(EVENTS_ERASED)
.map_err(|e| be(&e))?
.get((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.is_some();
if row.0 == kind && (row.3 == 0 || row.7 == run) && row.4 == 0 && !erased {
hit = Some((id, row));
break 'outer;
}
}
}
match hit {
None => None,
Some((id, (kd, payload, received, _, _, src, bid, _, ba, bb))) => {
events
.insert(
(tenant.as_str(), id.as_str()),
(
src.as_str(),
bid.as_str(),
kd.as_str(),
payload.as_str(),
received,
run.as_str(),
ts(at),
1u8,
0u8,
"",
ba.as_str(),
bb.as_str(),
),
)
.map_err(|e| be(&e))?;
drop(events);
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.remove((tenant.as_str(), received, id.as_str()))
.map_err(|e| be(&e))?;
w.open_table(EVENTS_CLAIMED)
.map_err(|e| be(&e))?
.insert((tenant.as_str(), run.as_str(), 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, &tenant, &id)?;
Some(BufferedEvent {
event: InboundEvent {
source: src,
id: bid,
kind: kd,
correlation,
payload: serde_json::from_str(&payload)?,
by: decode_minter(&ba, &bb)?,
},
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 tenant = self.tenant_name();
let id = event.dedup_key();
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 mut subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let hit = oldest_waiter(&w, &tenant, &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((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (src, bid, kd, pl, ra, _, _, hc, dead, _, ba, bb) = v.value();
(
kd.to_owned(),
pl.to_owned(),
ra,
hc,
dead,
src.to_owned(),
bid.to_owned(),
ba.to_owned(),
bb.to_owned(),
)
});
let claimable = row
.as_ref()
.is_some_and(|(_, _, _, hc, dead, ..)| *hc == 0 && *dead == 0);
if claimable {
let (kd, pl, ra, _, _, src, bid, ba, bb) =
row.expect("claimable implies present");
claim_row(
&mut events,
(&tenant, &id),
(&src, &bid, &kd, &pl, ra),
(&run, ts(at)),
(&ba, &bb),
)?;
drop(events);
w.open_table(EVENTS_LIVE)
.map_err(|e| be(&e))?
.remove((tenant.as_str(), ra, id.as_str()))
.map_err(|e| be(&e))?;
w.open_table(EVENTS_CLAIMED)
.map_err(|e| be(&e))?
.insert((tenant.as_str(), run.as_str(), id.as_str()), ())
.map_err(|e| be(&e))?;
let mut correlation = Vec::new();
let mut retired = Vec::new();
for e in subs
.range(
(tenant.as_str(), run.as_str(), effect.as_str(), "", "")
..=(
tenant.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 (_, _, _, _, sub_kind, created) = v.value();
correlation
.push(CorrelationKey::new(ns.to_owned(), val.to_owned()));
retired.push((
ns.to_owned(),
val.to_owned(),
sub_kind.to_owned(),
created,
));
}
w.open_table(PARKED)
.map_err(|e| be(&e))?
.remove((tenant.as_str(), run.as_str(), effect.as_str()))
.map_err(|e| be(&e))?;
for (ns, val, sub_kind, created) in retired {
subs.remove((
tenant.as_str(),
run.as_str(),
effect.as_str(),
ns.as_str(),
val.as_str(),
))
.map_err(|e| be(&e))?;
let mut by_key = w.open_table(SUBS_BY_KEY).map_err(|e| be(&e))?;
by_key
.remove((
tenant.as_str(),
sub_kind.as_str(),
ns.as_str(),
val.as_str(),
created,
run.as_str(),
effect.as_str(),
))
.map_err(|e| be(&e))?;
let mut by_time = w.open_table(SUBS_BY_TIME).map_err(|e| be(&e))?;
by_time
.remove((
tenant.as_str(),
created,
run.as_str(),
effect.as_str(),
ns.as_str(),
val.as_str(),
))
.map_err(|e| be(&e))?;
}
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 deliver_to(
&self,
target: RunId,
event: &InboundEvent,
at: Timestamp,
) -> Result<TargetedDelivery, StoreError> {
let tenant = self.tenant_name();
let run = target.to_string();
let id = event.dedup_key();
let bare_id = event.id.clone();
let source = event.source.clone();
let kind = event.kind.clone();
let payload = serde_json::to_string(&event.payload)?;
let keys = event.correlation.clone();
let (by_actor, by_basis) = encode_minter(event.by.as_ref());
self.with_db(move |db| {
let w = begin_write(db)?;
let outcome = {
let mut events = w.open_table(EVENTS).map_err(|e| be(&e))?;
let existing_claim = events
.get((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.map(|row| {
let (_, _, _, _, _, claimed_by, _, has_claim, _, _, _, _) = row.value();
(claimed_by.to_owned(), has_claim)
});
let subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut selected: Option<(String, String, u8, u32, String)> = None;
for row in subs
.range(
(tenant.as_str(), run.as_str(), "", "", "")
..=(tenant.as_str(), run.as_str(), MAX_STR, MAX_STR, MAX_STR),
)
.map_err(|e| be(&e))?
{
let (key, value) = row.map_err(|e| be(&e))?;
let (_, _, effect, namespace, value_key) = key.value();
let (case, has_case, step, phase, event_kind, _) = value.value();
if event_kind == kind
&& keys.iter().any(|candidate| {
candidate.namespace == namespace && candidate.value == value_key
})
{
selected = Some((
effect.to_owned(),
case.to_owned(),
has_case,
step,
phase.to_owned(),
));
break;
}
}
let Some((effect, case, has_case, step, phase)) = selected else {
drop(subs);
drop(events);
w.commit().map_err(|e| be(&e))?;
return Ok(if existing_claim.is_some() {
TargetedDelivery::Duplicate
} else {
TargetedDelivery::NotWaiting
});
};
let mut correlation = Vec::new();
for row in subs
.range(
(tenant.as_str(), run.as_str(), effect.as_str(), "", "")
..=(
tenant.as_str(),
run.as_str(),
effect.as_str(),
MAX_STR,
MAX_STR,
),
)
.map_err(|e| be(&e))?
{
let (key, _) = row.map_err(|e| be(&e))?;
let (_, _, _, namespace, value) = key.value();
correlation.push(CorrelationKey::new(namespace.to_owned(), value.to_owned()));
}
drop(subs);
let subscription_effect = effect.clone();
let subscription = Subscription {
run: target,
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,
};
if let Some((claimed_by, has_claim)) = existing_claim {
drop(events);
if has_claim == 1 && claimed_by == run {
TargetedDelivery::Matched(subscription)
} else {
TargetedDelivery::Duplicate
}
} else {
events
.insert(
(tenant.as_str(), id.as_str()),
(
source.as_str(),
bare_id.as_str(),
kind.as_str(),
payload.as_str(),
ts(at),
run.as_str(),
ts(at),
1u8,
0u8,
"",
by_actor.as_str(),
by_basis.as_str(),
),
)
.map_err(|e| be(&e))?;
drop(events);
index_correlation(&w, &tenant, &id, &keys, ts(at))?;
w.open_table(EVENTS_CLAIMED)
.map_err(|e| be(&e))?
.insert((tenant.as_str(), run.as_str(), id.as_str()), ())
.map_err(|e| be(&e))?;
w.open_table(PARKED)
.map_err(|e| be(&e))?
.insert(
(tenant.as_str(), run.as_str(), subscription_effect.as_str()),
ts(at),
)
.map_err(|e| be(&e))?;
TargetedDelivery::Matched(subscription)
}
};
w.commit().map_err(|e| be(&e))?;
Ok(outcome)
})
.await
}
async fn unsubscribe(&self, run: RunId, effect: EffectKey) -> Result<(), StoreError> {
let tenant = self.tenant_name();
let (run, effect) = (run.to_string(), effect.to_hex());
self.with_db(move |db| {
let w = begin_write(db)?;
drop_wait(&w, &tenant, &run, &effect)?;
shed_claimed(&w, &tenant, &run)?;
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
async fn unsubscribe_run(&self, run: RunId) -> Result<usize, StoreError> {
let tenant = self.tenant_name();
let run = run.to_string();
self.with_db(move |db| {
let w = begin_write(db)?;
let effects: std::collections::BTreeSet<String> = {
let subs = w.open_table(SUBS).map_err(|e| be(&e))?;
subs.range(
(tenant.as_str(), run.as_str(), "", "", "")
..=(tenant.as_str(), run.as_str(), MAX_STR, MAX_STR, MAX_STR),
)
.map_err(|e| be(&e))?
.map(|e| e.map(|(k, _)| k.value().2.to_owned()).map_err(|e| be(&e)))
.collect::<Result<_, _>>()?
};
for effect in &effects {
drop_wait(&w, &tenant, &run, effect)?;
}
shed_claimed(&w, &tenant, &run)?;
w.commit().map_err(|e| be(&e))?;
Ok(effects.len())
})
.await
}
async fn park_wait(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError> {
let tenant = self.tenant_name();
let sub = sub.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
register_wait(&w, &tenant, &sub, ts(at))?;
w.open_table(PARKED)
.map_err(|e| be(&e))?
.insert(
(
tenant.as_str(),
sub.run.to_string().as_str(),
sub.effect.to_hex().as_str(),
),
ts(at),
)
.map_err(|e| be(&e))?;
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
async fn parked_waits(&self, limit: usize) -> Result<Vec<Subscription>, StoreError> {
let tenant = self.tenant_name();
self.with_db(move |db| {
let w = begin_write(db)?;
let listed: Vec<(String, String)> = {
let parked = w.open_table(PARKED).map_err(|e| be(&e))?;
parked
.range((tenant.as_str(), "", "")..=(tenant.as_str(), MAX_STR, MAX_STR))
.map_err(|e| be(&e))?
.map(|e| {
e.map(|(k, _)| {
let (_, run, effect) = k.value();
(run.to_owned(), effect.to_owned())
})
.map_err(|e| be(&e))
})
.collect::<Result<_, _>>()?
};
let mut live = Vec::new();
for (run, effect) in listed {
if is_sealed(&w, &tenant, &run)? {
drop_wait(&w, &tenant, &run, &effect)?;
shed_claimed(&w, &tenant, &run)?;
} else {
live.push((run, effect));
}
}
let subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut out = Vec::new();
for (run, effect) in &live {
if out.len() >= limit {
break;
}
let (run, effect) = (run.as_str(), effect.as_str());
let mut found: Option<(String, u8, u32, String, String)> = None;
let mut correlation = Vec::new();
for row in subs
.range(
(tenant.as_str(), run, effect, "", "")
..=(tenant.as_str(), run, effect, MAX_STR, MAX_STR),
)
.map_err(|e| be(&e))?
{
let (key, value) = row.map_err(|e| be(&e))?;
let (_, _, _, ns, val) = key.value();
let (case, has_case, step, phase, kind, _) = value.value();
correlation.push(CorrelationKey::new(ns.to_owned(), val.to_owned()));
found.get_or_insert_with(|| {
(
case.to_owned(),
has_case,
step,
phase.to_owned(),
kind.to_owned(),
)
});
}
let Some((case, has_case, step, phase, kind)) = found else {
continue;
};
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,
correlation,
});
}
drop(subs);
w.commit().map_err(|e| be(&e))?;
Ok(out)
})
.await
}
async fn erase_payload(&self, source: &str, id: &str) -> Result<bool, StoreError> {
let tenant = self.tenant_name();
let key = InboundEvent {
source: source.to_owned(),
id: id.to_owned(),
kind: String::new(),
correlation: Vec::new(),
payload: serde_json::Value::Null,
by: None,
}
.dedup_key();
self.with_db(move |db| {
let w = begin_write(db)?;
let existed = {
let mut events = w.open_table(EVENTS).map_err(|e| be(&e))?;
let row = events
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (_, _, _, _, received, by, _, claimed, dead, _, _, _) = v.value();
(received, by.to_owned(), claimed == 1, dead == 1)
});
let existed = strip_payload(&mut events, &tenant, &key)?;
if existed {
w.open_table(EVENTS_ERASED)
.map_err(|e| be(&e))?
.insert((tenant.as_str(), key.as_str()), ())
.map_err(|e| be(&e))?;
}
if let Some((received, claimant, claimed, dead)) = row {
let undelivered_claim = claimed
&& w.open_table(EVENTS_CLAIMED)
.map_err(|e| be(&e))?
.remove((tenant.as_str(), claimant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.is_some();
if (!claimed && !dead) || undelivered_claim {
dead_letter_as_erased(&w, &mut events, &tenant, &key, received)?;
}
if undelivered_claim {
drop(events);
unpark_for(&w, &tenant, &claimant, &key)?;
}
}
existed
};
w.commit().map_err(|e| be(&e))?;
Ok(existed)
})
.await
}
async fn minter(
&self,
source: &str,
id: &str,
) -> Result<Option<crate::case::Minter>, StoreError> {
let tenant = self.tenant_name();
let key = crate::core::origin_key(source, id);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let Ok(events) = r.open_table(EVENTS) else {
return Ok(None);
};
let Some(row) = events
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (.., ba, bb) = v.value();
(ba.to_owned(), bb.to_owned())
})
else {
return Ok(None);
};
Ok(Some(decode_minter(&row.0, &row.1)?.map_or(
crate::case::Minter::Nobody,
crate::case::Minter::Operator,
)))
})
.await
}
async fn sweep_unclaimed(
&self,
older_than: Timestamp,
reason: &str,
) -> Result<usize, StoreError> {
let tenant = self.tenant_name();
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((tenant.as_str(), i64::MIN, "")..=(tenant.as_str(), 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((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (src, bid, kd, pl, ra, _, _, hc, d, _, ba, bb) = v.value();
(
kd.to_owned(),
pl.to_owned(),
ra,
hc,
d,
src.to_owned(),
bid.to_owned(),
ba.to_owned(),
bb.to_owned(),
)
})
else {
continue;
};
events
.insert(
(tenant.as_str(), id.as_str()),
(
row.5.as_str(),
row.6.as_str(),
row.0.as_str(),
row.1.as_str(),
row.2,
"",
0i64,
0u8,
1u8,
reason.as_str(),
row.7.as_str(),
row.8.as_str(),
),
)
.map_err(|e| be(&e))?;
live.remove((tenant.as_str(), at, id.as_str()))
.map_err(|e| be(&e))?;
dead.insert((tenant.as_str(), 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> {
let tenant = self.tenant_name();
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
.range((tenant.as_str(), i64::MIN, "")..=(tenant.as_str(), i64::MAX, MAX_STR))
.map_err(|e| be(&e))?
.rev()
{
if out.len() >= limit {
break;
}
let (k, _) = e.map_err(|e| be(&e))?;
let id = k.value().2;
let Some(v) = events.get((tenant.as_str(), id)).map_err(|e| be(&e))? else {
continue;
};
let (source, bare, kind, payload, received, _, _, _, _, reason, ba, bb) = v.value();
out.push(DeadLetter {
event: InboundEvent {
source: source.to_owned(),
id: bare.to_owned(),
kind: kind.to_owned(),
correlation: load_correlation(&corr, &tenant, id)?,
payload: serde_json::from_str(payload)?,
by: decode_minter(ba, bb)?,
},
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> {
let tenant = self.tenant_name();
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
.range(
(tenant.as_str(), i64::MIN, "", "", "", "")
..=(
tenant.as_str(),
i64::MAX,
MAX_STR,
MAX_STR,
MAX_STR,
MAX_STR,
),
)
.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((tenant.as_str(), 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
}
}
#[cfg(test)]
mod codec_tests {
use super::*;
#[test]
fn every_written_phase_decodes_to_the_value_that_wrote_it() {
for phase in [
crate::core::Phase::Forward,
crate::core::Phase::Compensating,
] {
assert_eq!(phase_from(phase.as_str()).expect("round trip"), phase);
}
}
#[test]
fn an_unreadable_subscription_phase_is_refused_rather_than_defaulted() {
for bad in ["", "Forward", "compensating ", "backward"] {
assert!(
matches!(phase_from(bad).err(), Some(StoreError::Corrupt { .. })),
"phase '{bad}' decoded instead of refusing"
);
}
}
}