use aion_core::{Event, WorkflowId};
use aion_store::{InMemoryStore, ReadableEventStore};
pub(crate) struct FlakyStore {
inner: InMemoryStore,
append_failures_to_skip: std::sync::atomic::AtomicU32,
remaining_append_failures: std::sync::atomic::AtomicU32,
read_failures_to_skip: std::sync::atomic::AtomicU32,
remaining_read_failures: std::sync::atomic::AtomicU32,
remaining_read_conflicts: std::sync::atomic::AtomicU32,
remaining_read_lost_ownership: std::sync::atomic::AtomicU32,
remaining_retirement_failures: std::sync::atomic::AtomicU32,
ack_lost_appends_to_skip: std::sync::atomic::AtomicU32,
remaining_ack_lost_appends: std::sync::atomic::AtomicU32,
remaining_range_read_failures: std::sync::atomic::AtomicU32,
remaining_fenced_appends: std::sync::atomic::AtomicU32,
}
impl FlakyStore {
pub(crate) fn new() -> Self {
Self {
inner: InMemoryStore::default(),
append_failures_to_skip: std::sync::atomic::AtomicU32::new(0),
remaining_append_failures: std::sync::atomic::AtomicU32::new(0),
read_failures_to_skip: std::sync::atomic::AtomicU32::new(0),
remaining_read_failures: std::sync::atomic::AtomicU32::new(0),
remaining_read_conflicts: std::sync::atomic::AtomicU32::new(0),
remaining_read_lost_ownership: std::sync::atomic::AtomicU32::new(0),
remaining_retirement_failures: std::sync::atomic::AtomicU32::new(0),
ack_lost_appends_to_skip: std::sync::atomic::AtomicU32::new(0),
remaining_ack_lost_appends: std::sync::atomic::AtomicU32::new(0),
remaining_range_read_failures: std::sync::atomic::AtomicU32::new(0),
remaining_fenced_appends: std::sync::atomic::AtomicU32::new(0),
}
}
pub(crate) fn lose_ack_of_next_appends(&self, count: u32) {
self.lose_ack_of_appends_after(0, count);
}
pub(crate) fn lose_ack_of_appends_after(&self, skip: u32, count: u32) {
self.ack_lost_appends_to_skip
.store(skip, std::sync::atomic::Ordering::Release);
self.remaining_ack_lost_appends
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fail_next_range_reads(&self, count: u32) {
self.remaining_range_read_failures
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fence_next_appends(&self, count: u32) {
self.remaining_fenced_appends
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fail_next_appends(&self, count: u32) {
self.fail_appends_after(0, count);
}
pub(crate) fn fail_appends_after(&self, skip: u32, count: u32) {
self.append_failures_to_skip
.store(skip, std::sync::atomic::Ordering::Release);
self.remaining_append_failures
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fail_next_reads(&self, count: u32) {
self.fail_reads_after(0, count);
}
pub(crate) fn fail_reads_after(&self, skip: u32, count: u32) {
self.read_failures_to_skip
.store(skip, std::sync::atomic::Ordering::Release);
self.remaining_read_failures
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fail_next_reads_with_conflict(&self, count: u32) {
self.remaining_read_conflicts
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fail_next_reads_with_lost_ownership(&self, count: u32) {
self.remaining_read_lost_ownership
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) fn fail_next_retirements(&self, count: u32) {
self.remaining_retirement_failures
.store(count, std::sync::atomic::Ordering::Release);
}
pub(crate) async fn recorded_history(
&self,
workflow_id: &WorkflowId,
) -> Result<Vec<Event>, aion_store::StoreError> {
self.inner.read_history(workflow_id).await
}
pub(crate) fn unspent_append_failures(&self) -> u32 {
self.remaining_append_failures
.load(std::sync::atomic::Ordering::Acquire)
}
fn take_failure(budget: &std::sync::atomic::AtomicU32) -> bool {
budget
.fetch_update(
std::sync::atomic::Ordering::AcqRel,
std::sync::atomic::Ordering::Acquire,
|current| current.checked_sub(1),
)
.is_ok()
}
}
#[async_trait::async_trait]
impl aion_store::ReadableEventStore for FlakyStore {
async fn read_history(
&self,
workflow_id: &WorkflowId,
) -> Result<Vec<Event>, aion_store::StoreError> {
if Self::take_failure(&self.remaining_read_conflicts) {
return Err(aion_store::StoreError::SequenceConflict {
expected: 0,
found: 1,
});
}
if Self::take_failure(&self.remaining_read_lost_ownership) {
return Err(aion_store::StoreError::NotOwner { shard: 7 });
}
if !Self::take_failure(&self.read_failures_to_skip)
&& Self::take_failure(&self.remaining_read_failures)
{
return Err(aion_store::StoreError::Backend(
"transient history read failure injected by FlakyStore".to_owned(),
));
}
self.inner.read_history(workflow_id).await
}
async fn read_history_from(
&self,
workflow_id: &WorkflowId,
from_seq: u64,
) -> Result<Vec<Event>, aion_store::StoreError> {
if Self::take_failure(&self.remaining_range_read_failures) {
return Err(aion_store::StoreError::Backend(
"transient range read failure injected by FlakyStore".to_owned(),
));
}
self.inner.read_history_from(workflow_id, from_seq).await
}
async fn read_run_chain(
&self,
workflow_id: &WorkflowId,
) -> Result<Vec<aion_store::RunSummary>, aion_store::StoreError> {
self.inner.read_run_chain(workflow_id).await
}
async fn list_workflow_ids(&self) -> Result<Vec<WorkflowId>, aion_store::StoreError> {
self.inner.list_workflow_ids().await
}
async fn list_active(&self) -> Result<Vec<WorkflowId>, aion_store::StoreError> {
self.inner.list_active().await
}
async fn list_paused(&self) -> Result<Vec<WorkflowId>, aion_store::StoreError> {
self.inner.list_paused().await
}
async fn stream_heads(
&self,
) -> Result<Vec<aion_store::visibility::StreamHead>, aion_store::StoreError> {
self.inner.stream_heads().await
}
async fn query(
&self,
filter: &aion_core::WorkflowFilter,
) -> Result<Vec<aion_core::WorkflowSummary>, aion_store::StoreError> {
self.inner.query(filter).await
}
async fn schedule_timer(
&self,
workflow_id: &WorkflowId,
timer_id: &aion_core::TimerId,
fire_at: chrono::DateTime<chrono::Utc>,
armed_seq: u64,
) -> Result<(), aion_store::StoreError> {
self.inner
.schedule_timer(workflow_id, timer_id, fire_at, armed_seq)
.await
}
async fn retire_timer(
&self,
workflow_id: &WorkflowId,
timer_id: &aion_core::TimerId,
fire_at: chrono::DateTime<chrono::Utc>,
armed_seq: u64,
) -> Result<aion_store::TimerRetirement, aion_store::StoreError> {
if Self::take_failure(&self.remaining_retirement_failures) {
return Err(aion_store::StoreError::Backend(
"transient timer-row retirement failure injected by FlakyStore".to_owned(),
));
}
self.inner
.retire_timer(workflow_id, timer_id, fire_at, armed_seq)
.await
}
async fn expired_timers(
&self,
as_of: chrono::DateTime<chrono::Utc>,
) -> Result<Vec<aion_store::TimerEntry>, aion_store::StoreError> {
self.inner.expired_timers(as_of).await
}
}
#[async_trait::async_trait]
impl aion_store::WritableEventStore for FlakyStore {
async fn append(
&self,
token: aion_store::WriteToken,
workflow_id: &WorkflowId,
events: &[Event],
expected_seq: u64,
) -> Result<(), aion_store::StoreError> {
if Self::take_failure(&self.remaining_fenced_appends) {
return Err(aion_store::StoreError::NotOwner { shard: 7 });
}
if !Self::take_failure(&self.append_failures_to_skip)
&& Self::take_failure(&self.remaining_append_failures)
{
return Err(aion_store::StoreError::Backend(
"transient append failure injected by FlakyStore".to_owned(),
));
}
self.inner
.append(token, workflow_id, events, expected_seq)
.await?;
if !Self::take_failure(&self.ack_lost_appends_to_skip)
&& Self::take_failure(&self.remaining_ack_lost_appends)
{
return Err(aion_store::StoreError::Backend(
"append landed but its acknowledgement was lost, injected by FlakyStore".to_owned(),
));
}
Ok(())
}
}
#[async_trait::async_trait]
impl aion_store::PackageStore for FlakyStore {
async fn put_package(
&self,
record: aion_store::PackageRecord,
) -> Result<(), aion_store::StoreError> {
self.inner.put_package(record).await
}
async fn put_package_with_routes(
&self,
record: aion_store::PackageRecord,
route_workflow_types: &[String],
) -> Result<(), aion_store::StoreError> {
self.inner
.put_package_with_routes(record, route_workflow_types)
.await
}
async fn list_packages(
&self,
) -> Result<Vec<aion_store::PackageRecord>, aion_store::StoreError> {
self.inner.list_packages().await
}
async fn delete_package(
&self,
workflow_type: &str,
content_hash: &str,
) -> Result<(), aion_store::StoreError> {
self.inner.delete_package(workflow_type, content_hash).await
}
async fn put_package_route(
&self,
workflow_type: &str,
content_hash: &str,
) -> Result<(), aion_store::StoreError> {
self.inner
.put_package_route(workflow_type, content_hash)
.await
}
async fn list_package_routes(
&self,
) -> Result<Vec<aion_store::PackageRouteRecord>, aion_store::StoreError> {
self.inner.list_package_routes().await
}
}