khive-db 0.10.0

SQLite storage backend: entities, edges, notes, events, FTS5, sqlite-vec vectors.
Documentation
use super::*;
use rusqlite::types::Value as SqlValue;

const ACK_ROW_SQL: &str = concat!(
    "SELECT delivery_attempt_id,sender_agent_id,logical_message_id,binding,disposition,",
    "state,created_at,updated_at,attempt_count,not_before,retirement_reason ",
    "FROM comm_ack_work WHERE delivery_attempt_id=?1",
);

fn acknowledgement_row(backend: &StorageBackend, attempt: Uuid) -> Option<Vec<SqlValue>> {
    backend
        .pool()
        .writer()
        .unwrap()
        .conn()
        .query_row(ACK_ROW_SQL, [attempt.to_string()], |row| {
            (0..11).map(|column| row.get(column)).collect()
        })
        .optional()
        .unwrap()
}

fn receipt_row(backend: &StorageBackend) -> Vec<SqlValue> {
    backend
        .pool()
        .writer()
        .unwrap()
        .conn()
        .query_row("SELECT * FROM comm_recipient_replay", [], |row| {
            (0..row.as_ref().column_count())
                .map(|column| row.get(column))
                .collect()
        })
        .unwrap()
}

#[tokio::test]
async fn due_acknowledgement_preserves_binding_and_finish_is_idempotent() {
    let (backend, input) = fixture();
    let store = RecipientTransportStore::new(backend.pool_arc());
    store.commit(input.clone()).await.unwrap();
    let entries = store.list_due_acknowledgements(i64::MAX, 10).await.unwrap();
    assert_eq!(
        entries.len(),
        1,
        "one committed delivery must have one due entry"
    );
    assert_eq!(entries[0].delivery_attempt_id, input.delivery_attempt_id);
    assert_eq!(
        entries[0].binding, input.binding,
        "journal must retain exact binding"
    );
    assert_eq!(entries[0].disposition, RecipientDisposition::Stored);
    assert_eq!(entries[0].attempt_count, 0);
    assert_eq!(entries[0].not_before, None);

    assert!(store
        .finish_acknowledgement(input.delivery_attempt_id)
        .await
        .unwrap());
    assert!(
        store
            .list_due_acknowledgements(i64::MAX, 10)
            .await
            .unwrap()
            .is_empty(),
        "finished acknowledgement must not be listed"
    );
    let finished = acknowledgement_row(&backend, input.delivery_attempt_id);
    assert!(!store
        .finish_acknowledgement(input.delivery_attempt_id)
        .await
        .unwrap());
    assert_eq!(
        acknowledgement_row(&backend, input.delivery_attempt_id),
        finished,
        "a second finish must not update even its timestamp"
    );
    assert!(!store.finish_acknowledgement(Uuid::new_v4()).await.unwrap());
}

#[tokio::test]
async fn acknowledgement_retry_waits_until_not_before_and_counts_failed_tries() {
    let (backend, input) = fixture();
    let store = RecipientTransportStore::new(backend.pool_arc());
    store.commit(input.clone()).await.unwrap();
    let now = chrono::Utc::now().timestamp_micros();
    let not_before = now + 1_000_000;
    assert!(store
        .record_acknowledgement_failed_try(input.delivery_attempt_id, not_before)
        .await
        .unwrap());
    assert!(
        store
            .list_due_acknowledgements(now, 10)
            .await
            .unwrap()
            .is_empty(),
        "an acknowledgement before its not-before time must not be due"
    );
    let entries = store
        .list_due_acknowledgements(not_before, 10)
        .await
        .unwrap();
    assert_eq!(
        entries.len(),
        1,
        "an acknowledgement is due exactly at not-before"
    );
    assert_eq!(
        entries[0].attempt_count, 1,
        "one failed try must increment the counter"
    );
    assert_eq!(entries[0].not_before, Some(not_before));
    assert_eq!(entries[0].binding, input.binding);
    assert_eq!(entries[0].disposition, RecipientDisposition::Stored);
}

#[tokio::test]
async fn due_acknowledgements_are_oldest_first_bounded_and_indexed() {
    let (backend, first) = fixture();
    let (_, second) = fixture();
    let store = RecipientTransportStore::new(backend.pool_arc());
    store.commit(first.clone()).await.unwrap();
    store.commit(second.clone()).await.unwrap();
    {
        let guard = backend.pool().writer().unwrap();
        for (id, created_at) in [
            (first.delivery_attempt_id, 20),
            (second.delivery_attempt_id, 10),
        ] {
            guard
                .conn()
                .execute(
                    "UPDATE comm_ack_work SET created_at=?1 WHERE delivery_attempt_id=?2",
                    params![created_at, id.to_string()],
                )
                .unwrap();
        }
    }
    let one = store.list_due_acknowledgements(i64::MAX, 1).await.unwrap();
    assert_eq!(one.len(), 1, "due listing must obey the requested limit");
    assert_eq!(one[0].delivery_attempt_id, second.delivery_attempt_id);
    let both = store.list_due_acknowledgements(i64::MAX, 2).await.unwrap();
    assert_eq!(
        both.iter()
            .map(|entry| entry.delivery_attempt_id)
            .collect::<Vec<_>>(),
        vec![second.delivery_attempt_id, first.delivery_attempt_id],
        "due listing must return the oldest entry first"
    );
    assert!(store
        .list_due_acknowledgements(i64::MAX, 0)
        .await
        .unwrap()
        .is_empty());
    backend
        .pool()
        .writer()
        .unwrap()
        .conn()
        .execute("UPDATE comm_ack_work SET created_at=10", [])
        .unwrap();
    let tied = store.list_due_acknowledgements(i64::MAX, 2).await.unwrap();
    let mut expected = [first.delivery_attempt_id, second.delivery_attempt_id];
    expected.sort();
    assert_eq!(
        tied.iter()
            .map(|entry| entry.delivery_attempt_id)
            .collect::<Vec<_>>(),
        expected,
        "equal timestamps must have a stable attempt-identifier order"
    );
    let guard = backend.pool().writer().unwrap();
    let mut statement = guard
        .conn()
        .prepare(&format!("EXPLAIN QUERY PLAN {ACK_DUE_SQL}"))
        .unwrap();
    let plan = statement
        .query_map(params![i64::MAX, 2], |row| row.get::<_, String>(3))
        .unwrap()
        .collect::<rusqlite::Result<Vec<_>>>()
        .unwrap();
    assert!(
        plan.iter().any(|line| line.contains("idx_comm_ack_due")),
        "the production due query must use the due index: {plan:?}"
    );
    assert!(
        !plan.iter().any(|line| line.contains("USE TEMP B-TREE")),
        "the due index must provide the listing order: {plan:?}"
    );
}

#[tokio::test]
async fn acknowledged_and_retired_attempts_remain_unchanged_on_replay() {
    for retired in [false, true] {
        let (backend, input) = fixture();
        let store = RecipientTransportStore::new(backend.pool_arc());
        store.commit(input.clone()).await.unwrap();
        store
            .record_acknowledgement_failed_try(input.delivery_attempt_id, 99)
            .await
            .unwrap();
        if retired {
            assert!(store
                .retire_acknowledgement(
                    input.delivery_attempt_id,
                    AcknowledgementRetirementReason::PermanentTransport,
                )
                .await
                .unwrap());
        } else {
            assert!(store
                .finish_acknowledgement(input.delivery_attempt_id)
                .await
                .unwrap());
        }
        let terminal = acknowledgement_row(&backend, input.delivery_attempt_id);
        assert!(
            terminal.is_some(),
            "terminal acknowledgement must be retained"
        );
        assert!(!store.commit(input.clone()).await.unwrap().created);
        assert_eq!(
            acknowledgement_row(&backend, input.delivery_attempt_id),
            terminal,
            "replay must preserve terminal state, retry counter and timestamps"
        );
        assert!(store
            .list_due_acknowledgements(i64::MAX, 10)
            .await
            .unwrap()
            .is_empty());
        assert!(!store
            .record_acknowledgement_failed_try(input.delivery_attempt_id, 100)
            .await
            .unwrap());
        assert_eq!(
            acknowledgement_row(&backend, input.delivery_attempt_id),
            terminal,
            "a failed-try call must not mutate a terminal row"
        );
    }
}

#[tokio::test]
async fn retiring_acknowledgement_keeps_reason_binding_message_and_receipt() {
    let (backend, input) = fixture();
    let store = RecipientTransportStore::new(backend.pool_arc());
    let committed = store.commit(input.clone()).await.unwrap();
    let note_store = SqlNoteStore::new(backend.pool_arc(), false);
    let note_id = committed.note_id.unwrap();
    let note_before = note_store.get_note(note_id).await.unwrap();
    let receipt_before = receipt_row(&backend);
    let original = acknowledgement_row(&backend, input.delivery_attempt_id).unwrap();
    assert!(store
        .retire_acknowledgement(
            input.delivery_attempt_id,
            AcknowledgementRetirementReason::PermanentTransport,
        )
        .await
        .unwrap());
    let retained = acknowledgement_row(&backend, input.delivery_attempt_id);
    assert!(
        retained.is_some(),
        "retirement must retain the acknowledgement row"
    );
    let retained = retained.unwrap();
    assert_eq!(
        &retained[..5],
        &original[..5],
        "retirement must keep identity and binding"
    );
    assert_eq!(retained[5], SqlValue::Text("retired".into()));
    assert_eq!(retained[10], SqlValue::Text("permanent_transport".into()));
    assert!(
        store
            .list_due_acknowledgements(i64::MAX, 10)
            .await
            .unwrap()
            .is_empty(),
        "retired acknowledgements must never be retried"
    );
    assert_eq!(note_store.get_note(note_id).await.unwrap(), note_before);
    assert_eq!(
        receipt_row(&backend),
        receipt_before,
        "retirement must retain the receipt"
    );
    assert!(!store
        .retire_acknowledgement(
            input.delivery_attempt_id,
            AcknowledgementRetirementReason::PermanentTransport,
        )
        .await
        .unwrap());
    assert_eq!(
        acknowledgement_row(&backend, input.delivery_attempt_id),
        Some(retained)
    );
    assert!(!store
        .finish_acknowledgement(input.delivery_attempt_id)
        .await
        .unwrap());
}

#[tokio::test]
async fn acknowledgement_retry_bookkeeping_survives_file_backed_restart() {
    let dir = tempfile::tempdir().unwrap();
    let path = dir.path().join("acknowledgements.db");
    let (_, input) = fixture();
    let not_before = chrono::Utc::now().timestamp_micros() + 1_000_000;
    let pending_before;
    {
        let backend = StorageBackend::sqlite_for_test(&path).unwrap();
        backend.pool().run_migrations().unwrap();
        let store = RecipientTransportStore::new(backend.pool_arc());
        store.commit(input.clone()).await.unwrap();
        store
            .record_acknowledgement_failed_try(input.delivery_attempt_id, not_before)
            .await
            .unwrap();
        pending_before = acknowledgement_row(&backend, input.delivery_attempt_id);
    }
    let reopened = StorageBackend::sqlite_for_test(&path).unwrap();
    reopened.pool().run_migrations().unwrap();
    let store = RecipientTransportStore::new(reopened.pool_arc());
    assert_eq!(
        acknowledgement_row(&reopened, input.delivery_attempt_id),
        pending_before,
        "reopening must preserve the complete pending journal row"
    );
    let due = store
        .list_due_acknowledgements(not_before, 10)
        .await
        .unwrap();
    assert_eq!(due.len(), 1);
    assert_eq!(
        due[0].attempt_count, 1,
        "the failed-try counter must survive restart"
    );
    assert_eq!(
        due[0].not_before,
        Some(not_before),
        "not-before must survive restart"
    );
    assert_eq!(due[0].binding, input.binding);
    assert_eq!(due[0].disposition, RecipientDisposition::Stored);
    assert!(store
        .list_due_acknowledgements(not_before - 1, 10)
        .await
        .unwrap()
        .is_empty());
}

#[tokio::test]
async fn acknowledgement_failed_try_counter_overflow_keeps_the_row_unchanged() {
    let (backend, input) = fixture();
    let store = RecipientTransportStore::new(backend.pool_arc());
    store.commit(input.clone()).await.unwrap();
    backend
        .pool()
        .writer()
        .unwrap()
        .conn()
        .execute(
            "UPDATE comm_ack_work SET attempt_count=?1 WHERE delivery_attempt_id=?2",
            params![i64::MAX, input.delivery_attempt_id.to_string()],
        )
        .unwrap();
    let before = acknowledgement_row(&backend, input.delivery_attempt_id);
    assert!(
        store
            .record_acknowledgement_failed_try(input.delivery_attempt_id, 100)
            .await
            .is_err(),
        "retry counter overflow must be refused"
    );
    assert_eq!(
        acknowledgement_row(&backend, input.delivery_attempt_id),
        before,
        "a refused counter overflow must not change the retry time or timestamp"
    );
}