use std::fmt::Debug;
use std::time::Duration;
use async_trait::async_trait;
use crate::core::{Digest, Epoch, RunId, Seq, StoreError};
use super::{Append, Record};
#[derive(Debug, Clone, PartialEq)]
pub struct WaitingRun {
pub run: RunId,
pub reason: crate::core::SuspendReason,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Head {
pub seq: Seq,
pub hash: Digest,
}
impl Head {
#[must_use]
pub const fn genesis() -> Self {
Self {
seq: 0,
hash: Digest::ZERO,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Lease {
pub run: RunId,
pub owner: String,
pub epoch: Epoch,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Checkpoint {
pub origin: String,
pub size: u64,
pub root: Digest,
}
use super::note::{b64, unb64};
impl Checkpoint {
#[must_use]
pub fn to_note(&self) -> String {
format!(
"{}\n{}\n{}\n",
self.origin,
self.size,
b64(self.root.as_bytes())
)
}
#[must_use]
pub fn is_coherent(&self) -> bool {
self.size != 0 || self.root == crate::core::merkle::empty_root()
}
pub fn from_note(note: &str) -> Result<Self, StoreError> {
let bad = |what: &str| StoreError::Backend(format!("checkpoint note: {what}"));
let body = note.strip_suffix('\n').ok_or_else(|| {
bad("the note does not end in a newline, which is part of what \
gets signed")
})?;
let mut parts = body.split('\n');
let origin = parts.next().ok_or_else(|| bad("no origin"))?;
let size = parts.next().ok_or_else(|| bad("no size"))?;
let root = parts.next().ok_or_else(|| bad("no root"))?;
if parts.next().is_some() {
return Err(bad(
"the note carries more than the three lines tlog-checkpoint defines; \
extra lines are refused rather than ignored, because a parser that \
ignores them lets two different signed texts name one checkpoint",
));
}
if origin.is_empty() {
return Err(bad("the origin is empty, so the note names no log"));
}
if size.is_empty()
|| !size.bytes().all(|b| b.is_ascii_digit())
|| (size.len() > 1 && size.starts_with('0'))
{
return Err(bad(
"the size is not a canonical decimal number — no sign, no leading zero, \
no surrounding space, because each is a second spelling of one log",
));
}
let size = size
.parse::<u64>()
.map_err(|e| bad(&format!("size is not a number: {e}")))?;
let root = unb64(root).ok_or_else(|| bad("root is not canonical RFC 4648 base64"))?;
let root: [u8; 32] = root.try_into().map_err(|_| bad("root is not 32 bytes"))?;
Ok(Self {
origin: origin.to_owned(),
size,
root: Digest::from_bytes(root),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Inclusion {
pub index: u64,
pub size: u64,
pub seal: Digest,
pub proof: Vec<Digest>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Cancellation {
pub actor: crate::core::Operator,
pub reason: String,
}
#[async_trait]
pub trait JournalStore: Send + Sync + Debug {
async fn append(&self, epoch: Epoch, batch: Vec<Append>) -> Result<Vec<Record>, StoreError>;
fn is_shared(&self) -> bool;
fn atomic(&self) -> Option<&dyn crate::journal::AtomicJournal> {
None
}
async fn read(&self, run: RunId, from: Seq) -> Result<Vec<Record>, StoreError>;
async fn read_page(
&self,
run: RunId,
from: Seq,
limit: usize,
) -> Result<Vec<Record>, StoreError>;
async fn runs_by_outcome(&self, outcome: &str, limit: usize) -> Result<Vec<RunId>, StoreError>;
async fn count_by_outcome(&self, outcome: &str) -> Result<u64, StoreError>;
async fn admitted_as(&self, key: &str) -> Result<Option<RunId>, StoreError>;
async fn forget_admissions(
&self,
older_than: crate::core::Timestamp,
) -> Result<usize, StoreError>;
async fn recent_runs(
&self,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError>;
async fn recent_runs_from(
&self,
source: &str,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError>;
async fn case_history(
&self,
case: crate::core::CaseId,
limit: usize,
) -> Result<Vec<Record>, StoreError>;
async fn head(&self, run: RunId) -> Result<Head, StoreError>;
async fn acquire(&self, run: RunId, owner: &str, ttl: Duration) -> Result<Lease, StoreError>;
async fn renew(
&self,
run: RunId,
owner: &str,
epoch: Epoch,
ttl: Duration,
) -> Result<Lease, StoreError>;
async fn abandoned_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError>;
async fn waiting_runs(&self, limit: usize) -> Result<Vec<WaitingRun>, StoreError>;
async fn release_lease(&self, run: RunId, epoch: Epoch) -> Result<(), StoreError>;
fn tenant(&self) -> &str {
crate::core::TenantId::DEFAULT
}
async fn seal(&self, run: RunId, epoch: Epoch, outcome: &str) -> Result<Digest, StoreError>;
async fn checkpoint(&self) -> Result<Checkpoint, StoreError>;
async fn consistency_proof(&self, old_size: u64) -> Result<Vec<Digest>, StoreError>;
async fn inclusion_proof(&self, run: RunId) -> Result<Option<Inclusion>, StoreError>;
async fn request_cancel(
&self,
run: RunId,
actor: &crate::core::Operator,
reason: &str,
) -> Result<bool, StoreError>;
async fn cancellation(&self, run: RunId) -> Result<Option<Cancellation>, StoreError>;
async fn verify(&self, run: RunId) -> Result<Digest, StoreError> {
let records = self.read(run, 1).await?;
Record::verify_chain(&records, Digest::ZERO)
}
}