use std::time::Duration;
use futures::StreamExt;
use tokio::time::timeout;
use eventuary_core::io::ConsumerGroupId;
use eventuary_core::io::cursor::{CursorKind, CursorOrder, EncodedCursor};
use eventuary_core::io::filter::EventFilter;
use eventuary_core::io::reader::{
CheckpointKey, CheckpointReader, CheckpointScope, CheckpointStore, CheckpointSubscription,
PartitionedCursor, PartitionedReader, PartitionedReaderConfig, PartitionedSubscription,
};
use eventuary_core::io::{CursorId, Reader, StreamId, Writer};
use eventuary_core::{Event, OrganizationId, Payload, StartFrom, StopAt};
use eventuary_sqlite::checkpoint_store::{SqliteCheckpointStore, SqliteCheckpointStoreConfig};
use eventuary_sqlite::database::{SqliteConn, SqliteDatabase};
use eventuary_sqlite::reader::{
SqliteCursor, SqliteReader, SqliteReaderConfig, SqliteSubscription,
};
use eventuary_sqlite::writer::{SqliteWriter, SqliteWriterConfig};
fn prepare_test_schema(conn: &SqliteConn) {
SqliteWriter::prepare_schema(conn, &SqliteWriterConfig::default()).unwrap();
SqliteCheckpointStore::<SqliteCursor>::prepare_schema(
conn,
&SqliteCheckpointStoreConfig::default(),
)
.unwrap();
}
fn ev(org: &str, ns: &str, topic: &str, key: &str) -> Event {
Event::builder(org, ns, topic, key, Payload::from_string("p"))
.unwrap()
.build()
.expect("valid event")
}
fn fast_config() -> SqliteReaderConfig {
SqliteReaderConfig {
poll_interval: Duration::from_millis(10),
..SqliteReaderConfig::default()
}
}
fn sub_for(org: &str) -> SqliteSubscription {
SqliteSubscription {
start: StartFrom::Earliest,
stop_at: StopAt::Never,
filter: EventFilter::for_organization(OrganizationId::new(org).unwrap()),
batch_size: Some(10),
limit: None,
..SqliteSubscription::default()
}
}
fn scope() -> CheckpointScope {
CheckpointScope::new(
ConsumerGroupId::new("workers").unwrap(),
StreamId::new("billing").unwrap(),
)
}
#[tokio::test]
async fn checkpoint_reader_over_sqlite_resumes_after_ack() {
let db = SqliteDatabase::open_in_memory().unwrap();
prepare_test_schema(&db.conn());
let writer = SqliteWriter::new(db.conn());
for i in 0..3 {
writer
.write(&ev("acme", "/x", "thing.happened", &format!("k{i}")))
.await
.unwrap();
}
let store = SqliteCheckpointStore::<SqliteCursor>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let source = SqliteReader::new(db.conn(), fast_config());
let checkpointed = CheckpointReader::new(source, store);
let mut stream = checkpointed
.read(CheckpointSubscription::new(sub_for("acme"), scope()))
.await
.unwrap();
let m0 = timeout(Duration::from_secs(5), stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(m0.event().key().as_str(), "k0");
m0.ack().await.unwrap();
let m1 = timeout(Duration::from_secs(5), stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(m1.event().key().as_str(), "k1");
m1.ack().await.unwrap();
drop(stream);
let source2 = SqliteReader::new(db.conn(), fast_config());
let store2 = SqliteCheckpointStore::<SqliteCursor>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let checkpointed2 = CheckpointReader::new(source2, store2);
let mut stream2 = checkpointed2
.read(CheckpointSubscription::new(sub_for("acme"), scope()))
.await
.unwrap();
let next = timeout(Duration::from_secs(5), stream2.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(next.event().key().as_str(), "k2");
}
#[tokio::test]
async fn checkpoint_over_partitioned_sqlite_stores_per_lane_offsets() {
let db = SqliteDatabase::open_in_memory().unwrap();
prepare_test_schema(&db.conn());
let writer = SqliteWriter::new(db.conn());
for i in 0..6 {
writer
.write(&ev("acme", "/x", "thing.happened", &format!("k{i}")))
.await
.unwrap();
}
let source = SqliteReader::new(db.conn(), fast_config());
let partitioned = PartitionedReader::source(
source,
PartitionedReaderConfig {
partition_count: std::num::NonZeroU16::new(4).unwrap(),
..PartitionedReaderConfig::default()
},
);
let store = SqliteCheckpointStore::<PartitionedCursor<SqliteCursor>>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let checkpointed = CheckpointReader::new(partitioned, store);
let inner = PartitionedSubscription::new(sub_for("acme"));
let mut stream = checkpointed
.read(CheckpointSubscription::new(inner, scope()))
.await
.unwrap();
for _ in 0..6 {
let msg = timeout(Duration::from_secs(5), stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
msg.ack().await.unwrap();
}
drop(stream);
let store2 = SqliteCheckpointStore::<PartitionedCursor<SqliteCursor>>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let rows = store2.load_scope(&scope()).await.unwrap();
assert!(!rows.is_empty(), "expected per-lane checkpoints persisted");
for (cursor_id, _cursor) in &rows {
assert!(
cursor_id.as_str().starts_with("partition:"),
"partitioned cursor must be tagged with a named cursor id"
);
}
}
#[tokio::test]
async fn checkpoint_reader_no_advance_on_nack() {
let db = SqliteDatabase::open_in_memory().unwrap();
prepare_test_schema(&db.conn());
let writer = SqliteWriter::new(db.conn());
writer
.write(&ev("acme", "/x", "thing.happened", "k0"))
.await
.unwrap();
let source = SqliteReader::new(db.conn(), fast_config());
let store = SqliteCheckpointStore::<SqliteCursor>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let checkpointed = CheckpointReader::new(source, store);
let mut stream = checkpointed
.read(CheckpointSubscription::new(sub_for("acme"), scope()))
.await
.unwrap();
let m0 = timeout(Duration::from_secs(5), stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
m0.nack().await.unwrap();
drop(stream);
let store2 = SqliteCheckpointStore::<SqliteCursor>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let rows = store2.load_scope(&scope()).await.unwrap();
assert!(rows.is_empty(), "nack must not commit checkpoint");
}
#[tokio::test]
async fn partitioned_reader_tags_partition_on_cursor() {
let db = SqliteDatabase::open_in_memory().unwrap();
prepare_test_schema(&db.conn());
let writer = SqliteWriter::new(db.conn());
for i in 0..8 {
writer
.write(&ev("acme", "/x", "thing.happened", &format!("k{i}")))
.await
.unwrap();
}
let source = SqliteReader::new(db.conn(), fast_config());
let partitioned = PartitionedReader::source(
source,
PartitionedReaderConfig {
partition_count: std::num::NonZeroU16::new(4).unwrap(),
..PartitionedReaderConfig::default()
},
);
let mut stream = partitioned
.read(PartitionedSubscription::new(sub_for("acme")))
.await
.unwrap();
let mut delivered = 0usize;
while delivered < 8 {
let msg = timeout(Duration::from_secs(5), stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert!(msg.cursor().partition().id() < 4);
msg.ack().await.unwrap();
delivered += 1;
}
assert_eq!(delivered, 8);
}
#[tokio::test]
async fn sqlite_checkpoint_store_persists_encoded_cursor() {
let db = SqliteDatabase::open_in_memory().unwrap();
prepare_test_schema(&db.conn());
let store: SqliteCheckpointStore<EncodedCursor> =
SqliteCheckpointStore::new(db.conn(), SqliteCheckpointStoreConfig::default());
let key = CheckpointKey::new(scope(), CursorId::global());
let cursor = EncodedCursor::from_json(
CursorId::global(),
CursorKind::new("eventuary.test.cursor.v1").unwrap(),
CursorOrder::from_u64(7),
&serde_json::json!({ "sequence": 7 }),
)
.unwrap();
store.commit(&key, cursor.clone()).await.unwrap();
let loaded = store.load(&key).await.unwrap().unwrap();
assert_eq!(loaded.kind().as_str(), "eventuary.test.cursor.v1");
assert_eq!(loaded.id_ref(), &CursorId::global());
assert_eq!(loaded.order(), &CursorOrder::from_u64(7));
assert_eq!(loaded, cursor);
}
#[tokio::test]
async fn checkpoint_over_partitioned_resumes_and_skips_acked_events() {
let db = SqliteDatabase::open_in_memory().unwrap();
prepare_test_schema(&db.conn());
let writer = SqliteWriter::new(db.conn());
for i in 0..8 {
writer
.write(&ev("acme", "/x", "thing.happened", &format!("k{i}")))
.await
.unwrap();
}
let scp = scope();
let mut acked_keys: Vec<String> = Vec::new();
{
let source = SqliteReader::new(db.conn(), fast_config());
let partitioned = PartitionedReader::source(
source,
PartitionedReaderConfig {
partition_count: std::num::NonZeroU16::new(4).unwrap(),
..PartitionedReaderConfig::default()
},
);
let store = SqliteCheckpointStore::<PartitionedCursor<SqliteCursor>>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let checkpointed = CheckpointReader::new(partitioned, store);
let inner = PartitionedSubscription::new(sub_for("acme"));
let mut stream = checkpointed
.read(CheckpointSubscription::new(inner, scp.clone()))
.await
.unwrap();
for _ in 0..4 {
let msg = timeout(Duration::from_secs(5), stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
acked_keys.push(msg.event().key().as_str().to_owned());
msg.ack().await.unwrap();
}
}
let mut resumed_keys: Vec<String> = Vec::new();
{
let source2 = SqliteReader::new(db.conn(), fast_config());
let partitioned2 = PartitionedReader::source(
source2,
PartitionedReaderConfig {
partition_count: std::num::NonZeroU16::new(4).unwrap(),
..PartitionedReaderConfig::default()
},
);
let store2 = SqliteCheckpointStore::<PartitionedCursor<SqliteCursor>>::new(
db.conn(),
SqliteCheckpointStoreConfig::default(),
);
let checkpointed2 = CheckpointReader::new(partitioned2, store2);
let inner2 = PartitionedSubscription::new(sub_for("acme"));
let mut stream2 = checkpointed2
.read(CheckpointSubscription::new(inner2, scp.clone()))
.await
.unwrap();
for _ in 0..4 {
let msg = timeout(Duration::from_secs(5), stream2.next())
.await
.unwrap()
.unwrap()
.unwrap();
resumed_keys.push(msg.event().key().as_str().to_owned());
msg.ack().await.unwrap();
}
}
let acked_set: std::collections::HashSet<&String> = acked_keys.iter().collect();
let resumed_set: std::collections::HashSet<&String> = resumed_keys.iter().collect();
let overlap: std::collections::HashSet<_> = acked_set.intersection(&resumed_set).collect();
assert!(
overlap.is_empty(),
"resumed should not re-deliver acked events; overlap={overlap:?}"
);
let mut combined: std::collections::HashSet<String> = acked_keys.into_iter().collect();
combined.extend(resumed_keys);
assert_eq!(
combined.len(),
8,
"combined total should be 8; got {}",
combined.len()
);
}