use super::*;
use std::time::Duration;
use tempfile::TempDir;
const TENANT: &str = "synthetic-archive-consistency";
fn event() -> Event {
Event::from_strings(
"synthetic.updated".into(),
"synthetic-entity".into(),
TENANT.into(),
serde_json::json!({"synthetic": true}),
None,
)
.unwrap()
}
#[test]
fn eviction_cannot_erase_an_unflushed_conditional_predecessor() {
let directory = TempDir::new().unwrap();
let store = EventStore::with_config(EventStoreConfig::with_persistence(directory.path()));
store
.ingest_with_expected_version(&event(), Some(0))
.unwrap();
store.evict_tenant(TENANT);
assert!(
store.is_tenant_loaded(TENANT),
"pending history must remain resident"
);
assert!(
store
.ingest_with_expected_version(&event(), Some(0))
.is_err()
);
assert_eq!(
store
.query(&QueryEventsRequest {
tenant_id: Some(TENANT.into()),
..Default::default()
})
.unwrap()
.len(),
1
);
assert_eq!(
store
.ingest_with_expected_version(&event(), Some(1))
.unwrap(),
2
);
store.flush_storage().unwrap();
store.evict_tenant(TENANT);
assert!(!store.is_tenant_loaded(TENANT));
assert_eq!(
store
.ingest_with_expected_version(&event(), Some(2))
.unwrap(),
3
);
}
#[test]
fn in_memory_store_cannot_evict_its_only_history() {
let store = EventStore::new();
store
.ingest_with_expected_version(&event(), Some(0))
.unwrap();
store.evict_tenant(TENANT);
assert_eq!(store.total_events(), 1);
assert_eq!(
store
.ingest_with_expected_version(&event(), Some(1))
.unwrap(),
2
);
}
#[test]
fn replicated_history_reaches_parquet_before_it_can_be_evicted() {
let directory = TempDir::new().unwrap();
let store = EventStore::with_config(EventStoreConfig::with_persistence(directory.path()));
let mut replicated = event();
replicated.version = 1;
store.ingest_replicated(&replicated).unwrap();
store.evict_tenant(TENANT);
assert_eq!(store.total_events(), 1);
store.flush_storage().unwrap();
store.evict_tenant(TENANT);
assert_eq!(store.total_events(), 0);
let restored = store
.query(&QueryEventsRequest {
tenant_id: Some(TENANT.into()),
..Default::default()
})
.unwrap();
assert_eq!(restored.len(), 1);
assert_eq!(restored[0].id, replicated.id);
assert_eq!(restored[0].version, 1);
}
#[test]
fn read_only_replication_never_buffers_archive_writes_or_discards_memory() {
let directory = TempDir::new().unwrap();
let store = EventStore::with_config(EventStoreConfig {
read_only: true,
..EventStoreConfig::with_persistence(directory.path())
});
store.ingest_replicated(&event()).unwrap();
assert!(
!store
.storage
.as_ref()
.unwrap()
.read()
.has_pending_tenant_events(TENANT)
);
store.evict_tenant(TENANT);
assert_eq!(store.total_events(), 1);
}
#[test]
fn eviction_does_not_wait_for_in_flight_storage_work() {
let directory = TempDir::new().unwrap();
let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
directory.path(),
)));
store
.ingest_with_expected_version(&event(), Some(0))
.unwrap();
store.flush_storage().unwrap();
let storage = store.storage.as_ref().unwrap().read();
let eviction_store = Arc::clone(&store);
let (finished_tx, finished_rx) = std::sync::mpsc::channel();
let eviction = std::thread::spawn(move || {
eviction_store.evict_tenant(TENANT);
finished_tx.send(()).unwrap();
});
let finished_before_release = finished_rx.recv_timeout(Duration::from_secs(1));
drop(storage);
eviction.join().unwrap();
assert!(finished_before_release.is_ok());
assert_eq!(store.total_events(), 1);
assert!(store.is_tenant_loaded(TENANT));
store.evict_tenant(TENANT);
assert_eq!(store.total_events(), 0);
assert!(!store.is_tenant_loaded(TENANT));
}
#[test]
fn resident_counter_counts_only_the_new_batch() {
let store = EventStore::new();
store.ingest_batch(vec![event(), event(), event()]).unwrap();
store.ingest_batch(vec![event(), event()]).unwrap();
assert_eq!(store.stats().total_ingested, 5);
assert_eq!(store.total_events(), 5);
}
#[tokio::test(flavor = "current_thread")]
async fn eviction_between_hydration_and_version_check_cannot_authorize_a_stale_write() {
let directory = TempDir::new().unwrap();
{
let store = EventStore::with_config(EventStoreConfig::with_persistence(directory.path()));
store
.ingest_with_expected_version(&event(), Some(0))
.unwrap();
store.flush_storage().unwrap();
}
let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
directory.path(),
)));
let gate_store = Arc::clone(&store);
let (held_tx, held_rx) = tokio::sync::oneshot::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let holder = std::thread::spawn(move || {
let _gate = gate_store.durability_gate.write();
held_tx.send(()).unwrap();
release_rx.recv_timeout(Duration::from_secs(3))
});
held_rx.await.unwrap();
let writer_store = Arc::clone(&store);
let writer = tokio::task::spawn_blocking(move || {
writer_store.ingest_with_expected_version(&event(), Some(0))
});
tokio::time::timeout(Duration::from_secs(1), async {
while !store.tenant_loader.is_complete(TENANT) {
tokio::time::sleep(Duration::from_millis(1)).await;
}
})
.await
.unwrap();
let eviction_store = Arc::clone(&store);
let (evicted_tx, mut evicted_rx) = tokio::sync::oneshot::channel();
let eviction = tokio::task::spawn_blocking(move || {
eviction_store.evict_tenant(TENANT);
let _ = evicted_tx.send(());
});
let _ = tokio::time::timeout(Duration::from_millis(100), &mut evicted_rx).await;
release_tx.send(()).unwrap();
let result = writer.await.unwrap();
eviction.await.unwrap();
holder.join().unwrap().unwrap();
assert!(
result.is_err(),
"eviction authorized a stale write: {result:?}"
);
assert_eq!(
store
.query(&QueryEventsRequest {
tenant_id: Some(TENANT.into()),
..Default::default()
})
.unwrap()
.len(),
1
);
}