khive-db 0.10.0

SQLite storage backend: entities, edges, notes, events, FTS5, sqlite-vec vectors.
Documentation
use super::{
    blob_gc_unowned_attachment_predicate, blob_root_key, claim_blob_gc_batch,
    parse_blob_gc_claim_rows, release_blob_gc_batch, FsBlobStore,
};
use crate::StorageBackend;
use khive_storage::BlobStore;
use rusqlite::{params, Connection};

fn unowned_refs(conn: &Connection, candidates: &[&str], quarantine_present: bool) -> Vec<String> {
    let statement = format!(
        "SELECT candidate.value FROM json_each(?1) AS candidate \
         WHERE {} ORDER BY CAST(candidate.key AS INTEGER)",
        blob_gc_unowned_attachment_predicate(quarantine_present)
    );
    let encoded = serde_json::to_string(candidates).expect("encode real candidate refs");
    let mut query = conn
        .prepare(&statement)
        .expect("prepare the production ownership predicate");
    let rows = query
        .query_map([encoded], |row| row.get::<_, String>(0))
        .expect("execute the production ownership predicate");
    rows.collect::<rusqlite::Result<Vec<_>>>()
        .expect("read the production ownership decisions")
}

#[tokio::test]
async fn quarantine_liveness_predicate_keeps_quarantined_blob_live() {
    let dir = tempfile::tempdir().expect("create an isolated object root");
    let store = FsBlobStore::new(dir.path().join("blobs"), 0).expect("open the object store");
    let canonical_bytes = b"canonical attachment owner";
    let quarantined_bytes = b"quarantined attachment owner";
    let orphan_bytes = b"independent unowned object";
    let canonical = store
        .put(canonical_bytes.to_vec())
        .await
        .expect("publish the canonical owner's object");
    let quarantined = store
        .put(quarantined_bytes.to_vec())
        .await
        .expect("publish the quarantined owner's object");
    let orphan = store
        .put(orphan_bytes.to_vec())
        .await
        .expect("publish the independent orphan");
    assert_ne!(canonical, quarantined);
    assert_ne!(canonical, orphan);
    assert_ne!(quarantined, orphan);
    for content_ref in [&canonical, &quarantined, &orphan] {
        assert!(
            store
                .exists(content_ref)
                .await
                .expect("check the real object"),
            "every candidate must name a published object"
        );
    }

    let conn = Connection::open_in_memory().expect("open an isolated ownership database");
    conn.execute_batch(
        "CREATE TABLE blob_gc_claims (
             root_key TEXT NOT NULL,
             content_ref TEXT NOT NULL,
             claimed_at INTEGER NOT NULL,
             PRIMARY KEY (root_key, content_ref)
         ) STRICT;",
    )
    .expect("create the claim table required by the migration fences");
    conn.execute_batch(include_str!("../../../sql/021-attachments-a-stage.sql"))
        .expect("apply the real attachment schema");
    conn.execute_batch(include_str!(
        "../../../sql/021-attachments-b-claim-fences.sql"
    ))
    .expect("apply the real attachment claim fences");

    let legacy_role = "message-attachment:0\u{85}";
    conn.execute(
        "INSERT INTO attachments
             (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at)
         VALUES (?1, 'note', 'message-attachment:0', ?2, NULL, ?3, 1)",
        params![
            "11111111-1111-4111-8111-111111111111",
            canonical.as_str(),
            canonical_bytes.len() as i64
        ],
    )
    .expect("register the canonical owner");
    conn.execute(
        "INSERT INTO attachments
             (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at)
         VALUES (?1, 'note', ?2, ?3, NULL, ?4, 2)",
        params![
            "22222222-2222-4222-8222-222222222222",
            legacy_role,
            quarantined.as_str(),
            quarantined_bytes.len() as i64
        ],
    )
    .expect("register a role accepted by the historical schema");

    let quarantine_tables: i64 = conn
        .query_row(
            "SELECT COUNT(*) FROM sqlite_master
             WHERE type = 'table' AND name = 'attachment_quarantine'",
            [],
            |row| row.get(0),
        )
        .expect("inspect the historical table set");
    assert_eq!(
        quarantine_tables, 0,
        "the false branch must use no quarantine table"
    );
    let candidates = [quarantined.as_str(), orphan.as_str(), canonical.as_str()];
    assert_eq!(
        unowned_refs(&conn, &candidates, false),
        vec![orphan.to_string()],
        "the historical production predicate must keep both attachment owners live"
    );

    conn.execute_batch(include_str!(
        "../../../sql/047-attachment-role-quarantine.sql"
    ))
    .expect("migrate the actual rejected role into durable quarantine ownership");
    let retained: (String, String, String, i64) = conn
        .query_row(
            "SELECT role, content_ref, reason, size_bytes FROM attachment_quarantine
             WHERE record_uuid = ?1",
            ["22222222-2222-4222-8222-222222222222"],
            |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
        )
        .expect("read the migrated owner independently of the liveness predicate");
    assert_eq!(
        retained,
        (
            legacy_role.to_string(),
            quarantined.to_string(),
            "invalid_role".to_string(),
            quarantined_bytes.len() as i64
        ),
        "the migration must preserve the rejected role and its object ref"
    );
    let attachment_refs: Vec<String> = conn
        .prepare("SELECT content_ref FROM attachments ORDER BY record_uuid")
        .expect("prepare the canonical owner census")
        .query_map([], |row| row.get(0))
        .expect("read the canonical owner census")
        .collect::<rusqlite::Result<_>>()
        .expect("collect the canonical owner census");
    assert_eq!(attachment_refs, vec![canonical.to_string()]);

    assert!(
        unowned_refs(&conn, &[canonical.as_str()], true).is_empty(),
        "the production predicate must keep the canonical attachment live"
    );
    assert_eq!(
        unowned_refs(&conn, &[orphan.as_str()], true),
        vec![orphan.to_string()],
        "the independent real orphan must remain collectible by the predicate"
    );
    assert!(
        unowned_refs(&conn, &[quarantined.as_str()], true).is_empty(),
        "the production predicate must keep the quarantined object live"
    );
    assert_eq!(
        unowned_refs(&conn, &candidates, true),
        vec![orphan.to_string()],
        "only the independent orphan may be selected after quarantine migration"
    );
}

#[tokio::test]
async fn quarantine_ownership_reaches_the_actual_gc_batch_claim_and_dry_run() {
    let dir = tempfile::tempdir().unwrap();
    let store = FsBlobStore::new(dir.path().join("blobs"), 0).unwrap();
    let canonical = store.put(b"canonical owner".to_vec()).await.unwrap();
    let quarantined = store.put(b"quarantine owner".to_vec()).await.unwrap();
    let orphan = store.put(b"real orphan".to_vec()).await.unwrap();
    let grace = store.put(b"grace orphan".to_vec()).await.unwrap();
    let backend = StorageBackend::memory().expect("isolated batch database");
    backend
        .prepare_core_schema()
        .expect("actual complete core schema");
    let owner = "33333333-3333-4333-8333-333333333333";
    let quarantine_owner = "44444444-4444-4444-8444-444444444444";
    let role = "quarantine-original\u{85}";
    {
        let writer = backend.pool().writer().unwrap();
        writer.conn().execute(
            "INSERT INTO attachments (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at) VALUES (?1, 'note', 'message-attachment:0', ?2, NULL, 15, 1)",
            params![owner, canonical.as_str()],
        ).unwrap();
        writer.conn().execute(
            "INSERT INTO attachment_quarantine (record_uuid, substrate, role, content_ref, media_type, size_bytes, created_at, reason) VALUES (?1, 'note', ?2, ?3, NULL, 16, 2, 'invalid_role')",
            params![quarantine_owner, role, quarantined.as_str()],
        ).unwrap();
    }
    let snapshot = || {
        let writer = backend.pool().writer().unwrap();
        let canonical_row: (String, String, String) = writer
            .conn()
            .query_row(
                "SELECT record_uuid, role, content_ref FROM attachments",
                [],
                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
            )
            .unwrap();
        let quarantine_row: (String, String, String, String) = writer
            .conn()
            .query_row(
                "SELECT record_uuid, role, content_ref, reason FROM attachment_quarantine",
                [],
                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
            )
            .unwrap();
        (canonical_row, quarantine_row)
    };
    let before = snapshot();
    assert_eq!(
        before.0,
        (
            owner.to_string(),
            "message-attachment:0".into(),
            canonical.to_string()
        )
    );
    assert_eq!(
        before.1,
        (
            quarantine_owner.to_string(),
            role.into(),
            quarantined.to_string(),
            "invalid_role".into()
        )
    );
    let sql = backend.sql();
    let root_key = blob_root_key(store.root());
    let candidates = vec![
        (quarantined.clone(), false),
        (orphan.clone(), false),
        (canonical.clone(), false),
        (grace.clone(), true),
    ];
    // Exercise the internal claim seam without weakening public sweep admission.
    let dry_run = claim_blob_gc_batch(sql.as_ref(), root_key.clone(), &candidates, true)
        .await
        .expect("actual quarantine-aware dry run");
    assert_eq!(
        dry_run.would_delete, 1,
        "only the independent orphan is eligible"
    );
    assert_eq!(dry_run.grace_period_skipped, 1);
    assert!(dry_run.claimed_rows.is_empty());
    let claims: i64 = backend
        .pool()
        .writer()
        .unwrap()
        .conn()
        .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
        .unwrap();
    assert_eq!(claims, 0, "dry run writes no claims");
    assert_eq!(snapshot(), before);

    let claimed = claim_blob_gc_batch(sql.as_ref(), root_key.clone(), &candidates, false)
        .await
        .expect("actual quarantine-aware claim");
    assert_eq!(claimed.would_delete, 1);
    assert_eq!(claimed.grace_period_skipped, 1);
    assert_eq!(
        parse_blob_gc_claim_rows(claimed.claimed_rows).unwrap(),
        vec![orphan]
    );
    assert_eq!(snapshot(), before, "claiming preserves both owners exactly");
    for (reference, bytes) in [
        (&canonical, b"canonical owner".as_slice()),
        (&quarantined, b"quarantine owner".as_slice()),
        (&grace, b"grace orphan".as_slice()),
    ] {
        assert_eq!(
            store.get_bounded_verified(reference, 64).await.unwrap(),
            bytes
        );
    }
    release_blob_gc_batch(sql.as_ref(), root_key.clone())
        .await
        .unwrap();
    {
        let writer = backend.pool().writer().unwrap();
        assert_eq!(
            writer
                .conn()
                .execute(
                    "DELETE FROM attachment_quarantine WHERE record_uuid = ?1",
                    [quarantine_owner],
                )
                .unwrap(),
            1
        );
    }
    let released = claim_blob_gc_batch(
        sql.as_ref(),
        root_key.clone(),
        &[(quarantined.clone(), false)],
        false,
    )
    .await
    .expect("detached quarantine ownership becomes claimable");
    assert_eq!(released.would_delete, 1);
    assert_eq!(
        parse_blob_gc_claim_rows(released.claimed_rows).unwrap(),
        vec![quarantined]
    );
    release_blob_gc_batch(sql.as_ref(), root_key).await.unwrap();
    let claims: i64 = backend
        .pool()
        .writer()
        .unwrap()
        .conn()
        .query_row("SELECT COUNT(*) FROM blob_gc_claims", [], |row| row.get(0))
        .unwrap();
    assert_eq!(claims, 0);
}