#![allow(dead_code)]
use std::future::Future;
use std::time::Duration;
use reliar_outbox::{
AcquireRequest, AcquiredBatch, FailedRecord, OutboxStats, OutboxStore, PoisonedRow,
PurgeReport, PurgeRequest, RecordRef, WorkerId,
};
#[derive(Clone)]
pub(crate) struct OverDeliveringStore<S> {
inner: S,
over_claim_batch_size: u32,
}
impl<S> OverDeliveringStore<S> {
pub(crate) fn new(inner: S, over_claim_batch_size: u32) -> Self {
Self {
inner,
over_claim_batch_size,
}
}
}
impl<S: OutboxStore> OutboxStore for OverDeliveringStore<S> {
type Error = S::Error;
fn acquire(
&self,
request: AcquireRequest,
) -> impl Future<Output = Result<AcquiredBatch, Self::Error>> + Send {
self.inner
.acquire(request.batch_size(self.over_claim_batch_size))
}
fn complete(
&self,
worker: &WorkerId,
items: &[RecordRef],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.complete(worker, items)
}
fn fail(
&self,
worker: &WorkerId,
items: &[FailedRecord],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.fail(worker, items)
}
fn release(
&self,
worker: &WorkerId,
items: &[RecordRef],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.release(worker, items)
}
fn extend_lease(
&self,
worker: &WorkerId,
items: &[RecordRef],
lease: Duration,
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.extend_lease(worker, items, lease)
}
fn purge(
&self,
request: PurgeRequest,
) -> impl Future<Output = Result<PurgeReport, Self::Error>> + Send {
self.inner.purge(request)
}
fn stats(&self) -> impl Future<Output = Result<OutboxStats, Self::Error>> + Send {
self.inner.stats()
}
}
#[derive(Clone, Default)]
pub(crate) struct CountingStore<S> {
inner: S,
calls: std::sync::Arc<std::sync::atomic::AtomicUsize>,
}
impl<S> CountingStore<S> {
pub(crate) fn new(inner: S) -> Self {
Self {
inner,
calls: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
}
}
pub(crate) fn claim_calls(&self) -> usize {
self.calls.load(std::sync::atomic::Ordering::SeqCst)
}
}
impl<S: OutboxStore> OutboxStore for CountingStore<S> {
type Error = S::Error;
fn acquire(
&self,
request: AcquireRequest,
) -> impl Future<Output = Result<AcquiredBatch, Self::Error>> + Send {
self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
self.inner.acquire(request)
}
fn complete(
&self,
worker: &WorkerId,
items: &[RecordRef],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.complete(worker, items)
}
fn fail(
&self,
worker: &WorkerId,
items: &[FailedRecord],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.fail(worker, items)
}
fn release(
&self,
worker: &WorkerId,
items: &[RecordRef],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.release(worker, items)
}
fn extend_lease(
&self,
worker: &WorkerId,
items: &[RecordRef],
lease: Duration,
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.extend_lease(worker, items, lease)
}
fn purge(
&self,
request: PurgeRequest,
) -> impl Future<Output = Result<PurgeReport, Self::Error>> + Send {
self.inner.purge(request)
}
fn stats(&self) -> impl Future<Output = Result<OutboxStats, Self::Error>> + Send {
self.inner.stats()
}
}
#[derive(Clone)]
pub(crate) struct PoisoningStore<S> {
inner: S,
n_healthy: usize,
}
impl<S> PoisoningStore<S> {
pub(crate) fn new(inner: S, n_healthy: usize) -> Self {
Self { inner, n_healthy }
}
}
impl<S: OutboxStore> OutboxStore for PoisoningStore<S> {
type Error = S::Error;
fn acquire(
&self,
request: AcquireRequest,
) -> impl Future<Output = Result<AcquiredBatch, Self::Error>> + Send {
let n_healthy = self.n_healthy;
let claim = self.inner.acquire(request);
async move {
let batch = claim.await?;
let mut records = batch.records;
let to_poison = records.split_off(n_healthy.min(records.len()));
let mut poisoned = batch.poisoned;
poisoned.extend(to_poison.iter().map(|record| {
PoisonedRow::new(
record.id,
record.envelope.id,
"harness fixture: poisoned by PoisoningStore",
)
}));
Ok(AcquiredBatch::new(records, poisoned))
}
}
fn complete(
&self,
worker: &WorkerId,
items: &[RecordRef],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.complete(worker, items)
}
fn fail(
&self,
worker: &WorkerId,
items: &[FailedRecord],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.fail(worker, items)
}
fn release(
&self,
worker: &WorkerId,
items: &[RecordRef],
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.release(worker, items)
}
fn extend_lease(
&self,
worker: &WorkerId,
items: &[RecordRef],
lease: Duration,
) -> impl Future<Output = Result<u64, Self::Error>> + Send {
self.inner.extend_lease(worker, items, lease)
}
fn purge(
&self,
request: PurgeRequest,
) -> impl Future<Output = Result<PurgeReport, Self::Error>> + Send {
self.inner.purge(request)
}
fn stats(&self) -> impl Future<Output = Result<OutboxStats, Self::Error>> + Send {
self.inner.stats()
}
}