use std::sync::Arc;
use async_trait::async_trait;
use crate::core::{EffectError, Epoch, RunId, StoreError};
use crate::journal::{
Append, AtomicJournal, AtomicTx, AtomicWork, JournalStore, Record, RecordKind,
};
#[derive(Debug)]
pub(crate) struct PassJournal {
inner: Arc<dyn JournalStore>,
run: RunId,
marker: tokio::sync::Mutex<Option<RecordKind>>,
}
impl PassJournal {
pub(crate) fn new(inner: Arc<dyn JournalStore>, run: RunId, marker: RecordKind) -> Self {
Self {
inner,
run,
marker: tokio::sync::Mutex::new(Some(marker)),
}
}
fn marker_for(&self, kind: &RecordKind, batch: &[Append]) -> Option<Append> {
let first = batch.iter().find(|a| a.run == self.run)?;
let mut marker = Append::new(self.run, kind.clone());
marker.case = first.case;
Some(marker)
}
fn leading_conclusion(&self, batch: &[Append]) -> Option<usize> {
let first = batch.iter().position(|a| a.run == self.run)?;
matches!(batch[first].kind, RecordKind::RunConcluded { .. }).then_some(first)
}
}
fn without_marker(mut written: Vec<Record>) -> Vec<Record> {
if !written.is_empty() {
written.remove(0);
}
written
}
#[async_trait]
impl JournalStore for PassJournal {
fn is_shared(&self) -> bool {
self.inner.is_shared()
}
fn seals(&self) -> bool {
self.inner.seals()
}
fn tenant(&self) -> &str {
self.inner.tenant()
}
fn atomic(&self) -> Option<&dyn AtomicJournal> {
self.inner.atomic().map(|_| self as &dyn AtomicJournal)
}
async fn append(&self, epoch: Epoch, batch: Vec<Append>) -> Result<Vec<Record>, StoreError> {
let mut pending = self.marker.lock().await;
let Some(marker) = pending
.as_ref()
.and_then(|kind| self.marker_for(kind, &batch))
else {
drop(pending);
return self.inner.append(epoch, batch).await;
};
if let Some(conclusion) = self.leading_conclusion(&batch) {
let written = self.inner.append(epoch, vec![marker]).await?;
*pending = None;
let mut batch = batch;
if let Some(marker) = written.last()
&& let RecordKind::RunConcluded { chain_head, .. } = &mut batch[conclusion].kind
&& *chain_head == marker.prev_hash
{
*chain_head = marker.hash;
}
drop(pending);
return self.inner.append(epoch, batch).await;
}
let mut with_marker = Vec::with_capacity(batch.len() + 1);
with_marker.push(marker);
with_marker.extend(batch);
let written = self.inner.append(epoch, with_marker).await?;
*pending = None;
Ok(without_marker(written))
}
async fn read(&self, run: RunId, from: crate::core::Seq) -> Result<Vec<Record>, StoreError> {
self.inner.read(run, from).await
}
async fn read_page(
&self,
run: RunId,
from: crate::core::Seq,
limit: usize,
) -> Result<Vec<Record>, StoreError> {
self.inner.read_page(run, from, limit).await
}
async fn acquire(
&self,
run: RunId,
owner: &str,
ttl: std::time::Duration,
) -> Result<crate::journal::Lease, StoreError> {
self.inner.acquire(run, owner, ttl).await
}
async fn renew(
&self,
run: RunId,
owner: &str,
epoch: Epoch,
ttl: std::time::Duration,
) -> Result<crate::journal::Lease, StoreError> {
self.inner.renew(run, owner, epoch, ttl).await
}
async fn release_lease(&self, run: RunId, epoch: Epoch) -> Result<(), StoreError> {
self.inner.release_lease(run, epoch).await
}
async fn abandoned_runs(&self, limit: usize) -> Result<Vec<RunId>, StoreError> {
self.inner.abandoned_runs(limit).await
}
async fn waiting_runs(
&self,
limit: usize,
) -> Result<Vec<crate::journal::WaitingRun>, StoreError> {
self.inner.waiting_runs(limit).await
}
async fn admitted_as(&self, key: &str) -> Result<Option<RunId>, StoreError> {
self.inner.admitted_as(key).await
}
async fn forget_admissions(
&self,
older_than: crate::core::Timestamp,
) -> Result<usize, StoreError> {
self.inner.forget_admissions(older_than).await
}
async fn runs_by_outcome(&self, outcome: &str, limit: usize) -> Result<Vec<RunId>, StoreError> {
self.inner.runs_by_outcome(outcome, limit).await
}
async fn count_by_outcome(&self, outcome: &str) -> Result<u64, StoreError> {
self.inner.count_by_outcome(outcome).await
}
async fn runs_by_id(
&self,
after: Option<RunId>,
limit: usize,
) -> Result<Vec<RunId>, StoreError> {
self.inner.runs_by_id(after, limit).await
}
async fn recent_runs(
&self,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError> {
self.inner.recent_runs(after, limit).await
}
async fn recent_runs_from(
&self,
source: &str,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError> {
self.inner.recent_runs_from(source, after, limit).await
}
async fn case_history(
&self,
case: crate::core::CaseId,
limit: usize,
) -> Result<Vec<Record>, StoreError> {
self.inner.case_history(case, limit).await
}
async fn head(&self, run: RunId) -> Result<crate::journal::Head, StoreError> {
self.inner.head(run).await
}
async fn seal(
&self,
run: RunId,
epoch: Epoch,
outcome: &str,
) -> Result<crate::core::Digest, StoreError> {
self.inner.seal(run, epoch, outcome).await
}
async fn checkpoint(&self) -> Result<crate::journal::Checkpoint, StoreError> {
self.inner.checkpoint().await
}
async fn consistency_proof(
&self,
old_size: u64,
) -> Result<Vec<crate::core::Digest>, StoreError> {
self.inner.consistency_proof(old_size).await
}
async fn inclusion_proof(
&self,
run: RunId,
) -> Result<Option<crate::journal::Inclusion>, StoreError> {
self.inner.inclusion_proof(run).await
}
async fn inclusion_proof_at(
&self,
run: RunId,
size: u64,
) -> Result<Option<crate::journal::Inclusion>, StoreError> {
self.inner.inclusion_proof_at(run, size).await
}
async fn log_positions(
&self,
runs: &[RunId],
) -> Result<Vec<Option<(u64, crate::core::Digest)>>, StoreError> {
self.inner.log_positions(runs).await
}
async fn request_cancel(
&self,
run: RunId,
actor: &crate::core::Operator,
reason: &str,
) -> Result<bool, StoreError> {
self.inner.request_cancel(run, actor, reason).await
}
async fn cancellation(
&self,
run: RunId,
) -> Result<Option<crate::journal::Cancellation>, StoreError> {
self.inner.cancellation(run).await
}
async fn verify(&self, run: RunId) -> Result<crate::core::Digest, StoreError> {
self.inner.verify(run).await
}
}
struct MarkedWork<'a> {
work: &'a dyn AtomicWork,
journal: &'a PassJournal,
kind: &'a RecordKind,
marked: std::sync::atomic::AtomicBool,
}
#[async_trait]
impl AtomicWork for MarkedWork<'_> {
async fn run(&self, tx: &dyn AtomicTx) -> Result<Vec<Append>, EffectError> {
let batch = self.work.run(tx).await?;
let Some(marker) = self.journal.marker_for(self.kind, &batch) else {
return Ok(batch);
};
self.marked
.store(true, std::sync::atomic::Ordering::Relaxed);
let mut with_marker = Vec::with_capacity(batch.len() + 1);
with_marker.push(marker);
with_marker.extend(batch);
Ok(with_marker)
}
}
#[async_trait]
impl AtomicJournal for PassJournal {
async fn append_atomic(
&self,
run: RunId,
epoch: Epoch,
work: &dyn AtomicWork,
) -> Result<Vec<Record>, StoreError> {
let Some(inner) = self.inner.atomic() else {
return Err(StoreError::Backend(
"the pass journal's store has no transaction a resource can join".to_owned(),
));
};
let mut pending = self.marker.lock().await;
let Some(kind) = pending.clone().filter(|_| run == self.run) else {
drop(pending);
return inner.append_atomic(run, epoch, work).await;
};
let marked = MarkedWork {
work,
journal: self,
kind: &kind,
marked: std::sync::atomic::AtomicBool::new(false),
};
let written = inner.append_atomic(run, epoch, &marked).await?;
if !marked.marked.load(std::sync::atomic::Ordering::Relaxed) {
return Ok(written);
}
*pending = None;
Ok(without_marker(written))
}
}