use super::*;
use crate::pool::PoolConfig;
use rusqlite::hooks::{AuthAction, AuthContext, Authorization, TransactionOperation};
use serial_test::serial;
use std::sync::atomic::{AtomicBool, Ordering};
fn deny_commit(ctx: AuthContext<'_>) -> Authorization {
match ctx.action {
AuthAction::Transaction {
operation: TransactionOperation::Unknown,
} => Authorization::Deny,
_ => Authorization::Allow,
}
}
fn deny_rollback(ctx: AuthContext<'_>) -> Authorization {
match ctx.action {
AuthAction::Transaction {
operation: TransactionOperation::Rollback,
} => Authorization::Deny,
_ => Authorization::Allow,
}
}
fn deny_count_function(ctx: AuthContext<'_>) -> Authorization {
match ctx.action {
AuthAction::Function { function_name } if function_name.eq_ignore_ascii_case("count") => {
Authorization::Deny
}
_ => Authorization::Allow,
}
}
pub(super) mod page_snapshot_seam {
use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::sync::Mutex;
struct Barrier {
operation: &'static str,
namespace: String,
reached_tx: SyncSender<()>,
proceed_rx: Receiver<()>,
}
static BARRIER: Mutex<Option<Barrier>> = Mutex::new(None);
pub(crate) fn install(
operation: &'static str,
namespace: String,
) -> (Receiver<()>, SyncSender<()>) {
let (reached_tx, reached_rx) = sync_channel(0);
let (proceed_tx, proceed_rx) = sync_channel(0);
*BARRIER.lock().unwrap() = Some(Barrier {
operation,
namespace,
reached_tx,
proceed_rx,
});
(reached_rx, proceed_tx)
}
pub(crate) fn uninstall() {
*BARRIER.lock().unwrap() = None;
}
pub(crate) fn hook(operation: &'static str, namespace: &str) {
let barrier = {
let mut guard = BARRIER.lock().unwrap();
match guard.as_ref() {
Some(barrier)
if barrier.operation == operation && barrier.namespace == namespace =>
{
guard.take()
}
_ => None,
}
};
let Some(barrier) = barrier else {
return;
};
let _ = barrier.reached_tx.send(());
let _ = barrier.proceed_rx.recv();
}
}
fn setup_pool() -> Arc<ConnectionPool> {
let config = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = Arc::new(ConnectionPool::new(config).unwrap());
{
let writer = pool.writer().unwrap();
writer
.conn()
.execute_batch(&format!("{NOTES_DDL}\n{TEST_ATTACHMENTS_DDL}"))
.unwrap();
}
pool
}
const TEST_ATTACHMENTS_DDL: &str = r#"
CREATE TABLE attachments (
record_uuid TEXT NOT NULL,
substrate TEXT NOT NULL CHECK (substrate IN ('entity', 'note')),
role TEXT NOT NULL,
content_ref TEXT NOT NULL,
media_type TEXT,
size_bytes INTEGER,
created_at INTEGER NOT NULL,
PRIMARY KEY (record_uuid, role)
);
"#;
fn setup_memory_store() -> SqlNoteStore {
SqlNoteStore::new(setup_pool(), false)
}
fn make_note(namespace: &str, kind: &str, content: &str) -> Note {
Note::new(namespace, kind, content)
}
fn keyed_note(namespace: &str, kind: &str, key: &str) -> Note {
let mut note = make_note(namespace, kind, key);
note.key = Some(key.to_string());
note
}
fn assert_note_keys(notes: &[Note], expected_len: usize) {
assert_eq!(notes.len(), expected_len);
for note in notes {
assert_eq!(note.key.as_deref(), Some(note.content.as_str()));
}
}
#[tokio::test]
async fn note_key_round_trips_through_every_note_read_projection() {
let store = setup_memory_store();
let mut notes = vec![
keyed_note("local", "memory", "first"),
keyed_note("local", "memory", "second"),
keyed_note("local", "memory", "third"),
];
for (index, note) in notes.iter_mut().enumerate() {
note.id = Uuid::from_u128(index as u128 + 1);
note.created_at = if index < 2 { 100 } else { 90 };
}
assert_eq!(store.upsert_notes(notes.clone()).await.unwrap().affected, 3);
let page = PageRequest {
offset: 0,
limit: 10,
};
let filter = NoteFilter {
kind: Some("memory".into()),
..Default::default()
};
let first = store.get_note(notes[0].id).await.unwrap().unwrap();
assert_eq!(first, notes[0]);
assert_eq!(
store.get_note_including_deleted(first.id).await.unwrap(),
Some(first.clone())
);
assert_note_keys(
&store
.get_notes_batch(&[notes[0].id, notes[1].id])
.await
.unwrap(),
2,
);
assert_note_keys(
&store
.query_notes("local", Some("memory"), page.clone())
.await
.unwrap()
.items,
3,
);
assert_note_keys(
&store
.query_notes_count_free("local", Some("memory"), page.clone())
.await
.unwrap()
.items,
3,
);
assert_note_keys(
&store
.query_notes_filtered("local", &filter, page.clone())
.await
.unwrap()
.items,
3,
);
assert_note_keys(
&store
.query_notes_filtered_count_free("local", &filter, page.clone())
.await
.unwrap()
.items,
3,
);
assert_note_keys(
&store
.query_notes_filtered_bounded("local", &filter, 10)
.await
.unwrap(),
3,
);
let seek_filter = NoteFilter {
after: Some(NoteSeekAfter {
created_at: first.created_at,
id: first.id,
}),
..filter.clone()
};
let after = store
.query_notes_filtered_count_free("local", &seek_filter, page)
.await
.unwrap();
assert_note_keys(&after.items, 2);
assert_eq!(after.items[0].id, notes[1].id);
assert_eq!(after.items[1].id, notes[2].id);
let first_page = store
.query_notes_filtered_after("local", &filter, None, 1)
.await
.unwrap();
assert_note_keys(&first_page.items, 1);
assert_eq!(first_page.items[0].id, notes[0].id);
let second_page = store
.query_notes_filtered_after("local", &filter, first_page.next_after, 10)
.await
.unwrap();
assert_note_keys(&second_page.items, 2);
assert_eq!(second_page.items[0].id, notes[1].id);
}
#[tokio::test]
async fn note_key_is_scoped_by_namespace_and_kind_and_released_by_delete() {
let store = setup_memory_store();
let first = keyed_note("local", "memory", "shared");
store.upsert_note(first.clone()).await.unwrap();
assert!(store
.upsert_note(keyed_note("local", "memory", "shared"))
.await
.is_err());
store
.upsert_note(keyed_note("other", "memory", "shared"))
.await
.unwrap();
store
.upsert_note(keyed_note("local", "reference", "shared"))
.await
.unwrap();
for _ in 0..2 {
store
.upsert_note(make_note("local", "memory", "unkeyed"))
.await
.unwrap();
}
assert!(store.delete_note(first.id, DeleteMode::Soft).await.unwrap());
let deleted = store
.get_note_including_deleted(first.id)
.await
.unwrap()
.unwrap();
assert_eq!(deleted.key.as_deref(), Some("shared"));
assert!(deleted.deleted_at.is_some());
let second = keyed_note("local", "memory", "shared");
store.upsert_note(second.clone()).await.unwrap();
assert_ne!(first.id, second.id);
assert!(store
.delete_note(second.id, DeleteMode::Hard)
.await
.unwrap());
assert!(store
.get_note_including_deleted(second.id)
.await
.unwrap()
.is_none());
let third = keyed_note("local", "memory", "shared");
store.upsert_note(third.clone()).await.unwrap();
assert_ne!(second.id, third.id);
assert_eq!(store.count_notes("local", Some("memory")).await.unwrap(), 3);
}
#[tokio::test]
async fn note_key_round_trips_through_insert_paths_without_changing_external_id_dedup() {
let store = setup_memory_store();
let first = keyed_note("local", "memory", "insert-only");
assert!(store.insert_note_if_absent(first.clone()).await.unwrap());
assert_eq!(store.get_note(first.id).await.unwrap(), Some(first.clone()));
assert!(!store.insert_note_if_absent(first).await.unwrap());
assert!(store
.insert_note_if_absent(keyed_note("local", "memory", "insert-only"))
.await
.is_err());
let second = keyed_note("local", "memory", "try-insert");
assert!(store.try_insert_note(second.clone()).await.unwrap());
assert_eq!(store.get_note(second.id).await.unwrap(), Some(second));
let error = store
.try_insert_note(keyed_note("local", "memory", "try-insert"))
.await
.unwrap_err();
assert!(error
.to_string()
.contains("constraint other than external_id dedup"));
{
let writer = store.pool.writer().unwrap();
writer
.conn()
.execute_batch(crate::migrations::MIGRATIONS[4].up)
.unwrap();
}
let props = serde_json::json!({"external_id": "external-1"});
assert!(store
.try_insert_note(make_note("local", "message", "first").with_properties(props.clone()))
.await
.unwrap());
assert!(!store
.try_insert_note(make_note("local", "message", "second").with_properties(props))
.await
.unwrap());
}
#[tokio::test]
async fn note_key_survives_existing_full_and_property_updates() {
let store = setup_memory_store();
let original = keyed_note("local", "memory", "immutable");
let id = original.id;
store.upsert_note(original.clone()).await.unwrap();
let sequence = store.note_sequence(id).await.unwrap();
for candidate_key in [None, Some("replacement".to_string())] {
let mut update = store.get_note(id).await.unwrap().unwrap();
update.key = candidate_key;
update.content = "updated".to_string();
update.updated_at += 1;
store.upsert_note(update).await.unwrap();
assert_eq!(store.get_note(id).await.unwrap().unwrap().key, original.key);
}
let mut batch_update = store.get_note(id).await.unwrap().unwrap();
batch_update.key = Some("batch-replacement".to_string());
batch_update.updated_at += 1;
assert_eq!(
store
.upsert_notes(vec![batch_update])
.await
.unwrap()
.affected,
1
);
let snapshot = store.get_note(id).await.unwrap().unwrap();
assert_eq!(snapshot.key, original.key);
let mut replacement = snapshot.clone();
replacement.key = None;
replacement.updated_at += 1;
assert!(store
.replace_note_if_unchanged(replacement, snapshot.updated_at, snapshot.deleted_at)
.await
.unwrap());
assert_eq!(store.get_note(id).await.unwrap().unwrap().key, original.key);
assert!(store
.update_note_properties(
id,
Some(serde_json::json!({"tags": ["one"]})),
snapshot.updated_at + 2
)
.await
.unwrap());
assert!(store
.set_note_property(
id,
"tags",
serde_json::json!(["two"]),
snapshot.updated_at + 3
)
.await
.unwrap());
let final_note = store.get_note(id).await.unwrap().unwrap();
assert_eq!(final_note.key, original.key);
assert_eq!(
final_note.properties.unwrap()["tags"],
serde_json::json!(["two"])
);
assert_eq!(store.note_sequence(id).await.unwrap(), sequence);
let mut unkeyed = make_note("local", "memory", "no identity");
store.upsert_note(unkeyed.clone()).await.unwrap();
unkeyed.key = Some("late identity".to_string());
store.upsert_note(unkeyed.clone()).await.unwrap();
assert_eq!(store.get_note(unkeyed.id).await.unwrap().unwrap().key, None);
}
#[tokio::test]
async fn test_upsert_and_get_note() {
let store = setup_memory_store();
let note = make_note("default", "observation", "Hello world");
let id = note.id;
store.upsert_note(note).await.unwrap();
let fetched = store.get_note(id).await.unwrap();
assert!(fetched.is_some());
let fetched = fetched.unwrap();
assert_eq!(fetched.id, id);
assert_eq!(fetched.content, "Hello world");
assert_eq!(fetched.kind, "observation");
}
#[tokio::test]
async fn replace_note_cas_requires_new_revision_strictly_greater_than_snapshot() {
let store = setup_memory_store();
let mut original = make_note("default", "observation", "original");
original.created_at = 100;
original.updated_at = 100;
let id = original.id;
store.upsert_note(original.clone()).await.unwrap();
for refused_revision in [99, 100] {
let mut replacement = original.clone();
replacement.content = format!("must-not-land-{refused_revision}");
replacement.updated_at = refused_revision;
assert!(
!store
.replace_note_if_unchanged(replacement, original.updated_at, original.deleted_at)
.await
.unwrap(),
"CAS must refuse replacement revision {refused_revision} when the snapshot revision is {}",
original.updated_at
);
assert_eq!(
store.get_note(id).await.unwrap().unwrap().content,
"original"
);
}
let mut advanced = original.clone();
advanced.content = "advanced".to_string();
advanced.updated_at = 101;
assert!(
store
.replace_note_if_unchanged(advanced, original.updated_at, original.deleted_at)
.await
.unwrap(),
"a strictly newer revision must still satisfy the CAS"
);
let persisted = store.get_note(id).await.unwrap().unwrap();
assert_eq!(persisted.content, "advanced");
assert_eq!(persisted.updated_at, 101);
}
#[tokio::test]
async fn insert_note_if_absent_reports_the_loser_and_leaves_the_winner_untouched() {
let store = setup_memory_store();
let first = make_note("default", "observation", "first writer");
let id = first.id;
assert!(
store.insert_note_if_absent(first).await.unwrap(),
"the first insert on an absent id must report that it inserted"
);
let mut second = make_note("default", "insight", "second writer");
second.id = id;
second.updated_at += 1;
assert!(
!store.insert_note_if_absent(second).await.unwrap(),
"a second insert on an id that now exists must report that it did NOT insert, \
rather than reporting success to a caller whose write did not land"
);
let persisted = store.get_note(id).await.unwrap().unwrap();
assert_eq!(
persisted.content, "first writer",
"the pre-existing row must survive byte-for-byte: overwriting it is the behaviour \
this primitive exists to avoid, and it is invisible to the return value alone"
);
assert_eq!(
persisted.kind, "observation",
"no column of the pre-existing row may be rewritten by the refused insert"
);
}
#[tokio::test]
async fn test_kind_roundtrip_all_variants() {
let store = setup_memory_store();
for kind in [
"observation",
"insight",
"question",
"decision",
"reference",
] {
let note = make_note("default", kind, "content");
let id = note.id;
store.upsert_note(note).await.unwrap();
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(fetched.kind, kind);
}
}
#[tokio::test]
async fn test_soft_delete() {
let store = setup_memory_store();
let note = make_note("default", "observation", "to be deleted");
let id = note.id;
store.upsert_note(note).await.unwrap();
let deleted = store.delete_note(id, DeleteMode::Soft).await.unwrap();
assert!(deleted);
let fetched = store.get_note(id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_hard_delete() {
let store = setup_memory_store();
let note = make_note("default", "observation", "to be hard deleted");
let id = note.id;
store.upsert_note(note).await.unwrap();
let deleted = store.delete_note(id, DeleteMode::Hard).await.unwrap();
assert!(deleted);
let fetched = store.get_note(id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn note_soft_delete_retains_attachments_and_hard_delete_removes_them() {
let pool = setup_pool();
let store = SqlNoteStore::new(pool.clone(), false);
let note = make_note("default", "observation", "attached note");
let id = note.id;
store.upsert_note(note).await.unwrap();
pool.writer()
.unwrap()
.conn()
.execute(
"INSERT INTO attachments \
(record_uuid, substrate, role, content_ref, created_at) \
VALUES (?1, 'note', 'content', ?2, 123)",
rusqlite::params![id.to_string(), "a".repeat(64)],
)
.unwrap();
assert!(store.delete_note(id, DeleteMode::Soft).await.unwrap());
let retained: i64 = pool
.reader()
.unwrap()
.conn()
.query_row(
"SELECT COUNT(*) FROM attachments WHERE record_uuid = ?1",
[id.to_string()],
|row| row.get(0),
)
.unwrap();
assert_eq!(retained, 1);
assert!(store.delete_note(id, DeleteMode::Hard).await.unwrap());
let removed: i64 = pool
.reader()
.unwrap()
.conn()
.query_row(
"SELECT COUNT(*) FROM attachments WHERE record_uuid = ?1",
[id.to_string()],
|row| row.get(0),
)
.unwrap();
assert_eq!(removed, 0);
}
#[tokio::test]
async fn test_namespace_isolation() {
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
for _ in 0..3 {
store
.upsert_note(make_note("ns1", "observation", "content"))
.await
.unwrap();
}
store
.upsert_note(make_note("ns2", "observation", "other"))
.await
.unwrap();
let count_ns1 = store.count_notes("ns1", None).await.unwrap();
assert_eq!(count_ns1, 3);
let count_ns2 = store.count_notes("ns2", None).await.unwrap();
assert_eq!(count_ns2, 1);
}
#[tokio::test]
async fn batched_namespace_note_count_exceeds_sqlite_variable_limit() {
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
let live_a = make_note("stats-a", "observation", "live-a");
let deleted_a = make_note("stats-a", "observation", "deleted-a");
let deleted_a_id = deleted_a.id;
let live_b = make_note("stats-b", "insight", "live-b");
store.upsert_note(live_a).await.unwrap();
store.upsert_note(deleted_a).await.unwrap();
store.upsert_note(live_b).await.unwrap();
assert!(store
.delete_note(deleted_a_id, DeleteMode::Soft)
.await
.unwrap());
let per_namespace_total = store.count_notes("stats-a", None).await.unwrap()
+ store.count_notes("stats-b", None).await.unwrap();
pool.writer()
.unwrap()
.conn()
.set_limit(rusqlite::limits::Limit::SQLITE_LIMIT_VARIABLE_NUMBER, 999)
.unwrap();
let mut namespaces = vec!["stats-a".to_string(), "stats-b".to_string()];
namespaces.extend((0..999).map(|i| format!("empty-{i}")));
assert_eq!(namespaces.len(), 1_001);
assert_eq!(
store
.count_notes_in_namespaces(&namespaces, None)
.await
.unwrap(),
per_namespace_total
);
assert_eq!(
store
.count_notes_in_namespaces(&namespaces, Some("observation"))
.await
.unwrap(),
1
);
assert_eq!(per_namespace_total, 2);
}
#[tokio::test]
async fn duplicate_namespace_across_chunk_boundary_is_not_double_counted() {
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
store
.upsert_note(make_note("stats-a", "observation", "live-a-1"))
.await
.unwrap();
store
.upsert_note(make_note("stats-a", "observation", "live-a-2"))
.await
.unwrap();
let per_namespace_total = store.count_notes("stats-a", None).await.unwrap();
assert_eq!(per_namespace_total, 2);
let mut namespaces = vec!["stats-a".to_string()];
namespaces.extend((0..500).map(|i| format!("empty-{i}")));
assert_eq!(namespaces.len(), 501);
namespaces.push("stats-a".to_string());
assert_eq!(namespaces.len(), 502);
assert_eq!(
store
.count_notes_in_namespaces(&namespaces, None)
.await
.unwrap(),
per_namespace_total
);
}
#[tokio::test]
async fn test_query_and_count_use_caller_namespace() {
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
store
.upsert_note(make_note("ns_a", "observation", "A"))
.await
.unwrap();
store
.upsert_note(make_note("ns_b", "insight", "B"))
.await
.unwrap();
let page_a = store
.query_notes("ns_a", None, PageRequest::default())
.await
.unwrap();
assert_eq!(page_a.items.len(), 1);
assert_eq!(page_a.items[0].content, "A");
assert_eq!(page_a.total, Some(1));
let page_b = store
.query_notes("ns_b", None, PageRequest::default())
.await
.unwrap();
assert_eq!(page_b.items.len(), 1);
assert_eq!(page_b.items[0].content, "B");
assert_eq!(page_b.total, Some(1));
let count_a = store.count_notes("ns_a", None).await.unwrap();
let count_b = store.count_notes("ns_b", None).await.unwrap();
assert_eq!(count_a, 1);
assert_eq!(count_b, 1);
}
#[derive(Clone, Copy, Debug)]
enum SnapshotPageQuery {
Basic,
Filtered,
}
impl SnapshotPageQuery {
fn operation(self) -> &'static str {
match self {
Self::Basic => "query_notes",
Self::Filtered => "query_notes_filtered",
}
}
}
async fn run_snapshot_page_query(
store: &SqlNoteStore,
query: SnapshotPageQuery,
namespace: &str,
filter: &NoteFilter,
) -> Result<Page<Note>, StorageError> {
let page = PageRequest {
offset: 0,
limit: 10,
};
match query {
SnapshotPageQuery::Basic => {
store
.query_notes(namespace, Some("observation"), page)
.await
}
SnapshotPageQuery::Filtered => store.query_notes_filtered(namespace, filter, page).await,
}
}
async fn assert_page_count_and_items_share_snapshot(query: SnapshotPageQuery) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join(format!("note-page-snapshot-{query:?}.db"));
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(path),
write_queue_enabled: Some(false),
..PoolConfig::for_test()
})
.unwrap(),
);
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
let journal_mode: String = writer
.conn()
.pragma_query_value(None, "journal_mode", |row| row.get(0))
.unwrap();
assert_eq!(journal_mode.to_ascii_lowercase(), "wal");
}
let store = Arc::new(SqlNoteStore::new(Arc::clone(&pool), true));
let namespace = format!("snapshot-{}", Uuid::new_v4());
let properties = serde_json::json!({"snapshot_case": true});
let mut initial = make_note_with_props(
&namespace,
"observation",
"present when count starts",
properties.clone(),
);
initial.created_at = 1;
let initial_id = initial.id;
store.upsert_note(initial).await.unwrap();
let filter = NoteFilter {
kind: Some("observation".to_string()),
property_filters: vec![khive_storage::note::PropertyFilter {
json_path: "$.snapshot_case".to_string(),
op: FilterOp::Eq,
value: SqlValue::Bool(true),
}],
..NoteFilter::default()
};
let (reached_rx, proceed_tx) =
page_snapshot_seam::install(query.operation(), namespace.clone());
let query_task = {
let store = Arc::clone(&store);
let namespace = namespace.clone();
let filter = filter.clone();
tokio::spawn(
async move { run_snapshot_page_query(&store, query, &namespace, &filter).await },
)
};
tokio::task::spawn_blocking(move || reached_rx.recv_timeout(std::time::Duration::from_secs(5)))
.await
.expect("waiting for the count-to-page seam must not panic")
.expect("query must reach the seam after reading its count");
let mut concurrent = make_note_with_props(
&namespace,
"observation",
"committed between count and page",
properties,
);
concurrent.created_at = 2;
let concurrent_id = concurrent.id;
store
.upsert_note(concurrent)
.await
.expect("WAL writer must commit while the page reader is parked");
assert!(
store.get_note(concurrent_id).await.unwrap().is_some(),
"a new reader must observe the committed row before the page reader resumes"
);
assert!(
!query_task.is_finished(),
"page query must remain parked until the test releases its production seam"
);
proceed_tx
.send(())
.expect("page query must still be waiting at the production seam");
let page = query_task
.await
.expect("page query task must not panic")
.expect("page query must succeed");
page_snapshot_seam::uninstall();
assert_eq!(page.total, Some(1));
assert_eq!(
page.items.iter().map(|note| note.id).collect::<Vec<_>>(),
vec![initial_id]
);
let after = run_snapshot_page_query(&store, query, &namespace, &filter)
.await
.unwrap();
assert_eq!(after.total, Some(2));
assert!(after.items.iter().any(|note| note.id == concurrent_id));
}
#[tokio::test]
#[serial]
async fn note_page_count_and_items_share_one_snapshot_during_concurrent_insert() {
for query in [SnapshotPageQuery::Basic, SnapshotPageQuery::Filtered] {
assert_page_count_and_items_share_snapshot(query).await;
}
}
#[tokio::test]
async fn filtered_count_free_page_runs_without_count_over_large_match_set() {
use khive_storage::note::PropertyFilter as NotePropFilter;
let dir = tempfile::tempdir().unwrap();
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(dir.path().join("note-count-free-large-filter.db")),
max_readers: 1,
write_queue_enabled: Some(false),
..PoolConfig::default()
})
.unwrap(),
);
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = SqlNoteStore::new(Arc::clone(&pool), false);
let namespace = format!("count-free-large-{}", Uuid::new_v4().simple());
let notes = (0..2_000)
.map(|index| {
let mut note = make_note_with_props(
&namespace,
"message",
&format!("message-{index}"),
serde_json::json!({"direction": "inbound"}),
);
note.created_at = index;
note
})
.collect();
let summary = store.upsert_notes(notes).await.unwrap();
assert_eq!(summary.failed, 0, "large filtered seed failed: {summary:?}");
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
}],
..NoteFilter::default()
};
{
let reader = pool.reader().unwrap();
reader.conn().authorizer(Some(deny_count_function)).unwrap();
}
let request = PageRequest {
offset: 0,
limit: 6,
};
let exact_control = store
.query_notes_filtered(&namespace, &filter, request.clone())
.await;
assert!(
exact_control.is_err(),
"control exact-count page must be rejected by the count authorizer"
);
let exact_unfiltered_control = store
.query_notes(&namespace, Some("message"), request.clone())
.await;
assert!(
exact_unfiltered_control.is_err(),
"control unfiltered exact-count page must be rejected by the count authorizer"
);
let page = store
.query_notes_filtered_count_free(&namespace, &filter, request.clone())
.await
.expect("count-free page must not invoke SQLite count");
assert_eq!(page.total, None);
assert_eq!(page.items.len(), 6, "caller receives its lookahead row");
assert_eq!(page.items[0].content, "message-1999");
assert_eq!(page.items[5].content, "message-1994");
let unfiltered_page = store
.query_notes_count_free(&namespace, Some("message"), request)
.await
.expect("unfiltered count-free page must not invoke SQLite count");
assert_eq!(unfiltered_page.total, None);
assert_eq!(unfiltered_page.items.len(), 6);
assert_eq!(unfiltered_page.items[0].content, "message-1999");
assert_eq!(unfiltered_page.items[5].content, "message-1994");
let reader = pool.reader().unwrap();
reader
.conn()
.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
.unwrap();
}
#[tokio::test]
#[serial]
async fn filtered_count_free_page_keeps_one_statement_snapshot_and_total_order() {
let dir = tempfile::tempdir().unwrap();
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(dir.path().join("note-count-free-snapshot.db")),
write_queue_enabled: Some(false),
..PoolConfig::for_test()
})
.unwrap(),
);
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = Arc::new(SqlNoteStore::new(Arc::clone(&pool), true));
let namespace = format!("count-free-snapshot-{}", Uuid::new_v4().simple());
let mut initial = Vec::new();
for content in ["a", "b", "c"] {
let mut note = make_note_with_props(
&namespace,
"message",
content,
serde_json::json!({"page_case": true}),
);
note.created_at = 10;
initial.push(note);
}
let mut expected_ids: Vec<_> = initial.iter().map(|note| note.id).collect();
expected_ids.sort();
store.upsert_notes(initial).await.unwrap();
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![khive_storage::note::PropertyFilter {
json_path: "$.page_case".to_string(),
op: FilterOp::Eq,
value: SqlValue::Bool(true),
}],
..NoteFilter::default()
};
let (reached_rx, proceed_tx) =
page_snapshot_seam::install("query_notes_filtered_count_free", namespace.clone());
let query_task = {
let store = Arc::clone(&store);
let namespace = namespace.clone();
let filter = filter.clone();
tokio::spawn(async move {
store
.query_notes_filtered_count_free(
&namespace,
&filter,
PageRequest {
offset: 0,
limit: 3,
},
)
.await
})
};
tokio::task::spawn_blocking(move || reached_rx.recv_timeout(std::time::Duration::from_secs(5)))
.await
.expect("waiting for the row-step seam must not panic")
.expect("count-free query must step its first row");
let mut concurrent = make_note_with_props(
&namespace,
"message",
"concurrent",
serde_json::json!({"page_case": true}),
);
concurrent.created_at = 11;
let concurrent_id = concurrent.id;
store
.upsert_note(concurrent)
.await
.expect("WAL writer must commit while the page statement is paused");
proceed_tx.send(()).unwrap();
let page = query_task.await.unwrap().unwrap();
page_snapshot_seam::uninstall();
assert_eq!(page.total, None);
assert_eq!(
page.items.iter().map(|note| note.id).collect::<Vec<_>>(),
expected_ids,
"the lookahead page must retain id-ascending tie order on its pinned snapshot"
);
let after = store
.query_notes_filtered_count_free(
&namespace,
&filter,
PageRequest {
offset: 0,
limit: 4,
},
)
.await
.unwrap();
assert_eq!(after.items[0].id, concurrent_id);
assert_eq!(
after.items[1..]
.iter()
.map(|note| note.id)
.collect::<Vec<_>>(),
expected_ids
);
}
#[derive(Clone, Copy, Debug)]
enum SnapshotCountQuery {
Exact,
Bounded,
}
impl SnapshotCountQuery {
fn operation(self) -> &'static str {
match self {
Self::Exact => "count_notes_filtered_in_snapshot",
Self::Bounded => "count_notes_filtered_bounded_in_snapshot",
}
}
}
async fn run_snapshot_count_query(
store: &SqlNoteStore,
query: SnapshotCountQuery,
namespace: &str,
filters: &[NoteFilter],
) -> Result<Vec<u64>, StorageError> {
match query {
SnapshotCountQuery::Exact => {
store
.count_notes_filtered_in_snapshot(namespace, filters)
.await
}
SnapshotCountQuery::Bounded => store
.count_notes_filtered_bounded_in_snapshot(namespace, filters, 10)
.await
.map(|counts| {
counts
.into_iter()
.map(|count| {
assert!(!count.saturated, "one-row partition cannot hit cap");
assert_eq!(count.cap, 10);
count.count
})
.collect()
}),
}
}
async fn assert_filtered_count_partitions_share_snapshot(query: SnapshotCountQuery) {
use khive_storage::note::PropertyFilter as NotePropFilter;
let dir = tempfile::tempdir().unwrap();
let path = dir
.path()
.join(format!("note-filter-count-snapshot-{query:?}.db"));
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(path),
write_queue_enabled: Some(false),
..PoolConfig::for_test()
})
.unwrap(),
);
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = Arc::new(SqlNoteStore::new(Arc::clone(&pool), true));
let namespace = format!("count-snapshot-{}", Uuid::new_v4().simple());
let mut note = make_note_with_props(
&namespace,
"message",
"addressed before concurrent update",
serde_json::json!({
"direction": "inbound",
"read": false,
"to_actor": "actor:a",
}),
);
note.created_at = 1;
let note_id = note.id;
store.upsert_note(note).await.unwrap();
let base_filters = vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
];
let partition_filter = |op| {
let mut property_filters = base_filters.clone();
property_filters.push(NotePropFilter {
json_path: "$.to_actor".to_string(),
op,
value: SqlValue::Text("actor:a".to_string()),
});
NoteFilter {
kind: Some("message".to_string()),
property_filters,
..NoteFilter::default()
}
};
let filters = vec![
partition_filter(FilterOp::EqOrMissingIndexed),
partition_filter(FilterOp::JsonTypeMissingOrNullIndexed),
];
let (reached_rx, proceed_tx) =
page_snapshot_seam::install(query.operation(), namespace.clone());
let mut query_task = {
let store = Arc::clone(&store);
let namespace = namespace.clone();
let filters = filters.clone();
tokio::spawn(
async move { run_snapshot_count_query(&store, query, &namespace, &filters).await },
)
};
let mut seam_task = tokio::task::spawn_blocking(move || {
reached_rx.recv_timeout(std::time::Duration::from_secs(60))
});
tokio::select! {
reached = &mut seam_task => {
if !matches!(&reached, Ok(Ok(()))) {
page_snapshot_seam::uninstall();
drop(proceed_tx);
query_task.abort();
let query_outcome = query_task.await;
if matches!(
&reached,
Ok(Err(std::sync::mpsc::RecvTimeoutError::Timeout))
) {
panic!(
"precondition timeout: {query:?} count snapshot seam exceeded its 60s hang watchdog; \
query outcome after cleanup: {query_outcome:?}"
);
}
panic!(
"precondition: {query:?} count snapshot seam waiter failed; \
seam outcome: {reached:?}; query outcome after cleanup: {query_outcome:?}"
);
}
}
outcome = &mut query_task => {
page_snapshot_seam::uninstall();
drop(proceed_tx);
let seam_outcome = seam_task.await;
panic!(
"precondition: {query:?} count query completed while waiting for its snapshot seam; \
query outcome: {outcome:?}; seam outcome: {seam_outcome:?}"
);
}
}
let writer = pool.writer().unwrap();
writer
.conn()
.execute(
"UPDATE notes SET properties = json_object('direction', 'inbound', 'read', 0) \
WHERE id = ?1",
[note_id.to_string()],
)
.unwrap();
proceed_tx
.send(())
.expect("count query must still be waiting at the production seam");
let counts = query_task
.await
.expect("count query task must not panic")
.expect("count query must succeed");
page_snapshot_seam::uninstall();
assert_eq!(counts, vec![1, 0], "both partitions must use one snapshot");
let after = run_snapshot_count_query(&store, query, &namespace, &filters)
.await
.unwrap();
assert_eq!(
after,
vec![0, 1],
"the committed update must appear afterward"
);
}
#[tokio::test]
#[serial]
async fn filtered_count_partitions_share_one_snapshot_during_concurrent_update() {
for query in [SnapshotCountQuery::Exact, SnapshotCountQuery::Bounded] {
assert_filtered_count_partitions_share_snapshot(query).await;
}
}
#[tokio::test]
async fn bounded_filtered_count_is_exact_at_cap_and_saturates_above_it() {
use khive_storage::BoundedCount;
let store = setup_memory_store();
let filter = NoteFilter {
kind: Some("message".to_string()),
..NoteFilter::default()
};
for index in 0..5 {
store
.upsert_note(make_note("default", "message", &format!("message-{index}")))
.await
.unwrap();
}
let exact = store
.count_notes_filtered_bounded_in_snapshot("default", std::slice::from_ref(&filter), 5)
.await
.unwrap();
assert_eq!(
exact,
vec![BoundedCount {
count: 5,
cap: 5,
saturated: false,
}],
"a population equal to the cap is still exact"
);
store
.upsert_note(make_note("default", "message", "over-cap"))
.await
.unwrap();
let saturated = store
.count_notes_filtered_bounded_in_snapshot("default", &[filter], 5)
.await
.unwrap();
assert_eq!(
saturated,
vec![BoundedCount {
count: 5,
cap: 5,
saturated: true,
}]
);
}
#[tokio::test]
async fn test_soft_delete_sets_status_deleted() {
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
let note = make_note("default", "observation", "to delete");
let id = note.id;
store.upsert_note(note).await.unwrap();
let deleted = store.delete_note(id, DeleteMode::Soft).await.unwrap();
assert!(deleted);
let writer = pool.writer().unwrap();
let status: String = writer
.conn()
.query_row(
"SELECT status FROM notes WHERE id = ?1",
[id.to_string()],
|r| r.get(0),
)
.unwrap();
assert_eq!(status, "deleted");
}
#[tokio::test]
async fn test_note_status_field_roundtrip() {
let store = setup_memory_store();
let note = make_note("default", "observation", "status test");
let id = note.id;
store.upsert_note(note).await.unwrap();
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(fetched.status, "active");
}
#[tokio::test]
async fn set_note_property_initializes_null_and_preserves_json_type() {
let store = setup_memory_store();
let note = make_note("default", "observation", "atomic property set");
let id = note.id;
let updated_at = note.updated_at + 1;
store.upsert_note(note).await.unwrap();
assert!(store
.set_note_property(
id,
"delivery.stamp",
serde_json::json!({ "channel": "email", "attempt": 1 }),
updated_at,
)
.await
.unwrap());
assert!(store
.set_note_property(id, "explicit_null", serde_json::Value::Null, updated_at + 1)
.await
.unwrap());
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(
fetched.properties,
Some(serde_json::json!({
"delivery.stamp": { "channel": "email", "attempt": 1 },
"explicit_null": null
})),
"keys must be literal top-level segments and values must keep their JSON types"
);
assert_eq!(fetched.updated_at, updated_at + 1);
}
#[tokio::test]
async fn concurrent_distinct_note_property_sets_both_survive() {
let store = Arc::new(setup_memory_store());
let note = make_note("default", "message", "concurrent atomic properties")
.with_properties(serde_json::json!({ "existing": "preserved" }));
let id = note.id;
let updated_at = note.updated_at;
store.upsert_note(note).await.unwrap();
let gate = Arc::new(tokio::sync::Barrier::new(3));
let left = {
let store = Arc::clone(&store);
let gate = Arc::clone(&gate);
tokio::spawn(async move {
gate.wait().await;
store
.set_note_property(
id,
"delivery_stamp",
serde_json::json!("email"),
updated_at + 1,
)
.await
})
};
let right = {
let store = Arc::clone(&store);
let gate = Arc::clone(&gate);
tokio::spawn(async move {
gate.wait().await;
store
.set_note_property(id, "ingest_marker", serde_json::json!(true), updated_at + 2)
.await
})
};
gate.wait().await;
assert!(left.await.unwrap().unwrap());
assert!(right.await.unwrap().unwrap());
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(
fetched.properties,
Some(serde_json::json!({
"existing": "preserved",
"delivery_stamp": "email",
"ingest_marker": true
})),
"one-statement property sets must not lose a different concurrent key"
);
}
#[tokio::test]
async fn set_note_property_refuses_non_object_document() {
let store = setup_memory_store();
let note = make_note("default", "observation", "scalar properties")
.with_properties(serde_json::json!(["not", "an", "object"]));
let id = note.id;
let original_updated_at = note.updated_at;
store.upsert_note(note).await.unwrap();
assert!(!store
.set_note_property(id, "read", serde_json::json!(true), original_updated_at + 1)
.await
.unwrap());
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(
fetched.properties,
Some(serde_json::json!(["not", "an", "object"]))
);
assert_eq!(fetched.updated_at, original_updated_at);
}
#[tokio::test]
async fn try_patch_note_property_refuses_scalar_document() {
use khive_storage::note::NoteFilter;
let store = setup_memory_store();
let note =
make_note("default", "message", "scalar properties").with_properties(serde_json::json!(1));
let id = note.id;
let original_updated_at = note.updated_at;
store.upsert_note(note).await.unwrap();
let matched = store
.try_patch_note_property(
id,
"default",
&NoteFilter::default(),
"$.read",
serde_json::json!(true),
original_updated_at + 1,
)
.await
.unwrap();
assert!(!matched, "a scalar properties document must not be patched");
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(fetched.properties, Some(serde_json::json!(1)));
assert_eq!(fetched.updated_at, original_updated_at);
}
#[tokio::test]
async fn try_patch_note_property_refuses_array_document() {
use khive_storage::note::NoteFilter;
let store = setup_memory_store();
let note = make_note("default", "message", "array properties")
.with_properties(serde_json::json!(["not", "an", "object"]));
let id = note.id;
let original_updated_at = note.updated_at;
store.upsert_note(note).await.unwrap();
let matched = store
.try_patch_note_property(
id,
"default",
&NoteFilter::default(),
"$.read",
serde_json::json!(true),
original_updated_at + 1,
)
.await
.unwrap();
assert!(!matched, "an array properties document must not be patched");
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(
fetched.properties,
Some(serde_json::json!(["not", "an", "object"]))
);
assert_eq!(fetched.updated_at, original_updated_at);
}
#[test]
fn text_prefix_upper_bound_increments_the_last_code_point() {
assert_eq!(text_prefix_upper_bound("email:").as_deref(), Some("email;"));
assert_eq!(text_prefix_upper_bound("a").as_deref(), Some("b"));
assert_eq!(text_prefix_upper_bound("a\u{10FFFF}").as_deref(), Some("b"));
assert_eq!(text_prefix_upper_bound(""), None);
assert_eq!(text_prefix_upper_bound("\u{10FFFF}"), None);
assert_eq!(
text_prefix_upper_bound("x\u{D7FF}").as_deref(),
Some("x\u{E000}")
);
}
#[tokio::test]
async fn text_starts_with_indexed_matches_prefix_only() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter};
use khive_storage::types::{PageRequest, SqlValue};
let store = setup_memory_store();
let seeded = [
(
"email:a@b.c",
serde_json::json!({"to_actor": "email:a@b.c"}),
),
("email:", serde_json::json!({"to_actor": "email:"})),
("emaik:zzz", serde_json::json!({"to_actor": "emaik:zzz"})),
("email;", serde_json::json!({"to_actor": "email;"})),
("telegram:1", serde_json::json!({"to_actor": "telegram:1"})),
("number", serde_json::json!({"to_actor": 5})),
("null", serde_json::json!({"to_actor": null})),
("missing", serde_json::json!({})),
];
for (name, properties) in seeded {
let mut note = Note::new("default", "message", name);
note.properties = Some(properties);
store.upsert_note(note).await.unwrap();
}
let filter = NoteFilter {
kind: Some("message".into()),
property_filters: vec![PropertyFilter {
json_path: "$.to_actor".into(),
op: FilterOp::TextStartsWithIndexed,
value: SqlValue::Text("email:".into()),
}],
..Default::default()
};
let page = store
.query_notes_filtered_count_free(
"default",
&filter,
PageRequest {
limit: 50,
offset: 0,
},
)
.await
.unwrap();
let mut names: Vec<_> = page.items.iter().map(|n| n.content.clone()).collect();
names.sort();
assert_eq!(names, vec!["email:", "email:a@b.c"]);
let every_text = NoteFilter {
kind: Some("message".into()),
property_filters: vec![PropertyFilter {
json_path: "$.to_actor".into(),
op: FilterOp::TextStartsWithIndexed,
value: SqlValue::Text(String::new()),
}],
..Default::default()
};
let page = store
.query_notes_filtered_count_free(
"default",
&every_text,
PageRequest {
limit: 50,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(
page.items.len(),
5,
"an empty prefix admits every text value and nothing else"
);
}
fn atomic_mark_read_filter() -> khive_storage::note::NoteFilter {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter};
use khive_storage::types::SqlValue;
NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
PropertyFilter {
json_path: "$.direction".to_string(),
op: FilterOp::NotInOrMissing(vec![SqlValue::Text("outbound".to_string())]),
value: SqlValue::Null,
},
PropertyFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrMissing,
value: SqlValue::Text("lambda:reader".to_string()),
},
],
..Default::default()
}
}
#[tokio::test]
async fn atomic_note_property_patch_rolls_back_when_one_target_is_ineligible() {
let store = setup_memory_store();
let eligible = make_note("local", "message", "eligible").with_properties(serde_json::json!({
"direction": "inbound",
"to_actor": "lambda:reader",
"read": false,
"preserve": "eligible",
}));
let ineligible =
make_note("local", "message", "ineligible").with_properties(serde_json::json!({
"direction": "outbound",
"to_actor": "lambda:reader",
"read": false,
"preserve": "ineligible",
}));
let eligible_id = eligible.id;
let ineligible_id = ineligible.id;
let updated_at = eligible.updated_at.max(ineligible.updated_at) + 1;
store.upsert_note(eligible).await.unwrap();
store.upsert_note(ineligible).await.unwrap();
let filter = atomic_mark_read_filter();
let error = store
.patch_note_property_atomic(
vec![eligible_id, ineligible_id],
"local",
&filter,
"$.read",
serde_json::json!(true),
updated_at,
)
.await
.expect_err("an ineligible target must abort the atomic patch");
assert!(
matches!(
&error,
StorageError::WriterTaskRequestFailed {
request_state: WriterTaskRequestState::TransactionRolledBack,
source,
} if matches!(source.as_ref(), StorageError::Conflict { message, .. }
if message.contains(&ineligible_id.to_string()))
),
"the conflict must name the first failing id {ineligible_id}; got {error:?}"
);
assert!(!error.is_retryable(), "a precondition conflict is terminal");
for (id, preserved) in [(eligible_id, "eligible"), (ineligible_id, "ineligible")] {
let stored = store.get_note(id).await.unwrap().unwrap();
let properties = stored.properties.unwrap();
assert_eq!(properties["read"], false);
assert_eq!(properties["preserve"], preserved);
}
store
.patch_note_property_atomic(
vec![eligible_id, eligible_id],
"local",
&filter,
"$.read",
serde_json::json!(true),
updated_at,
)
.await
.expect("deduplicated eligible targets commit together");
let stored = store.get_note(eligible_id).await.unwrap().unwrap();
let properties = stored.properties.unwrap();
assert_eq!(properties["read"], true);
assert_eq!(properties["preserve"], "eligible");
}
#[tokio::test]
async fn atomic_note_property_patch_writer_task_commits_and_rolls_back() {
let dir = tempfile::tempdir().unwrap();
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(dir.path().join("atomic-note-property-writer-task.db")),
write_queue_enabled: Some(true),
..PoolConfig::for_test()
})
.unwrap(),
);
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = SqlNoteStore::new(Arc::clone(&pool), true);
let eligible_properties = serde_json::json!({
"direction": "inbound",
"to_actor": "lambda:reader",
"read": false,
});
let first = make_note("local", "message", "first eligible")
.with_properties(eligible_properties.clone());
let second = make_note("local", "message", "second eligible")
.with_properties(eligible_properties.clone());
let ineligible =
make_note("local", "message", "ineligible").with_properties(serde_json::json!({
"direction": "outbound",
"to_actor": "lambda:reader",
"read": false,
}));
let first_id = first.id;
let second_id = second.id;
let ineligible_id = ineligible.id;
let updated_at = first
.updated_at
.max(second.updated_at)
.max(ineligible.updated_at)
+ 1;
for note in [first, second, ineligible] {
store.upsert_note(note).await.unwrap();
}
let filter = atomic_mark_read_filter();
store
.patch_note_property_atomic(
vec![first_id, second_id],
"local",
&filter,
"$.read",
serde_json::json!(true),
updated_at,
)
.await
.expect("all eligible rows commit through the writer task");
for id in [first_id, second_id] {
let stored = store.get_note(id).await.unwrap().unwrap();
assert_eq!(stored.properties.unwrap()["read"], true);
}
store
.update_note_properties(first_id, Some(eligible_properties), updated_at + 1)
.await
.unwrap();
let error = store
.patch_note_property_atomic(
vec![first_id, ineligible_id],
"local",
&filter,
"$.read",
serde_json::json!(true),
updated_at + 2,
)
.await
.expect_err("a later ineligible row must abort the writer-task transaction");
assert!(
matches!(
&error,
StorageError::WriterTaskRequestFailed {
request_state: WriterTaskRequestState::TransactionRolledBack,
source,
} if matches!(source.as_ref(), StorageError::Conflict { message, .. }
if message.contains(&ineligible_id.to_string()))
),
"the conflict must name the first failing id {ineligible_id}; got {error:?}"
);
assert_eq!(
store
.get_note(first_id)
.await
.unwrap()
.unwrap()
.properties
.unwrap()["read"],
false,
"the earlier eligible update must roll back"
);
assert_eq!(pool.writer_task_spawn_count(), 1);
}
#[tokio::test]
async fn set_note_property_rejects_nul_key_without_mutation() {
let store = setup_memory_store();
let note = make_note("default", "observation", "nul property key").with_properties(
serde_json::json!({
"a": 0,
"a\u{0000}b": 9,
}),
);
let id = note.id;
let original_updated_at = note.updated_at;
let original_properties = note.properties.clone();
store.upsert_note(note).await.unwrap();
let result = store
.set_note_property(
id,
"a\u{0000}b",
serde_json::json!(1),
original_updated_at + 1,
)
.await;
assert!(
matches!(result, Err(StorageError::InvalidInput { .. })),
"expected InvalidInput, got {result:?}"
);
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(fetched.properties, original_properties);
assert_eq!(fetched.updated_at, original_updated_at);
}
fn make_note_with_props(
namespace: &str,
kind: &str,
content: &str,
props: serde_json::Value,
) -> Note {
Note::new(namespace, kind, content).with_properties(props)
}
#[tokio::test]
async fn eq_or_missing_does_not_match_present_empty_string() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter};
use khive_storage::types::{PageRequest, SqlValue};
let store = setup_memory_store();
store
.upsert_note(make_note_with_props(
"default",
"message",
"present empty",
serde_json::json!({"to_actor": ""}),
))
.await
.unwrap();
store
.upsert_note(make_note_with_props(
"default",
"message",
"absent",
serde_json::json!({}),
))
.await
.unwrap();
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![PropertyFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrMissing,
value: SqlValue::Text("actor:a".to_string()),
}],
..Default::default()
};
let page = store
.query_notes_filtered(
"default",
&filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.total, Some(1));
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].content, "absent");
}
#[tokio::test]
async fn test_filtered_namespace_and_kind_isolation() {
let store = setup_memory_store();
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::{PageRequest, SqlValue};
let n1 = make_note_with_props(
"ns1",
"scheduled_event",
"event1",
serde_json::json!({"status": "pending", "trigger_at": "2027-01-01T00:00:00Z"}),
);
let n2 = make_note_with_props(
"ns1",
"scheduled_event",
"event2",
serde_json::json!({"status": "done", "trigger_at": "2027-01-02T00:00:00Z"}),
);
let n3 = make_note_with_props(
"ns2",
"scheduled_event",
"event3",
serde_json::json!({"status": "pending", "trigger_at": "2027-01-03T00:00:00Z"}),
);
store.upsert_note(n1).await.unwrap();
store.upsert_note(n2).await.unwrap();
store.upsert_note(n3).await.unwrap();
let filter = NoteFilter {
kind: Some("scheduled_event".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("pending".to_string()),
}],
order_by: None,
..Default::default()
};
let page = store
.query_notes_filtered("ns1", &filter, PageRequest::default())
.await
.unwrap();
assert_eq!(
page.items.len(),
1,
"only the pending ns1 event should appear"
);
assert_eq!(page.items[0].content, "event1");
assert_eq!(page.total, Some(1));
}
#[tokio::test]
async fn test_filtered_order_by_json_path_asc() {
let store = setup_memory_store();
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter, SortDir};
use khive_storage::types::{PageRequest, SqlValue};
let n3 = make_note_with_props(
"ns1",
"scheduled_event",
"third",
serde_json::json!({"status": "pending", "trigger_at": "2027-01-03T00:00:00Z"}),
);
let n1 = make_note_with_props(
"ns1",
"scheduled_event",
"first",
serde_json::json!({"status": "pending", "trigger_at": "2027-01-01T00:00:00Z"}),
);
let n2 = make_note_with_props(
"ns1",
"scheduled_event",
"second",
serde_json::json!({"status": "pending", "trigger_at": "2027-01-02T00:00:00Z"}),
);
store.upsert_note(n3).await.unwrap();
store.upsert_note(n1).await.unwrap();
store.upsert_note(n2).await.unwrap();
let filter = NoteFilter {
kind: Some("scheduled_event".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("pending".to_string()),
}],
order_by: Some(("$.trigger_at".to_string(), SortDir::Asc)),
..Default::default()
};
let page = store
.query_notes_filtered("ns1", &filter, PageRequest::default())
.await
.unwrap();
assert_eq!(page.items.len(), 3);
assert_eq!(page.items[0].content, "first");
assert_eq!(page.items[1].content, "second");
assert_eq!(page.items[2].content, "third");
}
#[tokio::test]
async fn filtered_default_order_is_stable_across_equal_timestamp_pages() {
let store = setup_memory_store();
let created_at = 1_750_000_000_000_000_i64;
let mut expected_ids = Vec::new();
for index in 0..317 {
let mut note = make_note("ns1", "observation", &format!("note-{index}"));
note.created_at = created_at;
expected_ids.push(note.id);
store.upsert_note(note).await.unwrap();
}
expected_ids.sort_unstable();
let mut actual_ids = Vec::new();
let page_size = 29_u32;
let mut offset = 0_u64;
loop {
let page = store
.query_notes_filtered(
"ns1",
&NoteFilter::default(),
PageRequest {
offset,
limit: page_size,
},
)
.await
.unwrap();
if page.items.is_empty() {
break;
}
offset += page.items.len() as u64;
actual_ids.extend(page.items.into_iter().map(|note| note.id));
}
assert_eq!(actual_ids, expected_ids);
}
#[tokio::test]
async fn query_notes_offset_sweep_covers_equal_created_at_exactly_once() {
let store = setup_memory_store();
let created_at = 1_750_000_000_000_000_i64;
let mut expected_ids = Vec::new();
for index in 0..211 {
let mut note = make_note("ns1", "observation", &format!("note-{index}"));
note.created_at = created_at;
expected_ids.push(note.id);
store.upsert_note(note).await.unwrap();
}
expected_ids.sort_unstable();
let mut actual_ids = Vec::new();
let page_size = 37_u32;
let mut offset = 0_u64;
loop {
let page = store
.query_notes(
"ns1",
None,
PageRequest {
offset,
limit: page_size,
},
)
.await
.unwrap();
if page.items.is_empty() {
break;
}
offset += page.items.len() as u64;
actual_ids.extend(page.items.into_iter().map(|note| note.id));
}
assert_eq!(actual_ids, expected_ids);
}
#[tokio::test]
async fn query_notes_filtered_custom_order_offset_sweep_is_total() {
let store = setup_memory_store();
let mut expected_ids = Vec::new();
for index in 0..113 {
let note = make_note_with_props(
"ns1",
"observation",
&format!("note-{index}"),
serde_json::json!({"rank": 7}),
);
expected_ids.push(note.id);
store.upsert_note(note).await.unwrap();
}
expected_ids.sort_unstable_by(|a, b| b.cmp(a));
let filter = NoteFilter {
order_by: Some(("$.rank".to_string(), SortDir::Desc)),
..Default::default()
};
let mut actual_ids = Vec::new();
let page_size = 23_u32;
let mut offset = 0_u64;
loop {
let page = store
.query_notes_filtered(
"ns1",
&filter,
PageRequest {
offset,
limit: page_size,
},
)
.await
.unwrap();
if page.items.is_empty() {
break;
}
offset += page.items.len() as u64;
actual_ids.extend(page.items.into_iter().map(|note| note.id));
}
assert_eq!(actual_ids, expected_ids);
}
#[tokio::test]
async fn test_filtered_soft_deleted_excluded() {
let store = setup_memory_store();
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::{DeleteMode, PageRequest, SqlValue};
let n = make_note_with_props(
"ns1",
"scheduled_event",
"to_delete",
serde_json::json!({"status": "pending"}),
);
let id = n.id;
store.upsert_note(n).await.unwrap();
store.delete_note(id, DeleteMode::Soft).await.unwrap();
let filter = NoteFilter {
kind: Some("scheduled_event".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("pending".to_string()),
}],
order_by: None,
..Default::default()
};
let page = store
.query_notes_filtered("ns1", &filter, PageRequest::default())
.await
.unwrap();
assert_eq!(page.items.len(), 0, "soft-deleted rows must not appear");
}
#[tokio::test]
async fn test_filtered_invalid_json_path_rejected() {
let store = setup_memory_store();
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::{PageRequest, SqlValue};
let filter = NoteFilter {
kind: None,
property_filters: vec![NotePropFilter {
json_path: "DROP TABLE notes".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("x".to_string()),
}],
order_by: None,
..Default::default()
};
let result = store
.query_notes_filtered("ns1", &filter, PageRequest::default())
.await;
assert!(
result.is_err(),
"invalid json_path must be rejected before SQL"
);
}
#[tokio::test]
async fn test_try_insert_note_pk_collision_returns_error_not_dedup() {
let store = setup_memory_store();
let mut note = make_note("ns1", "message", "original content");
let fixed_id = uuid::Uuid::parse_str("00000000-0000-0000-0000-000000000099").unwrap();
note.id = fixed_id;
let inserted = store
.try_insert_note(note.clone())
.await
.expect("first insert must succeed");
assert!(inserted, "first insert must return true");
let result = store.try_insert_note(note).await;
assert!(
result.is_err(),
"PK collision without external_id must return StorageError, not Ok(false)"
);
}
#[tokio::test]
async fn test_upsert_note_insert_and_seq_assignment_are_atomic() {
let store = setup_memory_store();
let fail_id = uuid::Uuid::parse_str("00000000-0000-0000-0000-0000000000aa").unwrap();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch(&format!(
"CREATE TRIGGER inject_seq_failure_upsert BEFORE INSERT ON notes_seq \
WHEN NEW.note_id = '{fail_id}' \
BEGIN SELECT RAISE(ABORT, 'injected failure for #827 atomicity test'); END;"
))
.unwrap();
}
let mut note = make_note("ns1", "message", "atomic test upsert_note");
note.id = fail_id;
let result = store.upsert_note(note).await;
assert!(
result.is_err(),
"the injected notes_seq trigger failure must surface as an error"
);
let fetched = store.get_note(fail_id).await.unwrap();
assert!(
fetched.is_none(),
"the note insert must roll back together with the failed sequence \
assignment, not strand the note without a notes_seq row: {fetched:?}"
);
}
#[tokio::test]
async fn test_upsert_note_is_true_upsert_no_delete_semantics() {
let store = setup_memory_store();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch(
"PRAGMA recursive_triggers = ON;
CREATE TABLE delete_fires (n INTEGER);
CREATE TRIGGER notes_delete_probe AFTER DELETE ON notes \
BEGIN INSERT INTO delete_fires VALUES (1); END;",
)
.unwrap();
}
let mut note = make_note("default", "observation", "v1");
let id = note.id;
let original_created_at = note.created_at;
store.upsert_note(note.clone()).await.unwrap();
note.content = "v2".to_string();
note.salience = Some(0.9);
note.updated_at += 1_000;
note.created_at += 1_000;
store.upsert_note(note).await.unwrap();
let fetched = store.get_note(id).await.unwrap().unwrap();
assert_eq!(
fetched.content, "v2",
"mutable fields must reflect the second upsert"
);
assert_eq!(fetched.salience, Some(0.9));
assert_eq!(
fetched.created_at, original_created_at,
"created_at must be preserved across an upsert of an existing row"
);
let (row_count, delete_fires): (i64, i64) = {
let writer = store.pool.try_writer().unwrap();
let conn = writer.conn();
let row_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM notes WHERE id = ?1",
rusqlite::params![id.to_string()],
|row| row.get(0),
)
.unwrap();
let delete_fires: i64 = conn
.query_row("SELECT COUNT(*) FROM delete_fires", [], |row| row.get(0))
.unwrap();
(row_count, delete_fires)
};
assert_eq!(
row_count, 1,
"upsert must update in place, never duplicate rows"
);
assert_eq!(
delete_fires, 0,
"upserting an existing row must not fire DELETE-path triggers"
);
store.delete_note(id, DeleteMode::Hard).await.unwrap();
let delete_fires_after: i64 = {
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.query_row("SELECT COUNT(*) FROM delete_fires", [], |row| row.get(0))
.unwrap()
};
assert_eq!(
delete_fires_after, 1,
"the probe trigger must fire on a genuine delete"
);
}
#[tokio::test]
async fn test_try_insert_note_insert_and_seq_assignment_are_atomic() {
let store = setup_memory_store();
let fail_id = uuid::Uuid::parse_str("00000000-0000-0000-0000-0000000000bb").unwrap();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch(&format!(
"CREATE TRIGGER inject_seq_failure_try_insert BEFORE INSERT ON notes_seq \
WHEN NEW.note_id = '{fail_id}' \
BEGIN SELECT RAISE(ABORT, 'injected failure for #827 atomicity test'); END;"
))
.unwrap();
}
let mut note = make_note("ns1", "message", "atomic test try_insert_note");
note.id = fail_id;
let result = store.try_insert_note(note).await;
assert!(
result.is_err(),
"the injected notes_seq trigger failure must surface as an error"
);
let fetched = store.get_note(fail_id).await.unwrap();
assert!(
fetched.is_none(),
"the note insert must roll back together with the failed sequence \
assignment, not strand the note without a notes_seq row: {fetched:?}"
);
}
#[tokio::test]
async fn pooled_transaction_commit_failure_with_verified_rollback_keeps_writer_usable() {
let store = setup_memory_store();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE tx_finalize (id INTEGER PRIMARY KEY)")
.unwrap();
}
let result = store
.with_writer_tx("test_pooled_commit", |conn| {
conn.execute("INSERT INTO tx_finalize (id) VALUES (1)", [])?;
conn.authorizer(Some(deny_commit))?;
Ok(())
})
.await;
assert!(
matches!(
&result,
Err(StorageError::WriterTaskRequestFailed {
request_state: WriterTaskRequestState::TransactionRolledBack,
source,
}) if matches!(source.as_ref(), StorageError::Pool { operation, .. }
if operation == "test_pooled_commit")
),
"a denied COMMIT followed by a verified rollback must report the commit error: {result:?}"
);
let writer = store
.pool
.try_writer()
.expect("verified rollback must leave the pooled writer usable");
let conn = writer.conn();
assert!(conn.is_autocommit());
conn.authorizer(None::<fn(AuthContext<'_>) -> Authorization>)
.unwrap();
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM tx_finalize", [], |row| row.get(0))
.unwrap();
assert_eq!(count, 0);
}
#[tokio::test]
async fn pooled_transaction_rollback_failure_reports_unknown_and_retires_writer() {
let store = setup_memory_store();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE tx_rollback_fault (id INTEGER PRIMARY KEY)")
.unwrap();
}
let result = store
.with_writer_tx(
"test_pooled_rollback",
|conn| -> Result<(), rusqlite::Error> {
conn.execute("INSERT INTO tx_rollback_fault (id) VALUES (1)", [])?;
conn.authorizer(Some(deny_rollback))?;
Err(rusqlite::Error::InvalidQuery)
},
)
.await;
assert!(
matches!(
result,
Err(StorageError::WriterTaskTerminated {
request_state: WriterTaskRequestState::SideEffectsUnknown,
})
),
"a failed rollback cannot claim that the attempted write did not land"
);
let checkout = store.pool.try_writer();
assert!(
matches!(checkout, Err(SqliteError::InvalidData(message)) if message.contains("retired")),
"a connection with an unverified rollback must never be checked out again"
);
let legacy = store.pool.legacy_conn();
let legacy_guard = legacy.lock();
let direct_probe = legacy_guard.query_row("SELECT 1", [], |row| row.get::<_, i64>(0));
assert!(
direct_probe.is_err(),
"the compatibility raw-connection handle must not bypass retirement quarantine"
);
}
#[tokio::test]
async fn pooled_transaction_panic_with_failed_rollback_reports_unknown_and_retires_writer() {
let store = setup_memory_store();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch("CREATE TABLE tx_panic_fault (id INTEGER PRIMARY KEY)")
.unwrap();
}
let result = store
.with_writer_tx("test_pooled_panic", |conn| -> Result<(), rusqlite::Error> {
conn.execute("INSERT INTO tx_panic_fault (id) VALUES (1)", [])?;
conn.authorizer(Some(deny_rollback))?;
panic!("intentional pooled transaction panic");
})
.await;
assert!(
matches!(
result,
Err(StorageError::WriterTaskTerminated {
request_state: WriterTaskRequestState::SideEffectsUnknown,
})
),
"a panic whose rollback cannot be verified must report unknown side effects"
);
let checkout = store.pool.try_writer();
assert!(
matches!(checkout, Err(SqliteError::InvalidData(message)) if message.contains("retired")),
"the pooled writer must retire after a transaction-body panic"
);
}
#[tokio::test]
async fn pooled_transaction_refuses_preexisting_non_autocommit_connection() {
let store = setup_memory_store();
{
let writer = store.pool.try_writer().unwrap();
writer.conn().execute_batch("BEGIN IMMEDIATE").unwrap();
}
let operation_ran = Arc::new(AtomicBool::new(false));
let operation_ran_in_closure = Arc::clone(&operation_ran);
let result = store
.with_writer_tx("test_pooled_preexisting_transaction", move |_conn| {
operation_ran_in_closure.store(true, Ordering::SeqCst);
Ok(())
})
.await;
assert!(
matches!(
result,
Err(StorageError::WriterTaskTerminated {
request_state: WriterTaskRequestState::SideEffectsUnknown,
})
),
"an inherited transaction has an unknown prior outcome and must fail closed"
);
assert!(!operation_ran.load(Ordering::SeqCst));
let checkout = store.pool.try_writer();
assert!(
matches!(checkout, Err(SqliteError::InvalidData(message)) if message.contains("retired")),
"the inherited non-autocommit connection must not be reused"
);
}
#[tokio::test]
async fn test_upsert_notes_batch_rolls_back_fully_on_mid_batch_seq_failure() {
let store = setup_memory_store();
let fail_id = uuid::Uuid::parse_str("00000000-0000-0000-0000-0000000000cc").unwrap();
{
let writer = store.pool.try_writer().unwrap();
writer
.conn()
.execute_batch(&format!(
"CREATE TRIGGER inject_seq_failure_batch BEFORE INSERT ON notes_seq \
WHEN NEW.note_id = '{fail_id}' \
AND EXISTS (SELECT 1 FROM notes_seq WHERE note_id = NEW.note_id) \
BEGIN SELECT RAISE(ABORT, 'injected mid-batch failure for #827 test'); END;"
))
.unwrap();
}
let mut note_ok = make_note(
"ns1",
"message",
"first note in batch, seq assignment succeeds",
);
let ok_id = note_ok.id;
let mut note_fail = make_note(
"ns1",
"message",
"second note in batch, seq assignment fails",
);
note_fail.id = fail_id;
note_ok.created_at = 1_000_000;
note_fail.created_at = 1_000_001;
let result = store.upsert_notes(vec![note_ok, note_fail]).await;
assert!(
result.is_err(),
"the injected mid-batch notes_seq trigger failure must surface as an error, not a \
partial BatchWriteSummary: {result:?}"
);
let fetched_ok = store.get_note(ok_id).await.unwrap();
assert!(
fetched_ok.is_none(),
"the whole batch must roll back -- the first note (whose own insert and seq \
assignment succeeded) must not survive a later note's failure in the same batch: \
{fetched_ok:?}"
);
let fetched_fail = store.get_note(fail_id).await.unwrap();
assert!(
fetched_fail.is_none(),
"the failed note must not survive either: {fetched_fail:?}"
);
let next_note = make_note("ns1", "message", "write after rolled-back batch");
let next_id = next_note.id;
store
.upsert_note(next_note)
.await
.expect("a write after the rolled-back batch must succeed, not hang on an open BEGIN");
let fetched_next = store.get_note(next_id).await.unwrap();
assert!(
fetched_next.is_some(),
"the post-rollback write must actually land"
);
}
#[tokio::test]
async fn page_offset_over_i64max_rejected() {
let store = setup_memory_store();
store
.upsert_note(make_note("ns1", "observation", "Hello world"))
.await
.unwrap();
let oversized = PageRequest {
offset: (i64::MAX as u64) + 1,
limit: 10,
};
let result = store.query_notes("ns1", None, oversized.clone()).await;
assert!(
matches!(result, Err(StorageError::InvalidInput { .. })),
"query_notes: expected InvalidInput, got {result:?}"
);
let filtered_result = store
.query_notes_filtered("ns1", &NoteFilter::default(), oversized.clone())
.await;
assert!(
matches!(filtered_result, Err(StorageError::InvalidInput { .. })),
"query_notes_filtered: expected InvalidInput, got {filtered_result:?}"
);
let count_free_result = store
.query_notes_filtered_count_free("ns1", &NoteFilter::default(), oversized)
.await;
assert!(
matches!(count_free_result, Err(StorageError::InvalidInput { .. })),
"query_notes_filtered_count_free: expected InvalidInput, got {count_free_result:?}"
);
}
#[tokio::test]
async fn upsert_notes_routes_through_writer_task_when_flag_enabled() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("write_queue_notes.db");
let pool_cfg = PoolConfig {
path: Some(path.clone()),
write_queue_enabled: Some(true),
..PoolConfig::for_test()
};
let pool = Arc::new(ConnectionPool::new(pool_cfg).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = SqlNoteStore::new(Arc::clone(&pool), true);
let n1 = make_note("default", "observation", "first");
let n2 = make_note("default", "observation", "second");
let id1 = n1.id;
let id2 = n2.id;
let summary = store.upsert_notes(vec![n1, n2]).await.unwrap();
assert_eq!(summary.attempted, 2);
assert_eq!(summary.affected, 2);
assert_eq!(summary.failed, 0);
assert!(store.get_note(id1).await.unwrap().is_some());
assert!(store.get_note(id2).await.unwrap().is_some());
assert_eq!(
pool.writer_task_spawn_count(),
1,
"the flag-ON path must actually spawn and use the writer task"
);
}
#[tokio::test]
async fn upsert_note_routes_through_writer_task_when_flag_enabled() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("write_queue_note_single.db");
let pool_cfg = PoolConfig {
path: Some(path.clone()),
write_queue_enabled: Some(true),
..PoolConfig::for_test()
};
let pool = Arc::new(ConnectionPool::new(pool_cfg).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = Arc::new(SqlNoteStore::new(Arc::clone(&pool), true));
let writer_task = pool
.writer_task_handle()
.unwrap()
.expect("writer task must be spawned with the flag on for a file-backed pool");
let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
let occupier = {
let writer_task = writer_task.clone();
tokio::spawn(async move {
writer_task
.send(move |_conn| {
let _ = started_tx.send(());
let _ = release_rx.blocking_recv();
Ok::<(), StorageError>(())
})
.await
})
};
started_rx
.await
.expect("occupier must signal it has started running inside the writer task");
assert_eq!(
writer_task.queue_depth(),
0,
"channel must start empty once the occupier has been dequeued and is running"
);
let note = make_note("default", "observation", "single-row write-queue routing");
let note_id = note.id;
let store_task = {
let store = Arc::clone(&store);
tokio::spawn(async move { store.upsert_note(note).await })
};
let mut saw_enqueued = false;
for _ in 0..100 {
if writer_task.queue_depth() >= 1 {
saw_enqueued = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
assert!(
saw_enqueued,
"upsert_note's write request never appeared in the writer task's \
channel while the occupier held the single drain slot — with_writer \
is not routing this single-row write through the shared writer task"
);
release_tx
.send(())
.expect("occupier must still be waiting on the release signal");
occupier
.await
.expect("occupier task must not panic")
.expect("occupier write must succeed");
store_task
.await
.expect("store task must not panic")
.expect("upsert_note must succeed once unblocked");
let fetched = store.get_note(note_id).await.unwrap();
assert!(
fetched.is_some(),
"note must be committed and readable after queuing behind the occupier"
);
}
#[tokio::test]
async fn upsert_note_reports_configured_write_queue_admission_deadline() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("note_store_bounded_admission.db");
let pool_cfg = PoolConfig {
path: Some(path),
write_queue_enabled: Some(true),
write_queue_capacity: 1,
write_admission_deadline_ms: 100,
..PoolConfig::for_test()
};
let pool = Arc::new(ConnectionPool::new(pool_cfg).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = SqlNoteStore::new(Arc::clone(&pool), true);
let writer_task = pool
.writer_task_handle()
.unwrap()
.expect("writer task must be enabled for the file-backed pool");
let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
let (release_tx, release_rx) = std::sync::mpsc::channel::<()>();
let a_task = {
let writer_task = writer_task.clone();
tokio::spawn(async move {
writer_task
.send(move |_conn| {
let _ = started_tx.send(());
release_rx.recv().expect("test must release request A");
Ok::<(), StorageError>(())
})
.await
})
};
let started = tokio::time::timeout(std::time::Duration::from_secs(5), started_rx).await;
if !matches!(started, Ok(Ok(()))) {
let _ = release_tx.send(());
panic!("request A did not start inside the writer task: {started:?}");
}
let b_task = {
let writer_task = writer_task.clone();
tokio::spawn(async move { writer_task.send(|_conn| Ok::<(), StorageError>(())).await })
};
let b_enqueued = tokio::time::timeout(std::time::Duration::from_secs(5), async {
while writer_task.queue_depth() != 1 {
tokio::task::yield_now().await;
}
})
.await;
if b_enqueued.is_err() {
let _ = release_tx.send(());
let _ = a_task.await;
let _ = b_task.await;
panic!("request B did not occupy the writer task's sole queue slot");
}
let note = make_note(
"default",
"observation",
"store-level bounded write admission regression",
);
let note_id = note.id;
let c_result =
tokio::time::timeout(std::time::Duration::from_secs(2), store.upsert_note(note)).await;
let release_result = release_tx.send(());
let a_result = tokio::time::timeout(std::time::Duration::from_secs(5), a_task).await;
let b_result = tokio::time::timeout(std::time::Duration::from_secs(5), b_task).await;
release_result.expect("request A must still be waiting for release");
a_result
.expect("request A did not complete after release")
.expect("request A task must not panic")
.expect("request A must complete successfully");
b_result
.expect("request B did not complete after request A")
.expect("request B task must not panic")
.expect("request B must complete successfully");
match c_result.expect("note-store admission waited beyond its bounded deadline") {
Err(StorageError::WriteQueueFull { timeout_ms }) => assert_eq!(timeout_ms, 100),
other => panic!("expected configured WriteQueueFull from SqlNoteStore, got {other:?}"),
}
assert!(
store.get_note(note_id).await.unwrap().is_none(),
"a queue-rejected note write must never execute after capacity returns"
);
}
#[test]
fn transactional_write_refreshes_writer_task_after_construction_outside_runtime() {
assert!(tokio::runtime::Handle::try_current().is_err());
let dir = tempfile::tempdir().unwrap();
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: Some(dir.path().join("note-late-writer-task.db")),
write_queue_enabled: Some(true),
..PoolConfig::for_test()
})
.unwrap(),
);
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(NOTES_DDL).unwrap();
}
let store = Arc::new(SqlNoteStore::new(Arc::clone(&pool), true));
tokio::runtime::Runtime::new()
.unwrap()
.block_on(async move {
let writer_task = pool
.writer_task_handle()
.unwrap()
.expect("file-backed pool must spawn its writer task inside the runtime");
let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
let occupier = {
let writer_task = writer_task.clone();
tokio::spawn(async move {
writer_task
.send(move |_conn| {
let _ = started_tx.send(());
let _ = release_rx.blocking_recv();
Ok::<(), StorageError>(())
})
.await
})
};
started_rx.await.unwrap();
let write = {
let store = Arc::clone(&store);
tokio::spawn(async move {
store
.upsert_note(make_note(
"default",
"observation",
"late transactional write",
))
.await
})
};
let mut saw_enqueued = false;
for _ in 0..100 {
if writer_task.queue_depth() >= 1 {
saw_enqueued = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
release_tx.send(()).unwrap();
occupier.await.unwrap().unwrap();
write.await.unwrap().unwrap();
assert!(
saw_enqueued,
"transactional note write bypassed the queue after construction cached no handle"
);
});
}
#[tokio::test]
async fn unread_probe_query_uses_partial_index() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
for i in 0..5 {
let unread = make_note_with_props(
"default",
"message",
&format!("unread {i}"),
serde_json::json!({"direction": "inbound", "to_actor": "actor:a"}),
);
store.upsert_note(unread).await.unwrap();
let read = make_note_with_props(
"default",
"message",
&format!("read {i}"),
serde_json::json!({"direction": "inbound", "to_actor": "actor:a", "read": true}),
);
store.upsert_note(read).await.unwrap();
}
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrMissingIndexed,
value: SqlValue::Text("actor:a".to_string()),
},
],
order_by: None,
..Default::default()
};
let rows = store
.query_notes_filtered_bounded("default", &filter, 1_000)
.await
.unwrap();
assert_eq!(rows.len(), 5, "exactly the unread rows must match");
let (where_sql, params) = build_note_filter_read_clause("default", &filter).unwrap();
let sql =
format!("SELECT id FROM notes{where_sql} ORDER BY created_at DESC, id ASC LIMIT 1001");
assert!(
where_sql.contains("json_type(properties, '$.read') != 'true'"),
"JsonTypeNeMissing must inline the validated json_type literal, got:\n{where_sql}"
);
let mut filter_without_json_type = filter.clone();
filter_without_json_type
.property_filters
.retain(|property| !matches!(property.op, FilterOp::JsonTypeNeMissing));
let (_, params_without_json_type) =
build_note_filter_read_clause("default", &filter_without_json_type).unwrap();
assert_eq!(
params.len(),
params_without_json_type.len(),
"JsonTypeNeMissing must not add a bind parameter"
);
let reader = pool.reader().unwrap();
let plan = |sql: &str, params: &[Box<dyn rusqlite::types::ToSql>]| -> String {
let mut stmt = reader
.conn()
.prepare(&format!("EXPLAIN QUERY PLAN {sql}"))
.unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let details: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(3))
.unwrap()
.map(Result::unwrap)
.collect();
assert!(!details.is_empty(), "EXPLAIN returned no plan rows");
details.join("\n")
};
let indexed_plan = plan(&sql, ¶ms);
assert!(
indexed_plan.contains("idx_notes_unread_probe_recipient_direction"),
"unread probe must be served by the partial index, got plan:\n{indexed_plan}"
);
reader
.conn()
.execute_batch("DROP INDEX idx_notes_unread_probe_recipient_direction")
.unwrap();
let error = reader.conn().prepare(&sql).err().unwrap();
assert!(
error
.to_string()
.contains("idx_notes_unread_probe_recipient_direction"),
"a missing pinned index must fail loudly: {error}"
);
}
#[tokio::test]
async fn unread_probe_legacy_partition_includes_null_and_uses_partial_index() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::types::SqlValue;
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
for (label, recipient) in [
("missing", None),
("null", Some(serde_json::Value::Null)),
("empty", Some(serde_json::Value::String(String::new()))),
("other", Some(serde_json::json!("actor:z"))),
] {
let mut properties = serde_json::json!({
"direction": "inbound",
"read": false,
});
if let Some(recipient) = recipient {
properties["to_actor"] = recipient;
}
store
.upsert_note(make_note_with_props(
"default", "message", label, properties,
))
.await
.unwrap();
}
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::JsonTypeMissingOrNullIndexed,
value: SqlValue::Text("actor:a".to_string()),
},
],
..NoteFilter::default()
};
let rows = store
.query_notes_filtered_bounded("default", &filter, 1_000)
.await
.unwrap();
assert_eq!(rows.len(), 2, "only missing and JSON-null recipients match");
let (where_sql, params) = build_note_filter_read_clause("default", &filter).unwrap();
let sql =
format!("SELECT id FROM notes{where_sql} ORDER BY created_at DESC, id ASC LIMIT 1001");
let reader = pool.reader().unwrap();
let mut stmt = reader
.conn()
.prepare(&format!("EXPLAIN QUERY PLAN {sql}"))
.unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let plan: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(3))
.unwrap()
.map(Result::unwrap)
.collect();
let plan = plan.join("\n");
assert!(
plan.contains("idx_notes_unread_probe_recipient_direction"),
"legacy recipient partition must use the partial index, got plan:\n{plan}"
);
}
#[tokio::test]
async fn unread_probe_work_is_bounded_by_callers_own_unread_rows() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
use rusqlite::StatementStatus;
let pool = setup_pool();
let store = SqlNoteStore::new(Arc::clone(&pool), false);
for i in 0..5 {
let mine = make_note_with_props(
"default",
"message",
&format!("mine {i}"),
serde_json::json!({"direction": "inbound", "to_actor": "actor:a"}),
);
store.upsert_note(mine).await.unwrap();
}
let seed_other = |n: usize| {
let writer = pool.writer().unwrap();
let conn = writer.conn();
let base: i64 = conn
.query_row("SELECT ifnull(max(created_at), 0) FROM notes", [], |row| {
row.get(0)
})
.unwrap();
conn.execute_batch("BEGIN").unwrap();
for i in 0..n {
let stamp = base + 1 + i as i64;
conn.execute(
"INSERT INTO notes (id, namespace, kind, content, properties, \
created_at, updated_at) \
VALUES (?1, 'default', 'message', ?2, \
json_object('direction', 'inbound', 'to_actor', 'actor:z'), \
?3, ?3)",
rusqlite::params![
uuid::Uuid::new_v4().to_string(),
format!("other {i}"),
stamp
],
)
.unwrap();
}
conn.execute_batch("COMMIT").unwrap();
};
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrMissingIndexed,
value: SqlValue::Text("actor:a".to_string()),
},
],
order_by: None,
..Default::default()
};
let (where_sql, params) = build_note_filter_read_clause("default", &filter).unwrap();
let sql = format!("SELECT id FROM notes{where_sql} ORDER BY created_at DESC, id ASC LIMIT 2");
let measure = || -> (usize, i32) {
let reader = pool.reader().unwrap();
let mut stmt = reader.conn().prepare(&sql).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let rows: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(0))
.unwrap()
.map(Result::unwrap)
.collect();
(rows.len(), stmt.get_status(StatementStatus::VmStep))
};
seed_other(200);
let (rows_small, steps_small) = measure();
assert_eq!(rows_small, 2, "probe must still see the caller's rows");
assert!(
steps_small > 0,
"instrument control: VM steps must register"
);
seed_other(1_800); let (rows_large, steps_large) = measure();
assert_eq!(rows_large, 2, "probe must still see the caller's rows");
assert!(
steps_large < steps_small.saturating_mul(3),
"unread probe scanned other recipients' backlog: \
{steps_small} VM steps at 200 irrelevant rows vs {steps_large} at 2000 \
(a recipient-scoped probe stays flat; a blind scan grows ~10x)"
);
}
#[tokio::test]
async fn unread_probe_bounded_by_inbound_not_by_recipients_own_outbound_history() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
use rusqlite::StatementStatus;
const CAP: i64 = 1_000;
fn count_filter() -> NoteFilter {
NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrMissingIndexed,
value: SqlValue::Text("actor:a".to_string()),
},
],
order_by: None,
..Default::default()
}
}
fn count_sql_and_params(filter: &NoteFilter) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
let (where_sql, mut params) = build_note_filter_read_clause("default", filter).unwrap();
params.push(Box::new(CAP + 1));
let limit_idx = params.len();
(
format!("SELECT COUNT(*) FROM (SELECT 1 FROM notes{where_sql} LIMIT ?{limit_idx})"),
params,
)
}
fn list_sql_and_params(filter: &NoteFilter) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
let (where_sql, mut params) = build_note_filter_read_clause("default", filter).unwrap();
params.push(Box::new(CAP + 1));
let limit_idx = params.len();
(
format!(
"SELECT id FROM notes{where_sql} ORDER BY created_at DESC, id ASC LIMIT ?{limit_idx}"
),
params,
)
}
fn seed_inbound(conn: &rusqlite::Connection, n: usize) {
for i in 0..n {
conn.execute(
"INSERT INTO notes (id, namespace, kind, content, properties, \
created_at, updated_at) \
VALUES (?1, 'default', 'message', ?2, \
json_object('direction', 'inbound', 'to_actor', 'actor:a'), \
?3, ?3)",
rusqlite::params![
uuid::Uuid::new_v4().to_string(),
format!("mine {i}"),
i as i64
],
)
.unwrap();
}
}
fn seed_outbound(conn: &rusqlite::Connection, n: usize) {
let base: i64 = conn
.query_row("SELECT ifnull(max(created_at), 0) FROM notes", [], |row| {
row.get(0)
})
.unwrap();
conn.execute_batch("BEGIN").unwrap();
for i in 0..n {
let stamp = base + 1 + i as i64;
conn.execute(
"INSERT INTO notes (id, namespace, kind, content, properties, \
created_at, updated_at) \
VALUES (?1, 'default', 'message', ?2, \
json_object('direction', 'outbound', 'to_actor', 'actor:a'), \
?3, ?3)",
rusqlite::params![
uuid::Uuid::new_v4().to_string(),
format!("sent-copy {i}"),
stamp
],
)
.unwrap();
}
conn.execute_batch("COMMIT").unwrap();
}
fn measure(
conn: &rusqlite::Connection,
sql: &str,
params: &[Box<dyn rusqlite::types::ToSql>],
) -> (i64, i32) {
let mut stmt = conn.prepare(sql).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let value: i64 = stmt.query_row(refs.as_slice(), |row| row.get(0)).unwrap();
(value, stmt.get_status(StatementStatus::VmStep))
}
fn measure_list(
conn: &rusqlite::Connection,
sql: &str,
params: &[Box<dyn rusqlite::types::ToSql>],
) -> (usize, i32) {
let mut stmt = conn.prepare(sql).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let rows: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(0))
.unwrap()
.map(Result::unwrap)
.collect();
let steps = stmt.get_status(StatementStatus::VmStep);
(rows.len(), steps)
}
let pool = setup_pool();
{
let writer = pool.writer().unwrap();
let conn = writer.conn();
seed_inbound(conn, 5);
}
let filter = count_filter();
let (count_sql, count_params) = count_sql_and_params(&filter);
let (list_sql, list_params) = list_sql_and_params(&filter);
{
let reader = pool.reader().unwrap();
let plan = |sql: &str, params: &[Box<dyn rusqlite::types::ToSql>]| -> String {
let mut stmt = reader
.conn()
.prepare(&format!("EXPLAIN QUERY PLAN {sql}"))
.unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|p| p.as_ref()).collect();
let details: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(3))
.unwrap()
.map(Result::unwrap)
.collect();
details.join("\n")
};
let count_plan = plan(&count_sql, &count_params);
assert!(
count_plan.contains("idx_notes_unread_probe_recipient_direction"),
"bounded-count query must be served by the direction-aware index, got:\n{count_plan}"
);
let list_plan = plan(&list_sql, &list_params);
assert!(
list_plan.contains("idx_notes_unread_probe_recipient_direction"),
"unread listing query must be served by the direction-aware index, got:\n{list_plan}"
);
}
{
let writer = pool.writer().unwrap();
let conn = writer.conn();
seed_outbound(conn, 3_000);
}
let (count_val_small, count_steps_small) = {
let reader = pool.reader().unwrap();
measure(reader.conn(), &count_sql, &count_params)
};
let (list_rows_small, list_steps_small) = {
let reader = pool.reader().unwrap();
measure_list(reader.conn(), &list_sql, &list_params)
};
assert_eq!(
count_val_small, 5,
"only the caller's unread inbound rows count"
);
assert_eq!(
list_rows_small, 5,
"only the caller's unread inbound rows list"
);
{
let writer = pool.writer().unwrap();
let conn = writer.conn();
seed_outbound(conn, 27_000); }
let (count_val_large, count_steps_large) = {
let reader = pool.reader().unwrap();
measure(reader.conn(), &count_sql, &count_params)
};
let (list_rows_large, list_steps_large) = {
let reader = pool.reader().unwrap();
measure_list(reader.conn(), &list_sql, &list_params)
};
assert_eq!(
count_val_large, 5,
"only the caller's unread inbound rows count"
);
assert_eq!(
list_rows_large, 5,
"only the caller's unread inbound rows list"
);
assert!(
count_steps_large < count_steps_small.saturating_mul(3),
"bounded-count work scaled with the recipient's own outbound history: \
{count_steps_small} VM steps at 3,000 outbound rows vs {count_steps_large} at 30,000 \
(a direction-scoped probe stays flat; a direction-blind scan grows ~10x)"
);
assert!(
list_steps_large < list_steps_small.saturating_mul(3),
"unread listing work scaled with the recipient's own outbound history: \
{list_steps_small} VM steps at 3,000 outbound rows vs {list_steps_large} at 30,000 \
(a direction-scoped probe stays flat; a direction-blind scan grows ~10x)"
);
let control_count_sql =
count_sql.replace(" INDEXED BY idx_notes_unread_probe_recipient_direction", "");
let control_pool = setup_pool();
{
let writer = control_pool.writer().unwrap();
let conn = writer.conn();
conn.execute_batch(
"DROP INDEX idx_notes_unread_probe_recipient_direction;
DROP INDEX idx_notes_message_recipient_direction;
CREATE INDEX idx_notes_unread_probe_recipient_control
ON notes(namespace, kind,
ifnull(json_extract(properties, '$.to_actor'), ''),
created_at DESC, id ASC)
WHERE (json_type(properties, '$.read') IS NULL
OR json_type(properties, '$.read') != 'true')
AND deleted_at IS NULL;",
)
.unwrap();
seed_inbound(conn, 5);
seed_outbound(conn, 3_000);
}
let (control_count_small, control_count_steps_small) = {
let reader = control_pool.reader().unwrap();
measure(reader.conn(), &control_count_sql, &count_params)
};
{
let writer = control_pool.writer().unwrap();
let conn = writer.conn();
seed_outbound(conn, 27_000);
}
let (control_count_large, control_count_steps_large) = {
let reader = control_pool.reader().unwrap();
measure(reader.conn(), &control_count_sql, &count_params)
};
assert_eq!(
control_count_small, 5,
"control: same correct value at 3,000 outbound"
);
assert_eq!(
control_count_large, 5,
"control: same correct value at 30,000 outbound"
);
assert!(
control_count_steps_large > control_count_steps_small.saturating_mul(3),
"control instrument did not reproduce the pre-fix growth: \
{control_count_steps_small} VM steps at 3,000 outbound rows vs \
{control_count_steps_large} at 30,000 under the direction-blind index shape \
(expected clear growth, proving the fixed measurement above is falsifiable)"
);
}
#[tokio::test]
async fn inbox_unread_listing_uses_recipient_index_not_direction_blind_scan() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
use rusqlite::StatementStatus;
fn listing_filter(to_actor_op: FilterOp) -> NoteFilter {
NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: to_actor_op,
value: SqlValue::Text("actor:a".to_string()),
},
],
order_by: None,
..Default::default()
}
}
fn listing_sql_and_params(
filter: &NoteFilter,
) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
let (where_sql, mut params) = build_note_filter_read_clause("default", filter).unwrap();
params.push(Box::new(1001_i64));
let limit_idx = params.len();
(
format!(
"SELECT id FROM notes{where_sql} ORDER BY created_at DESC, id ASC LIMIT ?{limit_idx}"
),
params,
)
}
fn plan(
conn: &rusqlite::Connection,
sql: &str,
params: &[Box<dyn rusqlite::types::ToSql>],
) -> String {
let mut stmt = conn.prepare(&format!("EXPLAIN QUERY PLAN {sql}")).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
stmt.query_map(refs.as_slice(), |row| row.get::<_, String>(3))
.unwrap()
.map(Result::unwrap)
.collect::<Vec<_>>()
.join("\n")
}
fn measure_list(
conn: &rusqlite::Connection,
sql: &str,
params: &[Box<dyn rusqlite::types::ToSql>],
) -> (usize, i32) {
let mut stmt = conn.prepare(sql).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let rows: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(0))
.unwrap()
.map(Result::unwrap)
.collect();
let steps = stmt.get_status(StatementStatus::VmStep);
(rows.len(), steps)
}
fn seed(
conn: &rusqlite::Connection,
to_actor: &str,
direction: &str,
n: usize,
start_stamp: i64,
) {
conn.execute_batch("BEGIN").unwrap();
for i in 0..n {
conn.execute(
"INSERT INTO notes (id, namespace, kind, content, properties, \
created_at, updated_at) \
VALUES (?1, 'default', 'message', ?2, \
json_object('direction', ?3, 'to_actor', ?4), ?5, ?5)",
rusqlite::params![
uuid::Uuid::new_v4().to_string(),
format!("{to_actor}-{direction}-{i}"),
direction,
to_actor,
start_stamp + i as i64,
],
)
.unwrap();
}
conn.execute_batch("COMMIT").unwrap();
}
fn build_pool_with_comm_indexes() -> Arc<ConnectionPool> {
let pool = setup_pool();
let writer = pool.writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_comm_message_direction \
ON notes(namespace, kind, json_extract(properties, '$.direction'), \
json_extract(properties, '$.read'), created_at DESC) \
WHERE deleted_at IS NULL;",
)
.unwrap();
drop(writer);
pool
}
{
let pool = build_pool_with_comm_indexes();
{
let writer = pool.writer().unwrap();
seed(writer.conn(), "actor:a", "inbound", 5, 0);
}
let filter = listing_filter(FilterOp::EqOrMissing);
let (sql, params) = listing_sql_and_params(&filter);
{
let reader = pool.reader().unwrap();
let control_plan = plan(reader.conn(), &sql, ¶ms);
assert!(
!control_plan.contains("idx_notes_unread_probe_recipient_direction"),
"control must reproduce the pre-fix shape (no recipient-index seek), got:\n{control_plan}"
);
}
let (_, steps_small) = {
let reader = pool.reader().unwrap();
measure_list(reader.conn(), &sql, ¶ms)
};
{
let writer = pool.writer().unwrap();
seed(writer.conn(), "actor:other", "inbound", 3_000, 10_000);
}
let (rows_after, steps_large) = {
let reader = pool.reader().unwrap();
measure_list(reader.conn(), &sql, ¶ms)
};
assert_eq!(
rows_after, 5,
"control: still only the caller's own unread rows in the result set"
);
assert!(
steps_large > steps_small.saturating_mul(3),
"control instrument did not reproduce the pre-fix growth: \
{steps_small} VM steps before vs {steps_large} after seeding another actor's \
3,000-row unread backlog (a recipient-blind scan should grow sharply)"
);
}
{
let pool = build_pool_with_comm_indexes();
{
let writer = pool.writer().unwrap();
seed(writer.conn(), "actor:a", "inbound", 5, 0);
}
let filter = listing_filter(FilterOp::EqOrLegacyIndexed);
let (sql, params) = listing_sql_and_params(&filter);
{
let reader = pool.reader().unwrap();
let fixed_plan = plan(reader.conn(), &sql, ¶ms);
assert!(
fixed_plan.contains("idx_notes_unread_probe_recipient_direction"),
"fixed listing query must be served by the recipient-scoped partial index, got:\n{fixed_plan}"
);
}
let (_, steps_small) = {
let reader = pool.reader().unwrap();
measure_list(reader.conn(), &sql, ¶ms)
};
{
let writer = pool.writer().unwrap();
seed(writer.conn(), "actor:other", "inbound", 3_000, 10_000);
}
let (rows_after, steps_large) = {
let reader = pool.reader().unwrap();
measure_list(reader.conn(), &sql, ¶ms)
};
assert_eq!(
rows_after, 5,
"fixed: still only the caller's own unread rows in the result set"
);
assert!(
steps_large < steps_small.saturating_mul(3),
"fixed listing work scaled with another actor's unread backlog: \
{steps_small} VM steps before vs {steps_large} after seeding another actor's \
3,000-row unread backlog (a recipient-scoped seek should stay flat)"
);
}
}
#[tokio::test]
async fn eq_or_legacy_indexed_matches_eq_or_missing_on_every_write_path_shape() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter};
use khive_storage::types::{PageRequest, SqlValue};
let store = setup_memory_store();
store
.upsert_note(make_note_with_props(
"default",
"message",
"exact match",
serde_json::json!({"to_actor": "actor:a"}),
))
.await
.unwrap();
store
.upsert_note(make_note_with_props(
"default",
"message",
"absent (legacy)",
serde_json::json!({}),
))
.await
.unwrap();
store
.upsert_note(make_note_with_props(
"default",
"message",
"explicit null",
serde_json::json!({"to_actor": null}),
))
.await
.unwrap();
store
.upsert_note(make_note_with_props(
"default",
"message",
"present empty string",
serde_json::json!({"to_actor": ""}),
))
.await
.unwrap();
store
.upsert_note(make_note_with_props(
"default",
"message",
"different recipient",
serde_json::json!({"to_actor": "actor:b"}),
))
.await
.unwrap();
for op in [FilterOp::EqOrMissing, FilterOp::EqOrLegacyIndexed] {
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![PropertyFilter {
json_path: "$.to_actor".to_string(),
op: op.clone(),
value: SqlValue::Text("actor:a".to_string()),
}],
..Default::default()
};
let page = store
.query_notes_filtered(
"default",
&filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
let mut contents: Vec<&str> = page.items.iter().map(|n| n.content.as_str()).collect();
contents.sort_unstable();
assert_eq!(
contents,
vec!["absent (legacy)", "exact match", "explicit null"],
"{op:?}: must match the exact recipient plus absent/null legacy rows, \
and exclude the present-empty-string and different-recipient rows"
);
}
}
fn plan_details(
conn: &rusqlite::Connection,
sql: &str,
params: &[Box<dyn rusqlite::types::ToSql>],
) -> String {
let mut stmt = conn.prepare(&format!("EXPLAIN QUERY PLAN {sql}")).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
stmt.query_map(refs.as_slice(), |row| row.get::<_, String>(3))
.unwrap()
.map(Result::unwrap)
.collect::<Vec<_>>()
.join("\n")
}
fn listing_sql_and_params(
filter: &khive_storage::note::NoteFilter,
) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
let (where_sql, mut params) = build_note_filter_read_clause("default", filter).unwrap();
params.push(Box::new(1_i64));
let limit_idx = params.len();
(
format!(
"SELECT id FROM notes{where_sql} ORDER BY created_at DESC, id ASC LIMIT ?{limit_idx}"
),
params,
)
}
#[tokio::test]
async fn gtd_tasks_status_and_assignee_query_uses_hot_property_index() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
for i in 0..20 {
store
.upsert_note(make_note_with_props(
"default",
"task",
&format!("task {i}"),
serde_json::json!({
"status": if i % 2 == 0 { "active" } else { "next" },
"assignee": if i % 3 == 0 { "lambda:a" } else { "lambda:b" },
}),
))
.await
.unwrap();
}
let filter = NoteFilter {
kind: Some("task".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("active".to_string()),
},
NotePropFilter {
json_path: "$.assignee".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("lambda:a".to_string()),
},
],
..Default::default()
};
let (sql, params) = listing_sql_and_params(&filter);
let indexed_plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
indexed_plan.contains("idx_notes_task_status")
|| indexed_plan.contains("idx_notes_task_assignee"),
"combined status+assignee query must be served by one of the new hot-property \
indexes, got plan:\n{indexed_plan}"
);
{
let writer = store.pool.writer().unwrap();
writer
.conn()
.execute_batch("DROP INDEX idx_notes_task_status; DROP INDEX idx_notes_task_assignee;")
.unwrap();
}
let control_plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
!control_plan.contains("idx_notes_task_status")
&& !control_plan.contains("idx_notes_task_assignee"),
"control: dropped indexes must vanish from the plan, got:\n{control_plan}"
);
}
#[tokio::test]
async fn gtd_next_status_in_with_assignee_query_uses_assignee_index() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
for i in 0..20 {
store
.upsert_note(make_note_with_props(
"default",
"task",
&format!("task {i}"),
serde_json::json!({
"status": if i % 2 == 0 { "active" } else { "done" },
"assignee": if i % 3 == 0 { "lambda:a" } else { "lambda:b" },
}),
))
.await
.unwrap();
}
let filter = NoteFilter {
kind: Some("task".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::In(vec![
SqlValue::Text("next".to_string()),
SqlValue::Text("active".to_string()),
]),
value: SqlValue::Null,
},
NotePropFilter {
json_path: "$.assignee".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("lambda:a".to_string()),
},
],
..Default::default()
};
let (sql, params) = listing_sql_and_params(&filter);
let indexed_plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
indexed_plan.contains("idx_notes_task_assignee"),
"the assignee equality term must be served by idx_notes_task_assignee even \
though the sibling status predicate is a multi-value IN, got plan:\n{indexed_plan}"
);
assert!(
!indexed_plan.contains("idx_notes_task_status"),
"got plan:\n{indexed_plan}"
);
}
#[tokio::test]
async fn comm_inbox_unread_default_query_still_uses_partial_unread_probe_index() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
for i in 0..20 {
store
.upsert_note(make_note_with_props(
"default",
"message",
&format!("message {i}"),
serde_json::json!({
"direction": "inbound",
"to_actor": "lambda:reader",
"read": false,
}),
))
.await
.unwrap();
}
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrLegacyIndexed,
value: SqlValue::Text("lambda:reader".to_string()),
},
],
order_by: None,
..Default::default()
};
let (sql, params) = listing_sql_and_params(&filter);
let plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
plan.contains("idx_notes_unread_probe_recipient_direction"),
"unread-default inbox listing must stay on the partial unread-probe index, got:\n{plan}"
);
assert!(
!plan.contains("idx_notes_task_status") && !plan.contains("idx_notes_task_assignee"),
"unrelated new indexes must not appear in an unread-probe plan, got:\n{plan}"
);
}
#[tokio::test]
async fn seek_after_paging_matches_offset_paging_exactly() {
use khive_storage::note::{NoteFilter, NoteSeekAfter};
use khive_storage::types::PageRequest;
let store = setup_memory_store();
for i in 0..900i64 {
let mut note = make_note("default", "widget", &format!("n{i}"));
note.created_at = i;
note.updated_at = i;
store.upsert_note(note).await.unwrap();
}
let filter = NoteFilter {
kind: Some("widget".to_string()),
..Default::default()
};
let mut offset_ids = Vec::new();
let mut offset = 0u64;
loop {
let page = store
.query_notes_filtered_count_free("default", &filter, PageRequest { limit: 97, offset })
.await
.unwrap();
if page.items.is_empty() {
break;
}
offset_ids.extend(page.items.iter().map(|n| n.id));
offset += 97;
}
let mut cursor_ids = Vec::new();
let mut cursor: Option<NoteSeekAfter> = None;
loop {
let mut page_filter = filter.clone();
page_filter.after = cursor;
let page = store
.query_notes_filtered_count_free(
"default",
&page_filter,
PageRequest {
limit: 97,
offset: 0,
},
)
.await
.unwrap();
if page.items.is_empty() {
break;
}
cursor = page.items.last().map(|n| NoteSeekAfter {
created_at: n.created_at,
id: n.id,
});
cursor_ids.extend(page.items.iter().map(|n| n.id));
}
assert_eq!(cursor_ids.len(), 900);
assert_eq!(
offset_ids, cursor_ids,
"keyset paging must reassemble the identical total order offset paging produces"
);
}
#[tokio::test]
async fn instant_ordered_window_excludes_outside_rows() {
use khive_storage::note::{FilterOp, PropertyFilter, SortDir};
let store = setup_memory_store();
for (id, at) in [
(1, "2098-12-31T23:59:59.999999999Z"),
(2, "2099-01-01T05:00:00-05:00"),
(3, "2099-01-01T10:00:00.000000001Z"),
(4, "2099-01-01T10:00:01Z"),
(5, "not-a-date"),
] {
let mut note = make_note("local", "scheduled_event", at);
note.id = Uuid::from_u128(id);
note.properties = Some(serde_json::json!({ "trigger_at": at }));
store.upsert_note(note).await.unwrap();
}
let filter = NoteFilter {
kind: Some("scheduled_event".into()),
property_filters: vec![
PropertyFilter {
json_path: "$.trigger_at".into(),
op: FilterOp::Rfc3339Valid,
value: SqlValue::Null,
},
PropertyFilter {
json_path: "$.trigger_at".into(),
op: FilterOp::Rfc3339Gte,
value: SqlValue::Text("2099-01-01T10:00:00Z".into()),
},
PropertyFilter {
json_path: "$.trigger_at".into(),
op: FilterOp::Rfc3339Lte,
value: SqlValue::Text("2099-01-01T10:00:00.000000001Z".into()),
},
],
order_by: Some(("$.trigger_at".into(), SortDir::Asc)),
order_by_instant: true,
..Default::default()
};
let page = store
.query_notes_filtered_count_free(
"local",
&filter,
PageRequest {
limit: 20,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(
page.items.iter().map(|note| note.id).collect::<Vec<_>>(),
vec![Uuid::from_u128(2), Uuid::from_u128(3)],
"SQL must return only the exact UTC window, including a crossing offset and nanosecond bound"
);
}
#[tokio::test]
async fn instant_window_handles_utc_years_outside_rfc3339_text_range() {
use khive_storage::note::{FilterOp, PropertyFilter, SortDir};
let store = setup_memory_store();
for (id, at) in [
(1, "0000-01-01T00:00:00+23:59"),
(2, "9999-12-31T23:59:59-23:59"),
] {
let mut note = make_note("local", "scheduled_event", at);
note.id = Uuid::from_u128(id);
note.properties = Some(serde_json::json!({ "trigger_at": at }));
store.upsert_note(note).await.unwrap();
}
for (id, at) in [
(1, "0000-01-01T00:00:00+23:59"),
(2, "9999-12-31T23:59:59-23:59"),
] {
let instant = at.parse::<chrono::DateTime<chrono::Utc>>().unwrap();
let filter = NoteFilter {
kind: Some("scheduled_event".into()),
property_filters: vec![
PropertyFilter {
json_path: "$.trigger_at".into(),
op: FilterOp::Rfc3339Gte,
value: SqlValue::Timestamp(instant),
},
PropertyFilter {
json_path: "$.trigger_at".into(),
op: FilterOp::Rfc3339Lte,
value: SqlValue::Timestamp(instant),
},
],
order_by: Some(("$.trigger_at".into(), SortDir::Asc)),
order_by_instant: true,
..Default::default()
};
let page = store
.query_notes_filtered_count_free(
"local",
&filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].id, Uuid::from_u128(id));
}
}
#[tokio::test]
async fn instant_seek_survives_inserts_around_cursor() {
use khive_storage::note::{FilterOp, NoteInstantSeekAfter, PropertyFilter, SortDir};
let store = setup_memory_store();
let insert = |id: u128, at: &'static str| {
let mut note = make_note("local", "scheduled_event", at);
note.id = Uuid::from_u128(id);
note.properties = Some(serde_json::json!({ "trigger_at": at }));
note
};
for note in [
insert(1, "2099-01-01T05:00:00-05:00"),
insert(2, "2099-01-01T10:00:00Z"),
insert(3, "2099-01-01T11:00:00Z"),
] {
store.upsert_note(note).await.unwrap();
}
let mut filter = NoteFilter {
kind: Some("scheduled_event".into()),
property_filters: vec![PropertyFilter {
json_path: "$.trigger_at".into(),
op: FilterOp::Rfc3339Valid,
value: SqlValue::Null,
}],
order_by: Some(("$.trigger_at".into(), SortDir::Asc)),
order_by_instant: true,
..Default::default()
};
let first = store
.query_notes_filtered_count_free(
"local",
&filter,
PageRequest {
limit: 1,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(first.items[0].id, Uuid::from_u128(1));
filter.after_instant = Some(NoteInstantSeekAfter {
value: "2099-01-01T05:00:00-05:00".into(),
id: Uuid::from_u128(1),
});
store
.upsert_note(insert(4, "2099-01-01T04:00:00-06:00"))
.await
.unwrap();
store
.upsert_note(insert(5, "2099-01-01T09:00:00-01:00"))
.await
.unwrap();
let second = store
.query_notes_filtered_count_free(
"local",
&filter,
PageRequest {
limit: 1,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(second.items[0].id, Uuid::from_u128(5));
filter.after_instant = Some(NoteInstantSeekAfter {
value: "2099-01-01T09:00:00-01:00".into(),
id: Uuid::from_u128(5),
});
let third = store
.query_notes_filtered_count_free(
"local",
&filter,
PageRequest {
limit: 2,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(
third.items.iter().map(|note| note.id).collect::<Vec<_>>(),
vec![Uuid::from_u128(2), Uuid::from_u128(3)]
);
}
#[tokio::test]
async fn seek_after_rejects_custom_order_by() {
use khive_storage::note::{NoteFilter, NoteSeekAfter, SortDir};
use khive_storage::types::PageRequest;
let store = setup_memory_store();
let filter = NoteFilter {
kind: Some("widget".to_string()),
order_by: Some(("$.priority".to_string(), SortDir::Asc)),
after: Some(NoteSeekAfter {
created_at: 0,
id: uuid::Uuid::nil(),
}),
..Default::default()
};
let err = store
.query_notes_filtered_count_free(
"default",
&filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.expect_err("after + custom order_by must be rejected");
assert!(
err.to_string().contains("order_by"),
"rejection must name the order_by conflict, got: {err}"
);
}
#[tokio::test]
async fn query_notes_filtered_rejects_seek_cursor() {
use khive_storage::note::{NoteFilter, NoteSeekAfter};
use khive_storage::types::PageRequest;
let store = setup_memory_store();
let filter = NoteFilter {
kind: Some("widget".to_string()),
after: Some(NoteSeekAfter {
created_at: 0,
id: uuid::Uuid::nil(),
}),
..Default::default()
};
let err = store
.query_notes_filtered(
"default",
&filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.expect_err("query_notes_filtered must reject NoteFilter.after");
assert!(
err.to_string().contains("query_notes_filtered_count_free"),
"rejection must point callers at the method that does honour after, got: {err}"
);
}
#[tokio::test]
async fn query_notes_filtered_rejects_seek_cursor_with_custom_order_by() {
use khive_storage::note::{NoteFilter, NoteSeekAfter, SortDir};
use khive_storage::types::PageRequest;
let store = setup_memory_store();
let filter = NoteFilter {
kind: Some("widget".to_string()),
order_by: Some(("$.priority".to_string(), SortDir::Asc)),
after: Some(NoteSeekAfter {
created_at: 0,
id: uuid::Uuid::nil(),
}),
..Default::default()
};
let err = store
.query_notes_filtered(
"default",
&filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.expect_err(
"query_notes_filtered must reject NoteFilter.after even with a custom order_by",
);
assert!(
err.to_string().contains("query_notes_filtered_count_free"),
"rejection must point callers at the method that does honour after, got: {err}"
);
}
#[tokio::test]
async fn seek_after_does_not_scale_with_rows_already_seen_unlike_offset() {
use rusqlite::StatementStatus;
let pool = setup_pool();
{
let writer = pool.writer().unwrap();
let conn = writer.conn();
conn.execute_batch("BEGIN").unwrap();
for i in 0..5_000i64 {
conn.execute(
"INSERT INTO notes (id, namespace, kind, content, properties, created_at, updated_at) \
VALUES (?1, 'default', 'task', ?2, '{\"status\":\"active\"}', ?3, ?3)",
rusqlite::params![uuid::Uuid::new_v4().to_string(), format!("n{i}"), i],
)
.unwrap();
}
conn.execute_batch("COMMIT").unwrap();
}
let filter = khive_storage::note::NoteFilter {
kind: Some("task".to_string()),
property_filters: vec![khive_storage::note::PropertyFilter {
json_path: "$.status".to_string(),
op: khive_storage::note::FilterOp::Eq,
value: khive_storage::types::SqlValue::Text("active".to_string()),
}],
..Default::default()
};
fn measure(
conn: &rusqlite::Connection,
sql: &str,
params: &[Box<dyn rusqlite::types::ToSql>],
) -> (usize, i32) {
let mut stmt = conn.prepare(sql).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
let rows: Vec<String> = stmt
.query_map(refs.as_slice(), |row| row.get::<_, String>(0))
.unwrap()
.map(Result::unwrap)
.collect();
let steps = stmt.get_status(StatementStatus::VmStep);
(rows.len(), steps)
}
let reader = pool.reader().unwrap();
let (offset_where, mut offset_params) =
build_note_filter_read_clause("default", &filter).unwrap();
offset_params.push(Box::new(100_i64));
let offset_limit_idx = offset_params.len();
offset_params.push(Box::new(4_900_i64));
let offset_offset_idx = offset_params.len();
let offset_sql = format!(
"SELECT id FROM notes{offset_where} ORDER BY created_at DESC, id ASC \
LIMIT ?{offset_limit_idx} OFFSET ?{offset_offset_idx}"
);
let (offset_rows, offset_steps) = measure(reader.conn(), &offset_sql, &offset_params);
assert_eq!(offset_rows, 100);
let (peek_where, mut peek_params) = build_note_filter_read_clause("default", &filter).unwrap();
peek_params.push(Box::new(4_900_i64));
let peek_limit_idx = peek_params.len();
let peek_sql =
format!("SELECT id, created_at FROM notes{peek_where} ORDER BY created_at DESC, id ASC LIMIT ?{peek_limit_idx}");
let boundary: (String, i64) = {
let mut stmt = reader.conn().prepare(&peek_sql).unwrap();
let refs: Vec<&dyn rusqlite::types::ToSql> =
peek_params.iter().map(|p| p.as_ref()).collect();
stmt.query_map(refs.as_slice(), |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?))
})
.unwrap()
.map(Result::unwrap)
.last()
.unwrap()
};
let after = khive_storage::note::NoteSeekAfter {
created_at: boundary.1,
id: uuid::Uuid::parse_str(&boundary.0).unwrap(),
};
let (tie_where, mut tie_params) = build_note_filter_read_clause("default", &filter).unwrap();
tie_params.push(Box::new(after.created_at));
let tie_ts_idx = tie_params.len();
tie_params.push(Box::new(after.id.to_string()));
let tie_id_idx = tie_params.len();
tie_params.push(Box::new(100_i64));
let tie_limit_idx = tie_params.len();
let tie_sql = format!(
"SELECT id FROM notes{tie_where} AND created_at = ?{tie_ts_idx} AND id > ?{tie_id_idx} \
ORDER BY id ASC LIMIT ?{tie_limit_idx}"
);
let (tie_rows, tie_steps) = measure(reader.conn(), &tie_sql, &tie_params);
let remaining = 100 - tie_rows as i64;
let (lt_where, mut lt_params) = build_note_filter_read_clause("default", &filter).unwrap();
lt_params.push(Box::new(after.created_at));
let lt_ts_idx = lt_params.len();
lt_params.push(Box::new(remaining));
let lt_limit_idx = lt_params.len();
let lt_sql = format!(
"SELECT id FROM notes{lt_where} AND created_at < ?{lt_ts_idx} \
ORDER BY created_at DESC, id ASC LIMIT ?{lt_limit_idx}"
);
let (lt_rows, lt_steps) = measure(reader.conn(), <_sql, <_params);
let after_rows = tie_rows + lt_rows;
let after_steps = tie_steps + lt_steps;
assert_eq!(
after_rows, offset_rows,
"both approaches must return the same number of rows for the same tail page"
);
let production_items =
fetch_notes_after(reader.conn(), "default", &filter, &after, 100).unwrap();
assert_eq!(production_items.len(), offset_rows);
assert!(
after_steps < offset_steps / 10,
"seeking from a cursor must avoid walking the 4,900 skipped rows an OFFSET \
re-walks on every call: {offset_steps} VM steps via OFFSET vs {after_steps} via a \
two-branch keyset seek (measured ~32x at this fixture size: 20317 vs 638)"
);
}
#[tokio::test]
async fn json_type_ne_missing_rejects_non_vocabulary_value() {
use khive_storage::note::PropertyFilter as NotePropFilter;
use khive_storage::note::{FilterOp, NoteFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true' OR '1'='1".to_string()),
}],
order_by: None,
..Default::default()
};
let err = store
.query_notes_filtered_bounded("default", &filter, 10)
.await
.expect_err("non-vocabulary json_type value must be rejected");
assert!(
err.to_string().contains("json_type"),
"rejection must name the json_type vocabulary, got: {err}"
);
}
#[tokio::test]
async fn comm_inbox_status_all_lt_branch_seeks_recipient_index() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter as NotePropFilter};
use khive_storage::types::SqlValue;
let pool = setup_pool();
{
let writer = pool.writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_comm_message_direction \
ON notes(namespace, kind, json_extract(properties, '$.direction'), \
json_extract(properties, '$.read'), created_at DESC) \
WHERE deleted_at IS NULL;
CREATE INDEX IF NOT EXISTS idx_comm_message_to_actor \
ON notes(namespace, kind, \
json_extract(properties, '$.to_actor'), \
json_extract(properties, '$.direction'), \
json_extract(properties, '$.read'), \
created_at DESC) \
WHERE deleted_at IS NULL;",
)
.unwrap();
}
{
let writer = pool.writer().unwrap();
let conn = writer.conn();
conn.execute_batch("BEGIN").unwrap();
for i in 0..6000i64 {
let to_actor = if i % 5 == 0 {
"lambda:target"
} else {
"lambda:other"
};
conn.execute(
"INSERT INTO notes (id, namespace, kind, content, properties, created_at, updated_at) \
VALUES (?1, 'default', 'message', ?2, \
json_object('direction','inbound','to_actor',?3,'read', (?4 % 2 = 0)), ?4, ?4)",
rusqlite::params![uuid::Uuid::new_v4().to_string(), format!("m{i}"), to_actor, i],
)
.unwrap();
}
conn.execute_batch("COMMIT").unwrap();
}
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrLegacyIndexed,
value: SqlValue::Text("lambda:target".to_string()),
},
],
..Default::default()
};
let reader = pool.reader().unwrap();
let (lt_where, mut lt_params) = build_note_filter_read_clause("default", &filter).unwrap();
lt_params.push(Box::new(3000_i64));
let lt_ts_idx = lt_params.len();
lt_params.push(Box::new(100_i64));
let lt_limit_idx = lt_params.len();
let lt_sql = format!(
"SELECT id FROM notes{lt_where} AND created_at < ?{lt_ts_idx} \
ORDER BY created_at DESC, id ASC LIMIT ?{lt_limit_idx}"
);
let plan = plan_details(reader.conn(), <_sql, <_params);
assert!(
plan.contains("USE TEMP B-TREE FOR ORDER BY"),
"the two recipient partitions still require a merged ordering, got:\n{plan}"
);
assert!(
plan.contains("idx_notes_message_recipient_direction")
&& plan.contains("namespace=? AND kind=? AND <expr>=? AND <expr>=? AND created_at<?"),
"the lt-branch must seek recipient, direction, and timestamp: {plan}"
);
}
#[tokio::test]
async fn candidate_created_at_id_seek_index_cannot_steal_pinned_unread_plan() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter as NotePropFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
for i in 0..20 {
store
.upsert_note(make_note_with_props(
"default",
"message",
&format!("message {i}"),
serde_json::json!({
"direction": "inbound",
"to_actor": "lambda:reader",
"read": false,
}),
))
.await
.unwrap();
}
{
let writer = store.pool.writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_notes_kind_created_seek \
ON notes(namespace, kind, created_at DESC, id ASC) \
WHERE deleted_at IS NULL;",
)
.unwrap();
}
let filter = NoteFilter {
kind: Some("message".to_string()),
property_filters: vec![
NotePropFilter {
json_path: "$.direction".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".to_string()),
},
NotePropFilter {
json_path: "$.read".to_string(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".to_string()),
},
NotePropFilter {
json_path: "$.to_actor".to_string(),
op: FilterOp::EqOrLegacyIndexed,
value: SqlValue::Text("lambda:reader".to_string()),
},
],
order_by: None,
..Default::default()
};
let (sql, params) = listing_sql_and_params(&filter);
let plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
plan.contains("idx_notes_unread_probe_recipient_direction"),
"the unread pin must survive a competing ordering index: {plan}"
);
assert!(!plan.contains("idx_notes_kind_created_seek"), "{plan}");
}
fn typed_recipient_unread_filter() -> NoteFilter {
use khive_storage::note::PropertyFilter;
NoteFilter {
kind: Some("message".into()),
property_filters: vec![
PropertyFilter {
json_path: "$.direction".into(),
op: FilterOp::Eq,
value: SqlValue::Text("inbound".into()),
},
PropertyFilter {
json_path: "$.read".into(),
op: FilterOp::JsonTypeNeMissing,
value: SqlValue::Text("true".into()),
},
PropertyFilter {
json_path: "$.to_actor".into(),
op: FilterOp::JsonTypeEq,
value: SqlValue::Text("text".into()),
},
PropertyFilter {
json_path: "$.to_actor".into(),
op: FilterOp::EqOrMissingIndexed,
value: SqlValue::Text("{}".into()),
},
],
..Default::default()
}
}
#[test]
fn typed_recipient_unread_pin_requires_complete_exact_shape() {
const TYPED: &str = " INDEXED BY idx_notes_unread_probe_recipient_type_direction WHERE ";
const UNTYPED: &str = " INDEXED BY idx_notes_unread_probe_recipient_direction WHERE ";
let filter = typed_recipient_unread_filter();
let (clause, params) = build_note_filter_read_clause("default", &filter).unwrap();
assert!(clause.starts_with(TYPED), "{clause}");
assert!(
clause.contains("json_type(properties, '$.to_actor') = ?4"),
"{clause}"
);
assert!(
clause.contains("ifnull(json_extract(properties, '$.to_actor'), '') = ?5"),
"{clause}"
);
assert_eq!(
params.len(),
5,
"the unread partial predicate stays literal"
);
let mut untyped = filter.clone();
untyped.property_filters.remove(2);
let (clause, _) = build_note_filter_read_clause("default", &untyped).unwrap();
assert!(clause.starts_with(UNTYPED), "{clause}");
let mut non_text = filter.clone();
non_text.property_filters[2].value = SqlValue::Text("object".into());
let (clause, _) = build_note_filter_read_clause("default", &non_text).unwrap();
assert!(clause.starts_with(UNTYPED), "{clause}");
let mut legacy = filter.clone();
legacy.property_filters[3].op = FilterOp::EqOrLegacyIndexed;
let (clause, _) = build_note_filter_read_clause("default", &legacy).unwrap();
assert!(clause.starts_with(UNTYPED), "{clause}");
let mut all_status = filter.clone();
all_status.property_filters.remove(1);
let (clause, _) = build_note_filter_read_clause("default", &all_status).unwrap();
assert!(
clause.starts_with(" INDEXED BY idx_notes_message_recipient_direction WHERE "),
"{clause}"
);
let mut no_direction = filter.clone();
no_direction.property_filters.remove(0);
let (clause, _) = build_note_filter_read_clause("default", &no_direction).unwrap();
assert!(clause.starts_with(" WHERE "), "{clause}");
let (where_sql, _) = build_note_filter_where("default", &filter).unwrap();
let unsafe_where = where_sql.replace(
"json_type(properties, '$.to_actor') = ?4",
"(json_type(properties, '$.to_actor') = ?4 OR 1 = 1)",
);
assert_eq!(
comm_filter_index_clause(&filter, &unsafe_where),
" INDEXED BY idx_notes_unread_probe_recipient_direction"
);
}
#[tokio::test]
async fn typed_recipient_unread_count_seeks_type_recipient_and_direction() {
let store = setup_memory_store();
let mut notes = vec![make_note_with_props(
"default",
"message",
"string recipient",
serde_json::json!({"direction":"inbound", "to_actor":"{}", "read":false}),
)];
for i in 0..128 {
for (label, properties) in [
(
"object alias",
serde_json::json!({"direction":"inbound", "to_actor":{}, "read":false}),
),
(
"other recipient",
serde_json::json!({"direction":"inbound", "to_actor":"other", "read":false}),
),
(
"outbound",
serde_json::json!({"direction":"outbound", "to_actor":"{}", "read":false}),
),
(
"read",
serde_json::json!({"direction":"inbound", "to_actor":"{}", "read":true}),
),
(
"legacy",
serde_json::json!({"direction":"inbound", "read":false}),
),
] {
notes.push(make_note_with_props(
"default",
"message",
&format!("{label} {i}"),
properties,
));
}
}
store.upsert_notes(notes).await.unwrap();
{
let writer = store.pool.writer().unwrap();
writer.conn().execute_batch(
"CREATE INDEX idx_notes_kind_created_seek ON notes(namespace, kind, created_at DESC, id ASC) WHERE deleted_at IS NULL;"
).unwrap();
}
let filter = typed_recipient_unread_filter();
let counts = store
.count_notes_filtered_bounded_in_snapshot("default", std::slice::from_ref(&filter), 1)
.await
.unwrap();
assert_eq!(counts.len(), 1);
assert_eq!(counts[0].count, 1);
assert_eq!(counts[0].cap, 1);
assert!(
!counts[0].saturated,
"malformed aliases must not saturate the real count"
);
let (where_sql, mut params) = build_note_filter_read_clause("default", &filter).unwrap();
assert!(
where_sql.starts_with(" INDEXED BY idx_notes_unread_probe_recipient_type_direction WHERE "),
"{where_sql}"
);
params.push(Box::new(1001_i64));
let sql = format!(
"SELECT COUNT(*) FROM (SELECT 1 FROM notes{where_sql} LIMIT ?{})",
params.len()
);
{
let reader = store.pool.reader().unwrap();
let plan = plan_details(reader.conn(), &sql, ¶ms);
let typed_seek = plan
.lines()
.find(|line| {
line.contains("SEARCH notes")
&& line.contains("idx_notes_unread_probe_recipient_type_direction")
})
.unwrap_or_else(|| panic!("typed index must serve a seek: {plan}"));
assert!(
typed_seek.contains("namespace=? AND kind=? AND <expr>=? AND <expr>=? AND <expr>=?"),
"type, recipient and direction must all be seek keys: {plan}"
);
assert!(!plan.contains("idx_notes_kind_created_seek"), "{plan}");
let refs: Vec<&dyn rusqlite::types::ToSql> =
params.iter().map(|value| value.as_ref()).collect();
let count: i64 = reader
.conn()
.query_row(&sql, refs.as_slice(), |row| row.get(0))
.unwrap();
assert_eq!(
count, 1,
"the explained statement must execute with its real binds"
);
}
{
let writer = store.pool.writer().unwrap();
writer
.conn()
.execute_batch("DROP INDEX idx_notes_unread_probe_recipient_type_direction")
.unwrap();
let error = writer
.conn()
.prepare(&sql)
.expect_err("missing pin must refuse preparation");
assert!(
error
.to_string()
.contains("idx_notes_unread_probe_recipient_type_direction"),
"{error}"
);
}
let mut untyped = filter;
untyped.property_filters.remove(2);
let (old_sql, old_params) = listing_sql_and_params(&untyped);
let old_plan = plan_details(store.pool.reader().unwrap().conn(), &old_sql, &old_params);
assert!(
old_plan.contains("idx_notes_unread_probe_recipient_direction"),
"{old_plan}"
);
}
#[tokio::test]
async fn narrowing_hot_property_index_to_kind_task_is_not_chosen() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter as NotePropFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
for i in 0..20 {
store
.upsert_note(make_note_with_props(
"default",
"task",
&format!("task {i}"),
serde_json::json!({ "status": if i % 2 == 0 { "active" } else { "next" } }),
))
.await
.unwrap();
}
{
let writer = store.pool.writer().unwrap();
writer
.conn()
.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_notes_task_status_narrowed \
ON notes(namespace, json_extract(properties, '$.status'), \
created_at DESC, id ASC) \
WHERE kind = 'task' AND deleted_at IS NULL;",
)
.unwrap();
}
let filter = NoteFilter {
kind: Some("task".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::Eq,
value: SqlValue::Text("active".to_string()),
}],
..Default::default()
};
let (sql, params) = listing_sql_and_params(&filter);
let plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
!plan.contains("idx_notes_task_status_narrowed"),
"expected arm: a bound `kind = ?` parameter cannot prove the narrowed \
partial index's WHERE clause, so the planner must not choose it, got:\n{plan}"
);
assert!(
plan.contains("idx_notes_task_status") || plan.contains("idx_notes_task_assignee"),
"the shipped (non-narrowed) indexes must still serve the query, got:\n{plan}"
);
}
#[tokio::test]
async fn legacy_gtd_tasks_not_in_or_missing_has_no_seek_index() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter as NotePropFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
let statuses = ["inbox", "next", "active", "done", "cancelled"];
for i in 0..20 {
store
.upsert_note(make_note_with_props(
"default",
"task",
&format!("task {i}"),
serde_json::json!({ "status": statuses[i % 5] }),
))
.await
.unwrap();
}
let filter = NoteFilter {
kind: Some("task".to_string()),
property_filters: vec![NotePropFilter {
json_path: "$.status".to_string(),
op: FilterOp::NotInOrMissing(vec![
SqlValue::Text("done".to_string()),
SqlValue::Text("cancelled".to_string()),
]),
value: SqlValue::Null,
}],
..Default::default()
};
let (sql, params) = listing_sql_and_params(&filter);
let plan = plan_details(store.pool.reader().unwrap().conn(), &sql, ¶ms);
assert!(
!plan.contains("idx_notes_task_status"),
"expected arm: an open NOT-IN exclusion cannot be served by an equality-\
keyed index seek, so idx_notes_task_status must not appear in the plan, \
got:\n{plan}"
);
}
#[tokio::test]
async fn text_in_or_non_text_filters_before_pagination_with_count_parity() {
use khive_storage::note::{FilterOp, NoteFilter, PropertyFilter};
use khive_storage::types::SqlValue;
let store = setup_memory_store();
let properties = [
serde_json::json!({}),
serde_json::json!({"status": null}),
serde_json::json!({"status": false}),
serde_json::json!({"status": 4}),
serde_json::json!({"status": []}),
serde_json::json!({"status": {}}),
serde_json::json!({"status": "inbox"}),
serde_json::json!({"status": "next"}),
serde_json::json!({"status": "archived"}),
serde_json::json!({"status": "done"}),
serde_json::json!({"status": ""}),
];
let mut ids = Vec::new();
for (index, properties) in properties.into_iter().enumerate() {
let mut note = make_note_with_props("default", "task", "predicate fixture", properties);
note.created_at = index as i64 + 1;
ids.push(note.id);
store.upsert_note(note).await.unwrap();
}
for (values, matching) in [
(
vec![
SqlValue::Text("inbox".into()),
SqlValue::Text("next".into()),
],
8,
),
(vec![], 6),
] {
let filter = NoteFilter {
kind: Some("task".into()),
property_filters: vec![PropertyFilter {
json_path: "$.status".into(),
op: FilterOp::TextInOrNonText(values),
value: SqlValue::Null,
}],
..Default::default()
};
for offset in 0..=matching {
let page = PageRequest {
limit: 1,
offset: offset as u64,
};
let counted = store
.query_notes_filtered("default", &filter, page.clone())
.await
.unwrap();
let count_free = store
.query_notes_filtered_count_free("default", &filter, page)
.await
.unwrap();
assert_eq!(counted.total, Some(matching as u64));
assert_eq!(count_free.total, None);
assert_eq!(counted.items, count_free.items);
if offset == matching {
assert!(count_free.items.is_empty());
} else {
assert_eq!(count_free.items.len(), 1);
assert_eq!(count_free.items[0].id, ids[matching - offset - 1]);
}
}
}
}