use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use serde_json::Value;
use crate::core::{Epoch, RunId, StoreError};
use crate::journal::{Append, AtomicJournal, AtomicTx, AtomicWork, JournalStore, Record, SqlValue};
pub type Statement = (String, Vec<SqlValue>);
#[derive(Debug)]
pub struct StagedAtomic {
inner: Arc<dyn JournalStore>,
applied: Mutex<Vec<Statement>>,
lose_ack: std::sync::atomic::AtomicBool,
}
impl StagedAtomic {
#[must_use]
pub fn wrap(inner: Arc<dyn JournalStore>) -> Arc<Self> {
Arc::new(Self {
inner,
applied: Mutex::new(Vec::new()),
lose_ack: std::sync::atomic::AtomicBool::new(false),
})
}
pub fn lose_commit_acknowledgement(&self) {
self.lose_ack
.store(true, std::sync::atomic::Ordering::SeqCst);
}
#[must_use]
pub fn applied(&self) -> Vec<Statement> {
self.applied.lock().expect("staged").clone()
}
}
#[derive(Debug, Default)]
struct Staging(Mutex<Vec<Statement>>);
#[async_trait]
impl AtomicTx for Staging {
async fn execute(&self, sql: &str, params: &[SqlValue]) -> Result<u64, StoreError> {
self.0
.lock()
.expect("staging")
.push((sql.to_owned(), params.to_vec()));
Ok(1)
}
async fn query(&self, _sql: &str, _params: &[SqlValue]) -> Result<Vec<Value>, StoreError> {
Err(StoreError::Backend(
"this fixture stages writes and does not answer queries — read from a \
real database, or do not read"
.to_owned(),
))
}
}
#[async_trait]
impl AtomicJournal for StagedAtomic {
async fn append_atomic(
&self,
_run: RunId,
epoch: Epoch,
work: &dyn AtomicWork,
) -> Result<Vec<Record>, StoreError> {
let staging = Staging::default();
let batch = work
.run(&staging)
.await
.map_err(|e| StoreError::Backend(e.to_string()))?;
let sealed = self.inner.append(epoch, batch).await?;
self.applied
.lock()
.expect("staged")
.extend(staging.0.lock().expect("staging").drain(..));
if self.lose_ack.load(std::sync::atomic::Ordering::SeqCst) {
return Err(StoreError::CommitUnknown {
detail: "the fixture dropped the acknowledgement".to_owned(),
});
}
Ok(sealed)
}
}
#[async_trait]
impl JournalStore for StagedAtomic {
fn is_shared(&self) -> bool {
self.inner.is_shared()
}
fn tenant(&self) -> &str {
self.inner.tenant()
}
fn atomic(&self) -> Option<&dyn AtomicJournal> {
Some(self)
}
async fn append(&self, epoch: Epoch, batch: Vec<Append>) -> Result<Vec<Record>, StoreError> {
self.inner.append(epoch, batch).await
}
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 admitted_as(&self, key: &str) -> Result<Option<crate::core::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<crate::core::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 recent_runs(
&self,
after: Option<(u64, RunId)>,
limit: usize,
) -> Result<Vec<(RunId, u64)>, StoreError> {
self.inner.recent_runs(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 request_cancel(
&self,
run: RunId,
actor: &str,
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
}
}