eventuary-sqlite 0.3.0-rc.1

SQLite event backend for eventuary
use std::num::NonZeroU32;
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::partition::{EventKeyPartitionKeyResolver, Fnv1a64PartitionHasher};
use eventuary_core::{Event, OrganizationId, Payload, StartFrom, StopAt};
use eventuary_sqlite::checkpoint::{SqliteCheckpointStore, SqliteCheckpointStoreConfig};
use eventuary_sqlite::database::{SqliteConn, SqliteDatabase};
use eventuary_sqlite::reader::{
    SqliteCursor, SqliteReader, SqliteReaderConfig, SqliteSubscription,
};
use eventuary_sqlite::writer::{SqlitePartitioningConfig, SqliteWriter, SqliteWriterConfig};

fn writer_config() -> SqliteWriterConfig {
    SqliteWriterConfig {
        partitioning: SqlitePartitioningConfig::inline(
            NonZeroU32::new(4).unwrap(),
            EventKeyPartitionKeyResolver::new(),
            Fnv1a64PartitionHasher,
        ),
        ..SqliteWriterConfig::default()
    }
}

fn prepare_test_schema(conn: &SqliteConn) {
    SqliteWriter::prepare_schema(conn, &writer_config()).unwrap();
    SqliteCheckpointStore::<SqliteCursor>::prepare_schema(
        conn,
        &SqliteCheckpointStoreConfig::default(),
    )
    .unwrap();
}

fn make_writer(conn: SqliteConn) -> SqliteWriter {
    SqliteWriter::new_with_config(conn, writer_config())
}

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 = make_writer(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 = make_writer(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::NonZeroU32::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 = make_writer(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");
}

// Contiguous-delivered-order ack semantics are unit-tested in
// `eventuary-core` against `PendingState` directly. An end-to-end test
// against a real source reader would require multiple in-flight messages
// per partition, but the current readers (source-cursor SQL readers and
// the lane scheduler) hold at most one in-flight per lane, so the
// scenario cannot be constructed without a synthetic test reader.

#[tokio::test]
async fn partitioned_reader_tags_partition_on_cursor() {
    let db = SqliteDatabase::open_in_memory().unwrap();
    prepare_test_schema(&db.conn());
    let writer = make_writer(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::NonZeroU32::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 = make_writer(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::NonZeroU32::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::NonZeroU32::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()
    );
}