use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use crate::memory::InMemoryStore;
use crate::package::{PackageRecord, PackageRouteRecord, PackageStore};
use crate::{
Event, OutboxRow, ReadableEventStore, RunSummary, StoreError, TimerEntry, TimerId,
WorkflowFilter, WorkflowId, WorkflowSummary, WritableEventStore, WriteToken,
};
pub const REFUSED_SHARD: usize = 3;
pub struct FencedHistoryStore {
inner: InMemoryStore,
refusing: Arc<AtomicBool>,
}
impl FencedHistoryStore {
#[must_use]
pub fn new() -> Self {
Self {
inner: InMemoryStore::default(),
refusing: Arc::new(AtomicBool::new(false)),
}
}
pub fn arm_fence(&self) {
self.refusing.store(true, Ordering::SeqCst);
}
pub fn disarm_fence(&self) {
self.refusing.store(false, Ordering::SeqCst);
}
fn fence(&self) -> Result<(), StoreError> {
if self.refusing.load(Ordering::SeqCst) {
return Err(StoreError::NotOwner {
shard: REFUSED_SHARD,
});
}
Ok(())
}
}
impl Default for FencedHistoryStore {
fn default() -> Self {
Self::new()
}
}
#[async_trait]
impl ReadableEventStore for FencedHistoryStore {
async fn read_history(&self, workflow_id: &WorkflowId) -> Result<Vec<Event>, StoreError> {
self.fence()?;
self.inner.read_history(workflow_id).await
}
async fn read_history_from(
&self,
workflow_id: &WorkflowId,
from_seq: u64,
) -> Result<Vec<Event>, StoreError> {
self.fence()?;
self.inner.read_history_from(workflow_id, from_seq).await
}
async fn read_run_chain(
&self,
workflow_id: &WorkflowId,
) -> Result<Vec<RunSummary>, StoreError> {
self.fence()?;
self.inner.read_run_chain(workflow_id).await
}
async fn list_active(&self) -> Result<Vec<WorkflowId>, StoreError> {
self.inner.list_active().await
}
async fn list_workflow_ids(&self) -> Result<Vec<WorkflowId>, StoreError> {
self.inner.list_workflow_ids().await
}
async fn list_paused(&self) -> Result<Vec<WorkflowId>, StoreError> {
self.inner.list_paused().await
}
async fn query(&self, filter: &WorkflowFilter) -> Result<Vec<WorkflowSummary>, StoreError> {
self.inner.query(filter).await
}
async fn schedule_timer(
&self,
workflow_id: &WorkflowId,
timer_id: &TimerId,
fire_at: DateTime<Utc>,
) -> Result<(), StoreError> {
self.inner
.schedule_timer(workflow_id, timer_id, fire_at)
.await
}
async fn expired_timers(&self, as_of: DateTime<Utc>) -> Result<Vec<TimerEntry>, StoreError> {
self.inner.expired_timers(as_of).await
}
fn set_owned_shards(&self, shards: Option<&[usize]>) {
self.inner.set_owned_shards(shards);
}
fn acquire_owned_shards(&self, shards: &[usize]) -> Result<(), StoreError> {
self.inner.acquire_owned_shards(shards)
}
fn acquire_owned_shard(&self, shard: usize) -> Result<(), StoreError> {
self.inner.acquire_owned_shard(shard)
}
fn extend_owned_shards(&self, shards: &[usize]) {
self.inner.extend_owned_shards(shards);
}
fn is_current_owner(&self, shard: usize) -> bool {
if self.refusing.load(Ordering::SeqCst) {
return false;
}
self.inner.is_current_owner(shard)
}
fn publish_shard_owner(&self, shard: usize) -> Result<(), StoreError> {
self.inner.publish_shard_owner(shard)
}
}
#[async_trait]
impl WritableEventStore for FencedHistoryStore {
async fn append(
&self,
token: WriteToken,
workflow_id: &WorkflowId,
events: &[Event],
expected_seq: u64,
) -> Result<(), StoreError> {
self.fence()?;
self.inner
.append(token, workflow_id, events, expected_seq)
.await
}
async fn append_with_outbox(
&self,
token: WriteToken,
workflow_id: &WorkflowId,
events: &[Event],
expected_seq: u64,
outbox_rows: &[OutboxRow],
) -> Result<(), StoreError> {
self.fence()?;
self.inner
.append_with_outbox(token, workflow_id, events, expected_seq, outbox_rows)
.await
}
async fn rearm_outbox_pending(&self, rows: &[OutboxRow]) -> Result<(), StoreError> {
self.inner.rearm_outbox_pending(rows).await
}
async fn settle_outbox_row_cancelled(&self, dispatch_key: &str) -> Result<(), StoreError> {
self.inner.settle_outbox_row_cancelled(dispatch_key).await
}
async fn settle_workflow_outbox_rows_cancelled(
&self,
workflow_id: &WorkflowId,
) -> Result<Vec<String>, StoreError> {
self.inner
.settle_workflow_outbox_rows_cancelled(workflow_id)
.await
}
}
#[async_trait]
impl PackageStore for FencedHistoryStore {
async fn put_package(&self, record: PackageRecord) -> Result<(), StoreError> {
self.inner.put_package(record).await
}
async fn put_package_with_routes(
&self,
record: PackageRecord,
route_workflow_types: &[String],
) -> Result<(), StoreError> {
self.inner
.put_package_with_routes(record, route_workflow_types)
.await
}
async fn list_packages(&self) -> Result<Vec<PackageRecord>, StoreError> {
self.inner.list_packages().await
}
async fn delete_package(
&self,
workflow_type: &str,
content_hash: &str,
) -> Result<(), StoreError> {
self.inner.delete_package(workflow_type, content_hash).await
}
async fn put_package_route(
&self,
workflow_type: &str,
content_hash: &str,
) -> Result<(), StoreError> {
self.inner
.put_package_route(workflow_type, content_hash)
.await
}
async fn list_package_routes(&self) -> Result<Vec<PackageRouteRecord>, StoreError> {
self.inner.list_package_routes().await
}
}
#[cfg(test)]
mod tests {
use super::{FencedHistoryStore, REFUSED_SHARD};
#[test]
fn an_armed_fence_claims_no_shard_at_all() {
let store = FencedHistoryStore::new();
let other_shard = REFUSED_SHARD + 2;
assert!(
store.is_current_owner(REFUSED_SHARD),
"control: an unarmed double must own the shard it will later refuse"
);
assert!(
store.is_current_owner(other_shard),
"control: an unarmed double must own the other shard too"
);
store.arm_fence();
assert!(
!store.is_current_owner(REFUSED_SHARD),
"the refused shard must not be claimed while the fence is armed"
);
assert!(
!store.is_current_owner(other_shard),
"no OTHER shard may be claimed either: the fence refuses reads of every shard, so \
claiming one is a state no real backend occupies"
);
store.disarm_fence();
assert!(
store.is_current_owner(other_shard),
"disarming must restore honest ownership"
);
}
use crate::{ReadableEventStore, StoreError, WorkflowId};
#[tokio::test]
async fn reads_are_honest_until_armed_and_refuse_after()
-> Result<(), Box<dyn std::error::Error>> {
let store = FencedHistoryStore::new();
let workflow_id = WorkflowId::new_v4();
assert!(
store.read_history(&workflow_id).await.is_ok(),
"an unarmed double must be indistinguishable from an honest store"
);
store.arm_fence();
assert!(
matches!(
store.read_history(&workflow_id).await,
Err(StoreError::NotOwner { shard }) if shard == REFUSED_SHARD
),
"an armed double must refuse history reads with NotOwner"
);
assert!(
matches!(
store.read_run_chain(&workflow_id).await,
Err(StoreError::NotOwner { shard }) if shard == REFUSED_SHARD
),
"an armed double must refuse run-chain reads with NotOwner"
);
assert!(
store.list_active().await.is_ok(),
"only history reads are fenced — a fixture still needs the listings"
);
store.disarm_fence();
assert!(
store.read_history(&workflow_id).await.is_ok(),
"disarming must restore honest reads"
);
Ok(())
}
#[tokio::test]
async fn appends_are_fenced_as_well_as_reads() -> Result<(), Box<dyn std::error::Error>> {
use crate::{WritableEventStore, WriteToken};
let store = FencedHistoryStore::new();
let workflow_id = WorkflowId::new_v4();
assert!(
store
.append(WriteToken::recorder(), &workflow_id, &[], 0)
.await
.is_ok(),
"control: an unarmed double must accept an append"
);
store.arm_fence();
assert!(
matches!(
store
.append(WriteToken::recorder(), &workflow_id, &[], 0)
.await,
Err(StoreError::NotOwner { shard }) if shard == REFUSED_SHARD
),
"an armed double owns no shard, so it must refuse the append with NotOwner"
);
assert!(
matches!(
store
.append_with_outbox(WriteToken::recorder(), &workflow_id, &[], 0, &[])
.await,
Err(StoreError::NotOwner { shard }) if shard == REFUSED_SHARD
),
"the outbox append is the same history write and must refuse identically"
);
store.disarm_fence();
assert!(
store
.append(WriteToken::recorder(), &workflow_id, &[], 0)
.await
.is_ok(),
"disarming must restore honest appends"
);
Ok(())
}
}