khive-db 0.10.0

SQLite storage backend: entities, edges, notes, events, FTS5, sqlite-vec vectors.
Documentation
use super::*;
use crate::pool::PoolConfig;
use khive_storage::usage::{scope, UsageContext};
use serde_json::json;
use std::sync::atomic::Ordering;

struct Fixture {
    store: SqlEventStore,
    pool: Arc<ConnectionPool>,
    _directory: Option<tempfile::TempDir>,
}

impl Fixture {
    fn new(file_backed: bool) -> Self {
        let directory = file_backed.then(|| tempfile::tempdir().unwrap());
        let pool = Arc::new(
            ConnectionPool::new(PoolConfig {
                path: directory.as_ref().map(|dir| dir.path().join("events.db")),
                write_queue_enabled: Some(false),
                write_routing_strict: false,
                ..PoolConfig::for_test()
            })
            .unwrap(),
        );
        {
            let writer = pool.writer().unwrap();
            writer.conn().execute_batch(EVENTS_DDL).unwrap();
        }
        let store = SqlEventStore::new_scoped(Arc::clone(&pool), file_backed, "local");
        let fixture = Self {
            store,
            pool,
            _directory: directory,
        };
        fixture.assert_fallback();
        fixture
    }

    fn assert_fallback(&self) {
        assert!(!self.pool.write_queue_active());
        assert!(self.pool.writer_task_handle().unwrap().is_none());
    }

    fn lose_next_reply(&self) {
        self.store
            .lose_next_fallback_reply
            .store(true, Ordering::SeqCst);
    }

    async fn assert_committed(&self, events: &[Event]) {
        for event in events {
            assert_eq!(
                self.store.get_event(event.id).await.unwrap().as_ref(),
                Some(event),
                "the lost reply must follow the actual fallback commit"
            );
        }
        assert!(!self.store.lose_next_fallback_reply.load(Ordering::SeqCst));
        self.assert_fallback();
    }
}

fn event() -> Event {
    Event::new(
        "local",
        "search",
        EventKind::SearchExecuted,
        SubstrateKind::Note,
        "agent:test",
    )
    .with_payload(json!({"result_kind": "note"}))
}

fn assert_unknown(error: StorageError) {
    assert!(matches!(
        error,
        StorageError::WriterTaskTerminated {
            request_state: WriterTaskRequestState::SideEffectsUnknown,
        }
    ));
}

async fn single_append_lost_reply(file_backed: bool) {
    let fixture = Fixture::new(file_backed);
    let context = UsageContext::new();
    let known = event();
    scope(context.clone(), fixture.store.append_event(known.clone()))
        .await
        .unwrap();
    assert_eq!(context.shipping_snapshot(), Some(json!({"event_rows": 1})));

    fixture.lose_next_reply();
    let lost = event();
    let error = scope(context.clone(), fixture.store.append_event(lost.clone()))
        .await
        .unwrap_err();
    assert_unknown(error);
    fixture.assert_committed(&[known, lost]).await;
    assert_eq!(context.snapshot(), json!({"event_rows": 1}));
    assert_eq!(context.shipping_snapshot(), None);
}

async fn batch_append_lost_reply(file_backed: bool) {
    let fixture = Fixture::new(file_backed);
    let context = UsageContext::new();
    let known = vec![event(), event()];
    let summary = scope(context.clone(), fixture.store.append_events(known.clone()))
        .await
        .unwrap();
    assert_eq!(summary.attempted, 2);
    assert_eq!(summary.affected, 2);
    assert_eq!(context.shipping_snapshot(), Some(json!({"event_rows": 2})));

    fixture.lose_next_reply();
    let lost = vec![event(), event(), event()];
    let error = scope(context.clone(), fixture.store.append_events(lost.clone()))
        .await
        .unwrap_err();
    assert_unknown(error);
    fixture.assert_committed(&known).await;
    fixture.assert_committed(&lost).await;
    assert_eq!(context.snapshot(), json!({"event_rows": 2}));
    assert_eq!(context.shipping_snapshot(), None);
}

#[tokio::test]
async fn memory_single_append_lost_reply_withholds_usage_after_commit() {
    single_append_lost_reply(false).await;
}

#[tokio::test]
async fn standalone_single_append_lost_reply_withholds_usage_after_commit() {
    single_append_lost_reply(true).await;
}

#[tokio::test]
async fn memory_batch_append_lost_reply_withholds_usage_after_commit() {
    batch_append_lost_reply(false).await;
}

#[tokio::test]
async fn standalone_batch_append_lost_reply_withholds_usage_after_commit() {
    batch_append_lost_reply(true).await;
}