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};
type EventRow<'a> = (
&'a str,
&'a str,
&'a str,
&'a str,
i64,
&'a str,
i64,
u8,
u8,
&'a str,
);
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 SUBS_BY_TIME: TableDefinition<(&str, 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, &'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),
) -> 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, ""),
)
.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(
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 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)
}
#[async_trait]
#[allow(clippy::too_many_lines)]
impl EventStore for RedbStore {
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 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,
"",
),
)
.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 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 = (
tenant.as_str(),
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(
(
tenant.as_str(),
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(
(
tenant.as_str(),
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 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, _) = v.value();
(
kd.to_owned(),
pl.to_owned(),
ra,
hc,
dead,
src.to_owned(),
bid.to_owned(),
claimed_by.to_owned(),
)
})
else {
continue;
};
if row.0 == kind && (row.3 == 0 || row.7 == run) && row.4 == 0 {
hit = Some((id, row));
break 'outer;
}
}
}
match hit {
None => None,
Some((id, (kd, payload, received, _, _, src, bid, _))) => {
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,
"",
),
)
.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))?;
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)?,
},
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(&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, _) = v.value();
(
kd.to_owned(),
pl.to_owned(),
ra,
hc,
dead,
src.to_owned(),
bid.to_owned(),
)
});
let claimable = row
.as_ref()
.is_some_and(|(_, _, _, hc, dead, _, _)| *hc == 0 && *dead == 0);
if claimable {
let (kd, pl, ra, _, _, src, bid) =
row.expect("claimable implies present");
claim_row(
&mut events,
(&tenant, &id),
(&src, &bid, &kd, &pl, ra),
(&run, ts(at)),
)?;
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))?;
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,
));
}
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();
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 = 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,
"",
),
)
.map_err(|e| be(&e))?;
drop(events);
index_correlation(&w, &tenant, &id, &keys, ts(at))?;
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)?;
{
let mut subs = w.open_table(SUBS).map_err(|e| be(&e))?;
let mut doomed = 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 (_, _, _, _, 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.as_str(),
run.as_str(),
effect.as_str(),
ns.as_str(),
val.as_str(),
))
.map_err(|e| be(&e))?;
by_key
.remove((
tenant.as_str(),
kind.as_str(),
ns.as_str(),
val.as_str(),
created,
run.as_str(),
effect.as_str(),
))
.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))?;
}
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.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, _) = v.value();
(
kd.to_owned(),
pl.to_owned(),
ra,
hc,
d,
src.to_owned(),
bid.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(),
),
)
.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) = 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)?,
},
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.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((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
}
}