use std::path::Path;
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use async_trait::async_trait;
use redb::{Database, ReadableDatabase, ReadableTable, TableDefinition};
use crate::core::{Digest, EffectKey, Epoch, RunId, Seq, StoreError};
use crate::journal::{Append, Cancellation, Head, JournalStore, Lease, Record, WaitingRun};
const JOURNAL: TableDefinition<(&str, u64), &[u8]> = TableDefinition::new("journal");
const JOURNAL_BY_CASE: TableDefinition<(&str, &str, &str, u64), ()> =
TableDefinition::new("journal_by_case");
const EFFECT_ONCE: TableDefinition<(&str, &str), u64> = TableDefinition::new("effect_once");
const RUN_LEASE: TableDefinition<&str, (&str, u64, u64)> = TableDefinition::new("run_lease");
const RUN_SEAL: TableDefinition<&str, (&str, &[u8], u64)> = TableDefinition::new("run_seal");
const RUN_BY_OUTCOME: TableDefinition<(&str, &str, u64), &str> =
TableDefinition::new("run_by_outcome");
const RUN_OUTCOME: TableDefinition<(&str, &str), (&str, u64)> = TableDefinition::new("run_outcome");
const RUN_ADMISSION: TableDefinition<(&str, &str), (&str, i64)> =
TableDefinition::new("run_admission");
const RUN_ACTIVITY: TableDefinition<(&str, u64, &str), ()> = TableDefinition::new("run_activity");
const RUN_LAST_ACTIVITY: TableDefinition<(&str, &str), u64> =
TableDefinition::new("run_last_activity");
const RUN_SOURCE: TableDefinition<(&str, &str), &str> = TableDefinition::new("run_source");
const RUN_BY_SOURCE: TableDefinition<(&str, &str, u64, &str), ()> =
TableDefinition::new("run_by_source");
const RUN_WAITING: TableDefinition<(&str, i64, &str), &str> = TableDefinition::new("run_waiting");
const RUN_WAITING_AT: TableDefinition<(&str, &str), i64> = TableDefinition::new("run_waiting_at");
const SEAL_LOG: TableDefinition<(&str, u64), &str> = TableDefinition::new("seal_log");
const RUN_CANCEL: TableDefinition<&str, (&str, &str, &str, u64)> =
TableDefinition::new("run_cancel");
const COUNTERS: TableDefinition<&str, u64> = TableDefinition::new("counters");
const NEXT_LOG_INDEX: &str = "next_log_index";
const NEXT_CONCLUSION: &str = "next_conclusion";
pub(super) const MAX_STR: &str = "\u{10FFFF}";
#[derive(Debug, Clone)]
pub struct RedbStore {
db: Arc<Database>,
signer: Option<Arc<dyn crate::core::Signer>>,
origin: String,
tenant: crate::core::TenantId,
upcaster: Arc<dyn crate::journal::Upcaster>,
}
impl RedbStore {
pub fn open(path: impl AsRef<Path>) -> Result<Self, StoreError> {
let db = Database::create(path).map_err(|e| be(&e))?;
Self::init(db)
}
pub fn open_in_memory() -> Result<Self, StoreError> {
let db = Database::builder()
.create_with_backend(redb::backends::InMemoryBackend::new())
.map_err(|e| be(&e))?;
Self::init(db)
}
fn init(db: Database) -> Result<Self, StoreError> {
let w = begin_write(&db)?;
{
w.open_table(JOURNAL).map_err(|e| be(&e))?;
w.open_table(JOURNAL_BY_CASE).map_err(|e| be(&e))?;
w.open_table(EFFECT_ONCE).map_err(|e| be(&e))?;
w.open_table(RUN_LEASE).map_err(|e| be(&e))?;
w.open_table(RUN_SEAL).map_err(|e| be(&e))?;
w.open_table(RUN_BY_OUTCOME).map_err(|e| be(&e))?;
w.open_table(RUN_OUTCOME).map_err(|e| be(&e))?;
w.open_table(RUN_WAITING).map_err(|e| be(&e))?;
w.open_table(RUN_WAITING_AT).map_err(|e| be(&e))?;
w.open_table(SEAL_LOG).map_err(|e| be(&e))?;
w.open_table(RUN_CANCEL).map_err(|e| be(&e))?;
w.open_table(COUNTERS).map_err(|e| be(&e))?;
super::redb_cases::create_tables(&w)?;
super::redb_events::create_tables(&w)?;
super::redb_tasks::create_tables(&w)?;
super::redb_timers::create_tables(&w)?;
super::redb_batches::create_tables(&w)?;
}
w.commit().map_err(|e| be(&e))?;
Ok(Self {
db: Arc::new(db),
signer: None,
origin: "agentplane".to_owned(),
tenant: crate::core::TenantId::default(),
upcaster: crate::journal::current_upcaster(),
})
}
#[must_use]
pub fn for_tenant(mut self, tenant: crate::core::TenantId) -> Self {
self.tenant = tenant;
self
}
fn log_origin(&self) -> String {
if self.tenant.as_str() == crate::core::TenantId::DEFAULT {
self.origin.clone()
} else {
format!("{}/{}", self.origin, self.tenant)
}
}
pub(super) fn run_key(&self, run: RunId) -> String {
run_key_in(self.tenant.as_str(), &run.to_string())
}
pub(super) fn tenant_name(&self) -> String {
self.tenant.to_string()
}
pub(super) fn tenant_str(&self) -> &str {
self.tenant.as_str()
}
#[must_use]
pub fn origin(mut self, origin: impl Into<String>) -> Self {
self.origin = origin.into();
self
}
#[must_use]
pub fn signing_as(mut self, signer: Arc<dyn crate::core::Signer>) -> Self {
self.signer = Some(signer);
self
}
#[must_use]
pub fn upcasting_with(mut self, upcaster: Arc<dyn crate::journal::Upcaster>) -> Self {
self.upcaster = upcaster;
self
}
pub(super) async fn with_db<T, F>(&self, f: F) -> Result<T, StoreError>
where
T: Send + 'static,
F: FnOnce(&Database) -> Result<T, StoreError> + Send + 'static,
{
let db = Arc::clone(&self.db);
tokio::task::spawn_blocking(move || f(&db))
.await
.map_err(|e| StoreError::Backend(format!("blocking pool: {e}")))?
}
async fn log_leaves(&self) -> Result<Vec<crate::core::merkle::LeafHash>, StoreError> {
let tenant = self.tenant.clone();
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let log = r.open_table(SEAL_LOG).map_err(|e| be(&e))?;
let seals = r.open_table(RUN_SEAL).map_err(|e| be(&e))?;
let mut out = Vec::new();
for entry in log
.range((tenant.as_str(), 0)..=(tenant.as_str(), u64::MAX))
.map_err(|e| be(&e))?
{
let (_, run) = entry.map_err(|e| be(&e))?;
let Some(seal) = seals.get(run.value()).map_err(|e| be(&e))? else {
continue;
};
let (_, head, _) = seal.value();
out.push(crate::core::merkle::leaf_hash(&digest(head)?));
}
Ok(out)
})
.await
}
async fn positions(&self, runs: &[RunId]) -> Result<Vec<Option<(u64, Digest)>>, StoreError> {
let keys: Vec<String> = runs.iter().map(|&run| self.run_key(run)).collect();
let tenant = self.tenant.clone();
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let log = r.open_table(SEAL_LOG).map_err(|e| be(&e))?;
let seals = r.open_table(RUN_SEAL).map_err(|e| be(&e))?;
let mut wanted: std::collections::HashMap<&str, Option<(u64, Digest)>> =
keys.iter().map(|k| (k.as_str(), None)).collect();
let mut rank = 0u64;
for entry in log
.range((tenant.as_str(), 0)..=(tenant.as_str(), u64::MAX))
.map_err(|e| be(&e))?
{
let (_, run_id) = entry.map_err(|e| be(&e))?;
let Some(seal) = seals.get(run_id.value()).map_err(|e| be(&e))? else {
continue;
};
if let Some(slot) = wanted.get_mut(run_id.value()) {
let (_, head, _) = seal.value();
*slot = Some((rank, digest(head)?));
}
rank += 1;
}
Ok(keys.iter().map(|k| wanted[k.as_str()]).collect())
})
.await
}
#[doc(hidden)]
#[cfg(feature = "testkit")]
pub async fn tamper_for_test(
&self,
run: RunId,
seq: Seq,
body: Vec<u8>,
) -> Result<(), StoreError> {
let key = self.run_key(run);
self.with_db(move |db| {
let w = begin_write(db)?;
{
let mut t = w.open_table(JOURNAL).map_err(|e| be(&e))?;
let existing = t
.get((key.as_str(), seq))
.map_err(|e| be(&e))?
.map(|v| v.value().to_vec());
if let Some(bytes) = existing {
let row = Row::decode(&bytes)?;
let tampered = Row { body, ..row };
t.insert((key.as_str(), seq), tampered.encode()?.as_slice())
.map_err(|e| be(&e))?;
}
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
#[doc(hidden)]
#[cfg(feature = "testkit")]
pub async fn delete_run_for_test(&self, run: RunId) -> Result<(), StoreError> {
let key = self.run_key(run);
self.with_db(move |db| {
let w = begin_write(db)?;
{
let mut j = w.open_table(JOURNAL).map_err(|e| be(&e))?;
let seqs: Vec<u64> = j
.range((key.as_str(), 0u64)..=(key.as_str(), u64::MAX))
.map_err(|e| be(&e))?
.filter_map(Result::ok)
.map(|(k, _)| k.value().1)
.collect();
for s in seqs {
j.remove((key.as_str(), s)).map_err(|e| be(&e))?;
}
w.open_table(RUN_SEAL)
.map_err(|e| be(&e))?
.remove(key.as_str())
.map_err(|e| be(&e))?;
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
}
struct Row {
body: Vec<u8>,
prev_hash: [u8; 32],
hash: [u8; 32],
key_id: Option<String>,
signature: Option<Vec<u8>>,
}
impl Row {
fn encode(&self) -> Result<Vec<u8>, StoreError> {
let mut out = Vec::with_capacity(self.body.len() + 96);
out.extend_from_slice(&self.prev_hash);
out.extend_from_slice(&self.hash);
match (&self.key_id, &self.signature) {
(Some(k), Some(s)) => {
out.push(1);
push_bytes(&mut out, k.as_bytes())?;
push_bytes(&mut out, s)?;
}
_ => out.push(0),
}
push_bytes(&mut out, &self.body)?;
Ok(out)
}
fn decode(raw: &[u8]) -> Result<Self, StoreError> {
let corrupt = |what: &str| StoreError::Corrupt {
seq: 0,
detail: format!("journal row truncated in {what}"),
};
if raw.len() < 65 {
return Err(corrupt("header"));
}
let prev_hash: [u8; 32] = raw[0..32].try_into().map_err(|_| corrupt("prev_hash"))?;
let hash: [u8; 32] = raw[32..64].try_into().map_err(|_| corrupt("hash"))?;
let mut at = 64;
let attested = match raw[at] {
0 => false,
1 => true,
_ => return Err(corrupt("the signature flag")),
};
at += 1;
let (key_id, signature) = if attested {
let (k, n) = take_bytes(raw, at).ok_or_else(|| corrupt("key_id"))?;
at = n;
let (s, n) = take_bytes(raw, at).ok_or_else(|| corrupt("signature"))?;
at = n;
(
Some(String::from_utf8(k).map_err(|_| corrupt("key_id utf-8"))?),
Some(s),
)
} else {
(None, None)
};
let (body, end) = take_bytes(raw, at).ok_or_else(|| corrupt("body"))?;
if end != raw.len() {
return Err(StoreError::Corrupt {
seq: 0,
detail: "journal row carries bytes after its body".to_owned(),
});
}
Ok(Self {
body,
prev_hash,
hash,
key_id,
signature,
})
}
fn into_record(self, upcaster: &dyn crate::journal::Upcaster) -> Result<Record, StoreError> {
let signature = self
.key_id
.zip(self.signature)
.map(|(key_id, signature)| crate::core::KeySignature { key_id, signature });
Record::from_stored_with(
upcaster,
self.body,
Digest::from_bytes(self.prev_hash),
Digest::from_bytes(self.hash),
signature,
)
}
}
fn push_bytes(out: &mut Vec<u8>, bytes: &[u8]) -> Result<(), StoreError> {
let len = u32::try_from(bytes.len()).map_err(|_| {
StoreError::Backend(format!(
"a row field of {} bytes has no u32 length",
bytes.len()
))
})?;
out.extend_from_slice(&len.to_le_bytes());
out.extend_from_slice(bytes);
Ok(())
}
fn take_bytes(raw: &[u8], at: usize) -> Option<(Vec<u8>, usize)> {
let end = at.checked_add(4)?;
let len = u32::from_le_bytes(raw.get(at..end)?.try_into().ok()?) as usize;
let stop = end.checked_add(len)?;
Some((raw.get(end..stop)?.to_vec(), stop))
}
fn digest(bytes: &[u8]) -> Result<Digest, StoreError> {
let b: [u8; 32] = bytes.try_into().map_err(|_| StoreError::Corrupt {
seq: 0,
detail: "a stored hash is not 32 bytes".into(),
})?;
Ok(Digest::from_bytes(b))
}
pub(super) fn begin_write(db: &Database) -> Result<redb::WriteTransaction, StoreError> {
let mut w = db.begin_write().map_err(|e| be(&e))?;
w.set_durability(redb::Durability::Immediate)
.map_err(|e| be(&e))?;
Ok(w)
}
pub(super) fn run_key_in(tenant: &str, run: &str) -> String {
format!("{tenant}/{run}")
}
pub(super) fn is_sealed(
w: &redb::WriteTransaction,
tenant: &str,
run: &str,
) -> Result<bool, StoreError> {
let seals = w.open_table(RUN_SEAL).map_err(|e| be(&e))?;
Ok(seals
.get(run_key_in(tenant, run).as_str())
.map_err(|e| be(&e))?
.is_some())
}
pub(super) fn be<E: std::fmt::Display>(e: &E) -> StoreError {
StoreError::Backend(e.to_string())
}
pub(super) fn decoded<T>(what: &str, raw: &str, parsed: Option<T>) -> Result<T, StoreError> {
parsed.ok_or_else(|| StoreError::Corrupt {
seq: 0,
detail: format!("unknown {what} '{raw}'"),
})
}
#[allow(clippy::disallowed_methods)]
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| d.as_secs())
}
fn lease_expiry(now: u64, ttl: Duration) -> Result<u64, StoreError> {
let secs = ttl.as_secs();
if secs == 0 {
return Err(StoreError::Backend(format!(
"a lease TTL of {ttl:?} is below this store's whole-second \
granularity and would round to zero — pass at least one second, \
or use the runtime's lease_ttl which enforces its own minimum"
)));
}
now.checked_add(secs).ok_or_else(|| {
StoreError::Backend(format!(
"a lease TTL of {ttl:?} overflows the expiry instant — the lease \
would wrap into the past and read as already expired"
))
})
}
fn refuse_unconcluded(run: &str, last: Option<&[u8]>) -> Result<(), StoreError> {
let kind = last
.and_then(|body| serde_json::from_slice::<serde_json::Value>(body).ok())
.and_then(|body| body.get("kind").and_then(|k| k.as_str().map(str::to_owned)));
if kind.as_deref() == Some("RunConcluded") {
return Ok(());
}
Err(StoreError::Backend(format!(
"run {run} cannot be sealed: its last record is {} rather than its conclusion",
kind.as_deref().unwrap_or("absent")
)))
}
fn head_of(
t: &impl ReadableTable<(&'static str, u64), &'static [u8]>,
run: &str,
) -> Result<Head, StoreError>
where
{
let last = t
.range((run, 0u64)..=(run, u64::MAX))
.map_err(|e| be(&e))?
.next_back();
match last {
None => Ok(Head::genesis()),
Some(entry) => {
let (k, v) = entry.map_err(|e| be(&e))?;
let row = Row::decode(v.value())?;
Ok(Head {
seq: k.value().1,
hash: Digest::from_bytes(row.hash),
})
}
}
}
fn epoch_after_history(
w: &redb::WriteTransaction,
run: &str,
upcaster: &dyn crate::journal::Upcaster,
) -> Result<Epoch, StoreError> {
let t = w.open_table(JOURNAL).map_err(|e| be(&e))?;
let last = t
.range((run, 0u64)..=(run, u64::MAX))
.map_err(|e| be(&e))?
.next_back();
match last {
None => Ok(1),
Some(entry) => {
let (_, v) = entry.map_err(|e| be(&e))?;
Ok(Row::decode(v.value())?.into_record(upcaster)?.body.epoch + 1)
}
}
}
#[allow(clippy::too_many_lines)]
#[async_trait]
impl JournalStore for RedbStore {
fn is_shared(&self) -> bool {
false
}
fn seals(&self) -> bool {
false
}
fn tenant(&self) -> &str {
self.tenant.as_str()
}
async fn append(&self, epoch: Epoch, batch: Vec<Append>) -> Result<Vec<Record>, StoreError> {
if batch.is_empty() {
return Ok(Vec::new());
}
let run = batch[0].run;
if let Some(a) = batch.iter().find(|a| a.run != run) {
return Err(StoreError::Backend(format!(
"batch spans runs {run} and {} — a batch is one run's atomic unit",
a.run
)));
}
let signer = self.signer.clone();
let key = self.run_key(run);
let run_id = run.to_string();
let tenant = self.tenant_name();
let upcaster = Arc::clone(&self.upcaster);
self.with_db(move |db| {
let w = begin_write(db)?;
let sealed = {
let mut journal = w.open_table(JOURNAL).map_err(|e| be(&e))?;
let mut by_case = w.open_table(JOURNAL_BY_CASE).map_err(|e| be(&e))?;
let mut once = w.open_table(EFFECT_ONCE).map_err(|e| be(&e))?;
{
let leases = w.open_table(RUN_LEASE).map_err(|e| be(&e))?;
if let Some(l) = leases.get(key.as_str()).map_err(|e| be(&e))? {
let (_, current, _) = l.value();
if epoch != current {
return Err(StoreError::Fenced {
run: key.clone(),
held: epoch,
current,
});
}
}
}
{
let seals = w.open_table(RUN_SEAL).map_err(|e| be(&e))?;
if let Some(seal) = seals.get(key.as_str()).map_err(|e| be(&e))? {
let (outcome, _, _) = seal.value();
return Err(StoreError::RunSealed {
run: key.clone(),
outcome: outcome.to_owned(),
});
}
}
let mut head = head_of(&journal, &key)?;
let mut sealed = Vec::with_capacity(batch.len());
let mut waiting: Option<crate::core::SuspendReason> = None;
let mut admitted_from: Option<String> = None;
for append in batch {
let (body, written) =
append.into_parts(head.seq + 1, epoch, upcaster.as_ref())?;
let effect_key = body.effect_key.map(EffectKey::to_hex);
let is_start =
matches!(body.kind, crate::journal::RecordKind::EffectStarted { .. });
let (conclusion, claimed) = match &body.kind {
crate::journal::RecordKind::RunConcluded { outcome, .. } => {
(Some(outcome.clone()), None)
}
crate::journal::RecordKind::RunAdmitted {
idempotency_key: Some(k),
..
} => (None, Some(k.clone())),
_ => (None, None),
};
waiting = match &body.kind {
crate::journal::RecordKind::RunSuspended { reason } => Some(reason.clone()),
_ => None,
};
let record = Record::seal_at(body, written, head.hash, signer.as_deref())?;
let seq = record.seq();
if is_start && let Some(ek) = &effect_key {
let prior = once
.insert((key.as_str(), ek.as_str()), seq)
.map_err(|e| be(&e))?;
if prior.is_some() {
return Err(record.effect_key().map_or_else(
|| StoreError::Backend("duplicate effect".into()),
StoreError::DuplicateEffect,
));
}
}
let row = Row {
body: record.raw().to_vec(),
prev_hash: *record.prev_hash.as_bytes(),
hash: *record.hash.as_bytes(),
key_id: record.signature.as_ref().map(|a| a.key_id.clone()),
signature: record.signature.as_ref().map(|a| a.signature.clone()),
};
journal
.insert((key.as_str(), seq), row.encode()?.as_slice())
.map_err(|e| be(&e))?;
if let Some(case) = record.body.case {
by_case
.insert(
(
tenant.as_str(),
case.to_string().as_str(),
key.as_str(),
seq,
),
(),
)
.map_err(|e| be(&e))?;
}
if let Some(source) = record.admission_source() {
admitted_from = Some(source.to_owned());
}
if let Some(k) = claimed {
let mut admissions = w.open_table(RUN_ADMISSION).map_err(|e| be(&e))?;
if let Some(held) = admissions
.get((tenant.as_str(), k.as_str()))
.map_err(|e| be(&e))?
.map(|v| v.value().0.to_owned())
{
return Err(StoreError::DuplicateAdmission { key: k, run: held });
}
admissions
.insert(
(tenant.as_str(), k.as_str()),
(run_id.as_str(), now_secs().cast_signed()),
)
.map_err(|e| be(&e))?;
}
if let Some(outcome) = conclusion {
let mut by_outcome = w.open_table(RUN_BY_OUTCOME).map_err(|e| be(&e))?;
let mut outcomes = w.open_table(RUN_OUTCOME).map_err(|e| be(&e))?;
let mut counters = w.open_table(COUNTERS).map_err(|e| be(&e))?;
if let Some(prior) = outcomes
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (o, ord) = v.value();
(o.to_owned(), ord)
})
{
by_outcome
.remove((tenant.as_str(), prior.0.as_str(), prior.1))
.map_err(|e| be(&e))?;
}
let counter = format!("{NEXT_CONCLUSION}/{tenant}");
let next = counters
.get(counter.as_str())
.map_err(|e| be(&e))?
.map_or(0, |v| v.value());
by_outcome
.insert((tenant.as_str(), outcome.as_str(), next), key.as_str())
.map_err(|e| be(&e))?;
outcomes
.insert((tenant.as_str(), key.as_str()), (outcome.as_str(), next))
.map_err(|e| be(&e))?;
counters
.insert(counter.as_str(), next + 1)
.map_err(|e| be(&e))?;
}
head = Head {
seq,
hash: record.hash,
};
sealed.push(record);
}
{
let mut index = w.open_table(RUN_WAITING).map_err(|e| be(&e))?;
let mut at = w.open_table(RUN_WAITING_AT).map_err(|e| be(&e))?;
if let Some(prior) = at
.remove((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|v| v.value())
{
index
.remove((tenant.as_str(), prior, key.as_str()))
.map_err(|e| be(&e))?;
}
if let Some(reason) = &waiting {
let until = reason.until().unix_timestamp();
let encoded = String::from_utf8(
crate::core::canon::to_bytes(reason)
.map_err(|e| StoreError::Backend(e.to_string()))?,
)
.map_err(|e| StoreError::Backend(e.to_string()))?;
index
.insert((tenant.as_str(), until, key.as_str()), encoded.as_str())
.map_err(|e| be(&e))?;
at.insert((tenant.as_str(), key.as_str()), until)
.map_err(|e| be(&e))?;
}
}
let updated = now_secs();
let mut activity = w.open_table(RUN_ACTIVITY).map_err(|e| be(&e))?;
let mut last = w.open_table(RUN_LAST_ACTIVITY).map_err(|e| be(&e))?;
let previous = last
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|value| value.value());
if let Some(previous) = previous {
activity
.remove((tenant.as_str(), previous, key.as_str()))
.map_err(|e| be(&e))?;
}
activity
.insert((tenant.as_str(), updated, key.as_str()), ())
.map_err(|e| be(&e))?;
last.insert((tenant.as_str(), key.as_str()), updated)
.map_err(|e| be(&e))?;
let mut sources = w.open_table(RUN_SOURCE).map_err(|e| be(&e))?;
if let Some(source) = &admitted_from {
sources
.insert((tenant.as_str(), key.as_str()), source.as_str())
.map_err(|e| be(&e))?;
}
let source = sources
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|value| value.value().to_owned());
if let Some(source) = source {
let mut by_source = w.open_table(RUN_BY_SOURCE).map_err(|e| be(&e))?;
if let Some(previous) = previous {
by_source
.remove((tenant.as_str(), source.as_str(), previous, key.as_str()))
.map_err(|e| be(&e))?;
}
by_source
.insert(
(tenant.as_str(), source.as_str(), updated, key.as_str()),
(),
)
.map_err(|e| be(&e))?;
}
sealed
};
w.commit().map_err(|e| be(&e))?;
Ok(sealed)
})
.await
}
async fn read(&self, run: RunId, from: Seq) -> Result<Vec<Record>, StoreError> {
self.read_page(run, from, usize::MAX).await
}
async fn read_page(
&self,
run: RunId,
from: Seq,
limit: usize,
) -> Result<Vec<Record>, StoreError> {
let key = self.run_key(run);
let upcaster = Arc::clone(&self.upcaster);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let t = r.open_table(JOURNAL).map_err(|e| be(&e))?;
let mut out = Vec::new();
for entry in t
.range((key.as_str(), from)..=(key.as_str(), u64::MAX))
.map_err(|e| be(&e))?
.take(limit)
{
let (_, v) = entry.map_err(|e| be(&e))?;
out.push(Row::decode(v.value())?.into_record(upcaster.as_ref())?);
}
Ok(out)
})
.await
}
async fn case_history(
&self,
case: crate::core::CaseId,
limit: usize,
) -> Result<Vec<Record>, StoreError> {
let tenant = self.tenant_name();
let case = case.to_string();
let upcaster = Arc::clone(&self.upcaster);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let idx = r.open_table(JOURNAL_BY_CASE).map_err(|e| be(&e))?;
let j = r.open_table(JOURNAL).map_err(|e| be(&e))?;
let mut out = Vec::new();
for entry in idx
.range(
(tenant.as_str(), case.as_str(), "", 0)
..=(tenant.as_str(), case.as_str(), MAX_STR, u64::MAX),
)
.map_err(|e| be(&e))?
{
if out.len() >= limit {
break;
}
let (k, _) = entry.map_err(|e| be(&e))?;
let (_, _, run, seq) = k.value();
if let Some(v) = j.get((run, seq)).map_err(|e| be(&e))? {
out.push(Row::decode(v.value())?.into_record(upcaster.as_ref())?);
}
}
Ok(out)
})
.await
}
async fn head(&self, run: RunId) -> Result<Head, StoreError> {
let key = self.run_key(run);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let t = r.open_table(JOURNAL).map_err(|e| be(&e))?;
head_of(&t, &key)
})
.await
}
async fn acquire(&self, run: RunId, owner: &str, ttl: Duration) -> Result<Lease, StoreError> {
let key = self.run_key(run);
let owner = owner.to_owned();
let owner_out = owner.clone();
let upcaster = Arc::clone(&self.upcaster);
let epoch = self
.with_db(move |db| {
let w = begin_write(db)?;
let epoch = {
let mut leases = w.open_table(RUN_LEASE).map_err(|e| be(&e))?;
let now = now_secs();
let expires = lease_expiry(now, ttl)?;
let existing = leases.get(key.as_str()).map_err(|e| be(&e))?.map(|v| {
let (o, e, x) = v.value();
(o.to_owned(), e, x)
});
let epoch = match existing {
None => epoch_after_history(&w, key.as_str(), upcaster.as_ref())?,
Some((_, epoch, expires_at)) if expires_at <= now => epoch + 1,
Some((held_by, epoch, expires_at)) => {
return Err(StoreError::LeaseHeld {
run: key.clone(),
owner: held_by,
epoch,
remaining_secs: expires_at.saturating_sub(now),
});
}
};
leases
.insert(key.as_str(), (owner.as_str(), epoch, expires))
.map_err(|e| be(&e))?;
epoch
};
w.commit().map_err(|e| be(&e))?;
Ok(epoch)
})
.await?;
Ok(Lease {
run,
owner: owner_out,
epoch,
})
}
async fn renew(
&self,
run: RunId,
owner: &str,
epoch: Epoch,
ttl: Duration,
) -> Result<Lease, StoreError> {
let key = self.run_key(run);
let owner = owner.to_owned();
let owner_out = owner.clone();
self.with_db(move |db| {
let w = begin_write(db)?;
{
let mut leases = w.open_table(RUN_LEASE).map_err(|e| be(&e))?;
let now = now_secs();
let expires = lease_expiry(now, ttl)?;
let existing = leases.get(key.as_str()).map_err(|e| be(&e))?.map(|v| {
let (o, e, x) = v.value();
(o.to_owned(), e, x)
});
match existing {
Some((held_by, held_epoch, expires_at))
if held_by == owner && held_epoch == epoch && expires_at > now =>
{
leases
.insert(key.as_str(), (owner.as_str(), epoch, expires))
.map_err(|e| be(&e))?;
}
_ => {
return Err(StoreError::LeaseNotHeld {
run: key.clone(),
epoch,
});
}
}
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await?;
Ok(Lease {
run,
owner: owner_out,
epoch,
})
}
async fn release_lease(&self, run: RunId, epoch: Epoch) -> Result<(), StoreError> {
let key = self.run_key(run);
self.with_db(move |db| {
let w = begin_write(db)?;
{
let mut leases = w.open_table(RUN_LEASE).map_err(|e| be(&e))?;
let held = leases
.get(key.as_str())
.map_err(|e| be(&e))?
.map(|v| v.value().1);
if held == Some(epoch) {
leases
.insert(key.as_str(), ("", epoch, 0))
.map_err(|e| be(&e))?;
}
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
}
async fn waiting_runs(&self, limit: usize) -> Result<Vec<WaitingRun>, StoreError> {
let tenant = self.tenant.to_string();
let prefix = format!("{}/", self.tenant);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let index = r.open_table(RUN_WAITING).map_err(|e| be(&e))?;
let mut out = Vec::new();
for entry in index
.range((tenant.as_str(), i64::MIN, "")..=(tenant.as_str(), i64::MAX, "\u{10FFFF}"))
.map_err(|e| be(&e))?
.take(limit)
{
let (k, v) = entry.map_err(|e| be(&e))?;
let (_, _, key) = k.value();
let id = key.strip_prefix(prefix.as_str()).unwrap_or(key);
let run = RunId::parse(id).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_waiting holds an unparsable run id '{id}': {e}"),
})?;
let reason = serde_json::from_str(v.value()).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_waiting holds an unreadable reason for {id}: {e}"),
})?;
out.push(WaitingRun { run, reason });
}
Ok(out)
})
.await
}
async fn abandoned_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError> {
let prefix = format!("{}/", self.tenant);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let leases = r.open_table(RUN_LEASE).map_err(|e| be(&e))?;
let now = now_secs();
let mut expired: Vec<(u64, RunId)> = Vec::new();
for entry in leases.range(prefix.as_str()..).map_err(|e| be(&e))? {
let (k, v) = entry.map_err(|e| be(&e))?;
let key = k.value();
if !key.starts_with(prefix.as_str()) {
break;
}
let (owner, _, expires_at) = v.value();
if owner.is_empty() || expires_at > now {
continue;
}
let id = &key[prefix.len()..];
let run = RunId::parse(id).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_lease holds an unparsable run id '{id}': {e}"),
})?;
expired.push((expires_at, run));
}
expired.sort_unstable_by_key(|(at, _)| *at);
Ok(expired.into_iter().take(limit).map(|(_, r)| r).collect())
})
.await
}
async fn seal(&self, run: RunId, epoch: Epoch, outcome: &str) -> Result<Digest, StoreError> {
let key = self.run_key(run);
let tenant = self.tenant.clone();
let outcome = outcome.to_owned();
self.with_db(move |db| {
let w = begin_write(db)?;
let head_hash = {
{
let leases = w.open_table(RUN_LEASE).map_err(|e| be(&e))?;
if let Some(l) = leases.get(key.as_str()).map_err(|e| be(&e))? {
let (_, current, _) = l.value();
if epoch != current {
return Err(StoreError::Fenced {
run: key.clone(),
held: epoch,
current,
});
}
}
}
let head = {
let journal = w.open_table(JOURNAL).map_err(|e| be(&e))?;
let last = journal
.range((key.as_str(), 0u64)..=(key.as_str(), u64::MAX))
.map_err(|e| be(&e))?
.next_back()
.transpose()
.map_err(|e| be(&e))?
.map(|(_, row)| Row::decode(row.value()))
.transpose()?;
refuse_unconcluded(&key, last.as_ref().map(|row| row.body.as_slice()))?;
head_of(&journal, &key)?
};
let mut seals = w.open_table(RUN_SEAL).map_err(|e| be(&e))?;
if seals.get(key.as_str()).map_err(|e| be(&e))?.is_none() {
seals
.insert(
key.as_str(),
(
outcome.as_str(),
head.hash.as_bytes().as_slice(),
now_secs(),
),
)
.map_err(|e| be(&e))?;
let counter = format!("{NEXT_LOG_INDEX}/{tenant}");
let mut counters = w.open_table(COUNTERS).map_err(|e| be(&e))?;
let next = counters
.get(counter.as_str())
.map_err(|e| be(&e))?
.map_or(0, |v| v.value());
w.open_table(SEAL_LOG)
.map_err(|e| be(&e))?
.insert((tenant.as_str(), next), key.as_str())
.map_err(|e| be(&e))?;
counters
.insert(counter.as_str(), next + 1)
.map_err(|e| be(&e))?;
}
head.hash
};
w.commit().map_err(|e| be(&e))?;
Ok(head_hash)
})
.await
}
async fn admitted_as(&self, key: &str) -> Result<Option<RunId>, StoreError> {
let tenant = self.tenant_name();
let key = key.to_owned();
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let admissions = match r.open_table(RUN_ADMISSION) {
Ok(t) => t,
Err(redb::TableError::TableDoesNotExist(_)) => return Ok(None),
Err(e) => return Err(be(&e)),
};
let Some(held) = admissions
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|v| v.value().0.to_owned())
else {
return Ok(None);
};
RunId::parse(&held)
.map(Some)
.map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_admission holds an unparseable run id '{held}': {e}"),
})
})
.await
}
async fn forget_admissions(
&self,
older_than: crate::core::Timestamp,
) -> Result<usize, StoreError> {
let tenant = self.tenant_name();
let cutoff = older_than.unix_timestamp();
self.with_db(move |db| {
let w = begin_write(db)?;
let removed = {
let mut admissions = match w.open_table(RUN_ADMISSION) {
Ok(t) => t,
Err(redb::TableError::TableDoesNotExist(_)) => return Ok(0),
Err(e) => return Err(be(&e)),
};
let stale: Vec<String> = admissions
.range((tenant.as_str(), "")..=(tenant.as_str(), MAX_STR))
.map_err(|e| be(&e))?
.filter_map(|entry| {
let (k, v) = entry.ok()?;
(v.value().1 < cutoff).then(|| k.value().1.to_owned())
})
.collect();
for key in &stale {
admissions
.remove((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?;
}
stale.len()
};
w.commit().map_err(|e| be(&e))?;
Ok(removed)
})
.await
}
async fn runs_by_outcome(&self, outcome: &str, limit: usize) -> Result<Vec<RunId>, StoreError> {
let tenant = self.tenant_name();
let outcome = outcome.to_owned();
let prefix = format!("{tenant}/");
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let idx = r.open_table(RUN_BY_OUTCOME).map_err(|e| be(&e))?;
let mut out = Vec::new();
for entry in idx
.range(
(tenant.as_str(), outcome.as_str(), 0)
..=(tenant.as_str(), outcome.as_str(), u64::MAX),
)
.map_err(|e| be(&e))?
.rev()
{
if out.len() >= limit {
break;
}
let (_, v) = entry.map_err(|e| be(&e))?;
let key = v.value();
let id = key
.strip_prefix(prefix.as_str())
.ok_or_else(|| StoreError::Corrupt {
seq: 0,
detail: format!(
"run_by_outcome points at '{key}', which is outside tenant '{tenant}'"
),
})?;
let run = RunId::parse(id).map_err(|e| StoreError::Corrupt {
seq: 0,
detail: format!("run_by_outcome holds an unparsable run id '{id}': {e}"),
})?;
out.push(run);
}
Ok(out)
})
.await
}
async fn count_by_outcome(&self, outcome: &str) -> Result<u64, StoreError> {
let tenant = self.tenant_name();
let outcome = outcome.to_owned();
let prefix = format!("{tenant}/");
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let idx = r.open_table(RUN_BY_OUTCOME).map_err(|e| be(&e))?;
let mut n: u64 = 0;
for entry in idx
.range(
(tenant.as_str(), outcome.as_str(), 0)
..=(tenant.as_str(), outcome.as_str(), u64::MAX),
)
.map_err(|e| be(&e))?
{
let (_, v) = entry.map_err(|e| be(&e))?;
let key = v.value();
if key.strip_prefix(prefix.as_str()).is_none() {
return Err(StoreError::Corrupt {
seq: 0,
detail: format!(
"run_by_outcome points at '{key}', which is outside tenant '{tenant}'"
),
});
}
n += 1;
}
Ok(n)
})
.await
}
async fn recent_runs(
&self,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError> {
let tenant = self.tenant_name();
let prefix = format!("{tenant}/");
let end = after.map(|(updated, run)| (updated, format!("{prefix}{run}")));
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let Ok(activity) = r.open_table(RUN_ACTIVITY) else {
return Ok(Vec::new());
};
let rows = match &end {
Some((updated, key)) => activity
.range((tenant.as_str(), 0, "")..(tenant.as_str(), *updated, key.as_str()))
.map_err(|e| be(&e))?,
None => activity
.range((tenant.as_str(), 0, "")..=(tenant.as_str(), u64::MAX, MAX_STR))
.map_err(|e| be(&e))?,
};
rows.rev()
.filter_map(|entry| match entry {
Ok((key, _)) => {
let (_, updated, stored) = key.value();
stored
.strip_prefix(prefix.as_str())
.and_then(|id| RunId::parse(id).ok())
.map(|run| Ok((run, updated)))
}
Err(error) => Some(Err(be(&error))),
})
.take(limit)
.collect()
})
.await
}
async fn runs_by_id(
&self,
after: Option<RunId>,
limit: usize,
) -> Result<Vec<RunId>, StoreError> {
use std::ops::Bound;
let tenant = self.tenant_name();
let prefix = format!("{tenant}/");
let start = after.map(|run| format!("{prefix}{run}"));
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let Ok(last) = r.open_table(RUN_LAST_ACTIVITY) else {
return Ok(Vec::new());
};
let from = match &start {
Some(key) => Bound::Excluded((tenant.as_str(), key.as_str())),
None => Bound::Included((tenant.as_str(), "")),
};
last.range((from, Bound::Included((tenant.as_str(), MAX_STR))))
.map_err(|e| be(&e))?
.filter_map(|entry| match entry {
Ok((key, _)) => {
let (_, stored) = key.value();
stored
.strip_prefix(prefix.as_str())
.and_then(|id| RunId::parse(id).ok())
.map(Ok)
}
Err(error) => Some(Err(be(&error))),
})
.take(limit)
.collect()
})
.await
}
async fn recent_runs_from(
&self,
source: &str,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError> {
let tenant = self.tenant_name();
let source = source.to_owned();
let prefix = format!("{tenant}/");
let end = after.map(|(updated, run)| (updated, format!("{prefix}{run}")));
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let by_source = match r.open_table(RUN_BY_SOURCE) {
Ok(t) => t,
Err(redb::TableError::TableDoesNotExist(_)) => return Ok(Vec::new()),
Err(e) => return Err(be(&e)),
};
let (t, s) = (tenant.as_str(), source.as_str());
let rows = match &end {
Some((updated, key)) => by_source
.range((t, s, 0, "")..(t, s, *updated, key.as_str()))
.map_err(|e| be(&e))?,
None => by_source
.range((t, s, 0, "")..=(t, s, u64::MAX, MAX_STR))
.map_err(|e| be(&e))?,
};
rows.rev()
.filter_map(|entry| match entry {
Ok((key, _)) => {
let (_, _, updated, stored) = key.value();
stored
.strip_prefix(prefix.as_str())
.and_then(|id| RunId::parse(id).ok())
.map(|run| Ok((run, updated)))
}
Err(error) => Some(Err(be(&error))),
})
.take(limit)
.collect()
})
.await
}
async fn checkpoint(&self) -> Result<crate::journal::Checkpoint, StoreError> {
let leaves = self.log_leaves().await?;
Ok(crate::journal::Checkpoint {
origin: self.log_origin(),
size: leaves.len() as u64,
root: crate::core::merkle::root(&leaves),
})
}
async fn consistency_proof(&self, old_size: u64) -> Result<Vec<Digest>, StoreError> {
let leaves = self.log_leaves().await?;
let old = usize::try_from(old_size).unwrap_or(usize::MAX);
if old > leaves.len() {
return Err(StoreError::Backend(format!(
"a checkpoint of size {old_size} is larger than this log ({}) — \
either it belongs to another plane, or runs were removed",
leaves.len()
)));
}
Ok(crate::core::merkle::consistency_proof(&leaves, old))
}
async fn inclusion_proof(
&self,
run: RunId,
) -> Result<Option<crate::journal::Inclusion>, StoreError> {
let leaves = self.log_leaves().await?;
let Some((index, seal)) = self.positions(&[run]).await?.pop().flatten() else {
return Ok(None);
};
Ok(Some(crate::journal::Inclusion {
index,
size: leaves.len() as u64,
seal,
proof: crate::core::merkle::inclusion_proof(
&leaves,
usize::try_from(index).unwrap_or(usize::MAX),
),
}))
}
async fn inclusion_proof_at(
&self,
run: RunId,
size: u64,
) -> Result<Option<crate::journal::Inclusion>, StoreError> {
let leaves = self.log_leaves().await?;
let prefix = usize::try_from(size)
.ok()
.and_then(|size| leaves.get(..size))
.ok_or_else(|| {
StoreError::Backend(format!(
"asked to prove at size {size} and the log holds {} leaves",
leaves.len()
))
})?;
let Some((index, seal)) = self.positions(&[run]).await?.pop().flatten() else {
return Ok(None);
};
let Some(at) = usize::try_from(index).ok().filter(|&at| at < prefix.len()) else {
return Ok(None);
};
Ok(Some(crate::journal::Inclusion {
index,
size,
seal,
proof: crate::core::merkle::inclusion_proof(prefix, at),
}))
}
async fn log_positions(
&self,
runs: &[RunId],
) -> Result<Vec<Option<(u64, Digest)>>, StoreError> {
self.positions(runs).await
}
async fn request_cancel(
&self,
run: RunId,
actor: &crate::core::Operator,
reason: &str,
) -> Result<bool, StoreError> {
let key = self.run_key(run);
let (who, basis, reason) = (
actor.actor().to_owned(),
actor.basis().as_str(),
reason.to_owned(),
);
self.with_db(move |db| {
let w = begin_write(db)?;
let first = {
let mut t = w.open_table(RUN_CANCEL).map_err(|e| be(&e))?;
if t.get(key.as_str()).map_err(|e| be(&e))?.is_some() {
false
} else {
t.insert(
key.as_str(),
(who.as_str(), basis, reason.as_str(), now_secs()),
)
.map_err(|e| be(&e))?;
true
}
};
w.commit().map_err(|e| be(&e))?;
Ok(first)
})
.await
}
async fn cancellation(&self, run: RunId) -> Result<Option<Cancellation>, StoreError> {
let key = self.run_key(run);
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let t = r.open_table(RUN_CANCEL).map_err(|e| be(&e))?;
t.get(key.as_str())
.map_err(|e| be(&e))?
.map(|v| {
let (who, basis, reason, _) = v.value();
Ok(Cancellation {
actor: super::decode_operator(who, basis, "run_cancel")?,
reason: reason.to_owned(),
})
})
.transpose()
})
.await
}
}
#[cfg(test)]
mod tests {
use super::*;
fn row(signed: bool) -> Vec<u8> {
Row {
body: b"{}".to_vec(),
prev_hash: [1; 32],
hash: [2; 32],
key_id: signed.then(|| "k".to_owned()),
signature: signed.then(|| vec![3; 64]),
}
.encode()
.expect("encodes")
}
#[test]
fn a_signature_flag_other_than_zero_or_one_is_corrupt() {
let mut row = row(false);
assert!(Row::decode(&row).is_ok_and(|r| r.signature.is_none()));
row[64] = 2;
assert!(
matches!(Row::decode(&row), Err(StoreError::Corrupt { .. })),
"a flag byte of 2 decoded as an unsigned row"
);
}
#[test]
fn trailing_bytes_after_a_journal_row_are_corrupt() {
let mut row = row(true);
assert!(Row::decode(&row).is_ok_and(|r| r.signature.is_some()));
row.push(0);
assert!(
matches!(Row::decode(&row), Err(StoreError::Corrupt { .. })),
"a row with bytes after its body decoded as whole"
);
}
#[tokio::test]
async fn checkpoint_origin_is_order_insensitive_and_idempotent() {
let tenant = crate::core::TenantId::new("acme").expect("tenant");
let tenant_then_origin = RedbStore::open_in_memory()
.expect("store")
.for_tenant(tenant.clone())
.origin("plane-1");
let origin_then_tenant = RedbStore::open_in_memory()
.expect("store")
.origin("plane-1")
.for_tenant(tenant.clone());
let twice = RedbStore::open_in_memory()
.expect("store")
.origin("plane-1")
.for_tenant(tenant.clone())
.for_tenant(tenant);
let a = tenant_then_origin.checkpoint().await.expect("checkpoint");
let b = origin_then_tenant.checkpoint().await.expect("checkpoint");
let c = twice.checkpoint().await.expect("checkpoint");
assert_eq!(a.origin, "plane-1/acme");
assert_eq!(b.origin, a.origin, "builder order changed the log's name");
assert_eq!(c.origin, a.origin, "for_tenant is not idempotent");
}
#[tokio::test]
async fn an_unparsable_run_id_is_reported_as_corruption_not_skipped() {
let store = RedbStore::open_in_memory().expect("store");
let stranded = RunId::generate();
let good_key = store.run_key(stranded);
store
.with_db({
let good_key = good_key.clone();
move |db| {
let w = begin_write(db)?;
{
w.open_table(RUN_LEASE)
.map_err(|e| be(&e))?
.insert(good_key.as_str(), ("worker", 1u64, 0u64))
.map_err(|e| be(&e))?;
w.open_table(RUN_BY_OUTCOME)
.map_err(|e| be(&e))?
.insert(("default", "failed", 0u64), good_key.as_str())
.map_err(|e| be(&e))?;
}
w.commit().map_err(|e| be(&e))?;
Ok(())
}
})
.await
.expect("plant the healthy rows");
assert_eq!(
store.abandoned_runs(10).await.expect("a clean scan"),
vec![stranded]
);
assert_eq!(
store
.runs_by_outcome("failed", 10)
.await
.expect("a clean scan"),
vec![stranded]
);
store
.with_db(|db| {
let w = begin_write(db)?;
{
w.open_table(RUN_LEASE)
.map_err(|e| be(&e))?
.insert("default/not-a-run-id", ("worker", 1u64, 0u64))
.map_err(|e| be(&e))?;
w.open_table(RUN_BY_OUTCOME)
.map_err(|e| be(&e))?
.insert(("default", "failed", 1u64), "default/not-a-run-id")
.map_err(|e| be(&e))?;
}
w.commit().map_err(|e| be(&e))?;
Ok(())
})
.await
.expect("plant the corrupt rows");
let err = store
.abandoned_runs(10)
.await
.expect_err("an unparsable lease row must refuse the sweep, not vanish from it");
assert!(
matches!(err, StoreError::Corrupt { .. }),
"the refusal must be Corrupt, so it is promoted past retry logic: {err}"
);
let err = store
.runs_by_outcome("failed", 10)
.await
.expect_err("an unparsable outcome row must refuse the listing, not vanish from it");
assert!(
matches!(err, StoreError::Corrupt { .. }),
"the refusal must be Corrupt, so it is promoted past retry logic: {err}"
);
}
}