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};
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_ACTIVITY: TableDefinition<(&str, u64, &str), ()> = TableDefinition::new("run_activity");
const RUN_LAST_ACTIVITY: TableDefinition<(&str, &str), u64> =
TableDefinition::new("run_last_activity");
const SEAL_LOG: TableDefinition<(&str, u64), &str> = TableDefinition::new("seal_log");
const RUN_CANCEL: TableDefinition<&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,
}
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> {
use std::sync::atomic::{AtomicU64, Ordering};
static N: AtomicU64 = AtomicU64::new(0);
let n = N.fetch_add(1, Ordering::Relaxed);
let pid = std::process::id();
let path = std::env::temp_dir().join(format!("agentplane-{pid}-{n}.redb"));
let _ = std::fs::remove_file(&path);
let db = Database::create(&path).map_err(|e| be(&e))?;
let _ = std::fs::remove_file(&path);
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(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(),
})
}
#[must_use]
pub fn for_tenant(mut self, tenant: crate::core::TenantId) -> Self {
self.origin = format!("{}/{}", self.origin, tenant);
self.tenant = tenant;
self
}
pub(super) fn run_key(&self, run: RunId) -> String {
format!("{}/{run}", self.tenant)
}
pub(super) fn tenant_name(&self) -> String {
self.tenant.to_string()
}
#[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
}
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<Digest>, 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 log_position(&self, run: RunId) -> Result<Option<(usize, Digest)>, StoreError> {
let key = self.run_key(run);
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 rank = 0usize;
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 run_id.value() == key {
let (_, head, _) = seal.value();
return Ok(Some((rank, digest(head)?)));
}
rank += 1;
}
Ok(None)
})
.await
}
#[doc(hidden)]
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)]
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) -> Vec<u8> {
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);
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 = raw[at] == 1;
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, _) = take_bytes(raw, at).ok_or_else(|| corrupt("body"))?;
Ok(Self {
body,
prev_hash,
hash,
key_id,
signature,
})
}
fn into_record(self) -> Result<Record, StoreError> {
let attestation = self
.key_id
.zip(self.signature)
.map(|(key_id, signature)| crate::core::Attestation { key_id, signature });
Record::from_stored_attested(
self.body,
Digest::from_bytes(self.prev_hash),
Digest::from_bytes(self.hash),
attestation,
)
}
}
fn push_bytes(out: &mut Vec<u8>, bytes: &[u8]) {
let len = u32::try_from(bytes.len()).unwrap_or(u32::MAX);
out.extend_from_slice(&len.to_le_bytes());
out.extend_from_slice(bytes);
}
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 be<E: std::fmt::Display>(e: &E) -> StoreError {
StoreError::Backend(e.to_string())
}
#[allow(clippy::disallowed_methods)]
fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| d.as_secs())
}
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),
})
}
}
}
#[allow(clippy::too_many_lines)]
#[async_trait]
impl JournalStore for RedbStore {
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 tenant = self.tenant_name();
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());
for append in batch {
let effect_key = append.effect_key.map(EffectKey::to_hex);
let body = append.into_body(head.seq + 1, epoch);
let is_start =
matches!(body.kind, crate::journal::RecordKind::EffectStarted { .. });
let conclusion = match &body.kind {
crate::journal::RecordKind::RunSealed { outcome, .. } => {
Some(outcome.clone())
}
_ => None,
};
let record = Record::seal_signed(body, 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.attestation.as_ref().map(|a| a.key_id.clone()),
signature: record.attestation.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(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 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))?;
if let Some(previous) = last
.get((tenant.as_str(), key.as_str()))
.map_err(|e| be(&e))?
.map(|value| value.value())
{
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))?;
sealed
};
w.commit().map_err(|e| be(&e))?;
Ok(sealed)
})
.await
}
async fn read(&self, run: RunId, from: Seq) -> Result<Vec<Record>, 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))?;
let mut out = Vec::new();
for entry in t
.range((key.as_str(), from)..=(key.as_str(), u64::MAX))
.map_err(|e| be(&e))?
{
let (_, v) = entry.map_err(|e| be(&e))?;
out.push(Row::decode(v.value())?.into_record()?);
}
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();
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()?);
}
}
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 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 = now + ttl.as_secs().max(1);
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 => 1,
Some((_, epoch, expires_at)) if expires_at <= now => epoch + 1,
Some((held_by, epoch, _)) if held_by == owner => epoch,
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 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 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))?;
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 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();
if let Some(id) = key.strip_prefix(prefix.as_str())
&& let Ok(run) = RunId::parse(id)
{
out.push(run);
}
}
Ok(out)
})
.await
}
async fn recent_runs(&self) -> Result<Vec<(RunId, u64)>, StoreError> {
let tenant = self.tenant_name();
let prefix = format!("{tenant}/");
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());
};
activity
.range((tenant.as_str(), 0, "")..=(tenant.as_str(), u64::MAX, MAX_STR))
.map_err(|e| be(&e))?
.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))),
})
.collect()
})
.await
}
async fn checkpoint(&self) -> Result<crate::journal::Checkpoint, StoreError> {
let leaves = self.log_leaves().await?;
Ok(crate::journal::Checkpoint {
origin: self.origin.clone(),
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.log_position(run).await? else {
return Ok(None);
};
Ok(Some(crate::journal::Inclusion {
index: index as u64,
size: leaves.len() as u64,
seal,
proof: crate::core::merkle::inclusion_proof(&leaves, index),
}))
}
async fn request_cancel(
&self,
run: RunId,
actor: &str,
reason: &str,
) -> Result<bool, StoreError> {
let key = self.run_key(run);
let (actor, reason) = (actor.to_owned(), 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(), (actor.as_str(), 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))?;
Ok(t.get(key.as_str()).map_err(|e| be(&e))?.map(|v| {
let (actor, reason, _) = v.value();
Cancellation {
actor: actor.to_owned(),
reason: reason.to_owned(),
}
}))
})
.await
}
}