use super::*;
use crate::pool::PoolConfig;
use serde_json::json;
fn setup_memory_store() -> SqlEventStore {
let config = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = Arc::new(ConnectionPool::new(config).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(EVENTS_DDL).unwrap();
}
SqlEventStore::new_scoped(pool, false, "default")
}
fn make_event(namespace: &str) -> Event {
Event::new(
namespace,
"search",
EventKind::SearchExecuted,
SubstrateKind::Note,
"agent:test",
)
.with_payload(json!({ "result_kind": "note" }))
}
#[tokio::test]
async fn append_usage_mark_distinguishes_typed_writer_outcomes() {
use khive_storage::usage::{scope, UsageContext, UsageUnit};
for state in [
WriterTaskRequestState::NotStarted,
WriterTaskRequestState::TransactionRolledBack,
WriterTaskRequestState::SideEffectsUnknown,
] {
for error in [
StorageError::WriterTaskTerminated {
request_state: state,
},
StorageError::WriterTaskRequestFailed {
request_state: state,
source: Box::new(StorageError::Pool {
operation: "append_event".into(),
message: "write failed".into(),
}),
},
] {
let ctx = UsageContext::new();
ctx.add(UsageUnit::EventRows, 1);
scope(ctx.clone(), async { mark_unknown_append_usage(&error) }).await;
assert_eq!(
ctx.shipping_snapshot().is_none(),
state == WriterTaskRequestState::SideEffectsUnknown,
"{error:?}"
);
assert_eq!(ctx.snapshot()["event_rows"], 1);
}
}
let ctx = UsageContext::new();
scope(ctx.clone(), async {
mark_unknown_append_usage(&StorageError::driver(
StorageCapability::Events,
"append_event",
std::io::Error::other("write failed"),
));
})
.await;
assert_eq!(ctx.shipping_snapshot(), Some(json!({})));
mark_unknown_append_usage(&StorageError::WriterTaskTerminated {
request_state: WriterTaskRequestState::SideEffectsUnknown,
});
}
async fn observations_for(store: &SqlEventStore, event_id: Uuid) -> Vec<EventObservation> {
let pool = Arc::clone(&store.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.reader().unwrap();
fetch_event_observations(guard.conn(), event_id).unwrap()
})
.await
.unwrap()
}
#[tokio::test]
async fn operation_attribution_round_trips_and_changes_idempotent_identity() {
use khive_storage::operation_context::scope_operation_attribution;
use khive_types::{OperationAttribution, RefResolution};
let store = setup_memory_store();
let direct = make_event("default");
assert_eq!((direct.op_index, direct.ref_resolution), (None, None));
store.append_event(direct.clone()).await.unwrap();
let attributed = scope_operation_attribution(
OperationAttribution {
op_index: u32::MAX,
ref_resolution: RefResolution::Resolved,
},
async { make_event("default") },
)
.await;
store.append_event(attributed.clone()).await.unwrap();
for expected in [&direct, &attributed] {
let stored = store.get_event(expected.id).await.unwrap().unwrap();
assert_eq!(&stored, expected);
}
let page = store
.query_events(
EventFilter::default(),
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert!(page.items.contains(&direct));
assert!(page.items.contains(&attributed));
let mut changed = attributed;
changed.op_index = Some(0);
let result = store.append_events_idempotent(vec![changed]).await.unwrap();
assert_eq!(result.rows[0], EventAppendDisposition::IdentityConflict);
}
const PRE_ATTRIBUTION_EVENTS_TABLE: &str = "CREATE TABLE events (
id TEXT PRIMARY KEY, namespace TEXT NOT NULL, verb TEXT NOT NULL, substrate TEXT NOT NULL,
actor TEXT NOT NULL, kind TEXT NOT NULL DEFAULT 'audit', outcome TEXT NOT NULL,
payload TEXT NOT NULL DEFAULT '{}', payload_schema_version INTEGER NOT NULL DEFAULT 1,
profile_state_version INTEGER, duration_us INTEGER NOT NULL DEFAULT 0, target_id TEXT,
session_id TEXT, aggregate_kind TEXT, aggregate_id TEXT, created_at INTEGER NOT NULL)";
#[tokio::test]
async fn store_schema_adds_operation_attribution_to_a_pre_attribution_events_table() {
use khive_storage::operation_context::scope_operation_attribution;
use khive_types::{OperationAttribution, RefResolution};
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: None,
..PoolConfig::default()
})
.unwrap(),
);
let legacy = make_event("default");
{
let writer = pool.writer().unwrap();
writer
.conn()
.execute_batch(PRE_ATTRIBUTION_EVENTS_TABLE)
.unwrap();
writer
.conn()
.execute(
"INSERT INTO events \
(id, namespace, verb, substrate, actor, kind, outcome, payload, created_at) \
VALUES (?1, ?2, ?3, 'note', ?4, 'search_executed', 'success', ?5, ?6)",
rusqlite::params![
legacy.id.to_string(),
legacy.namespace,
legacy.verb,
legacy.actor,
legacy.payload.to_string(),
legacy.created_at,
],
)
.unwrap();
ensure_events_schema(writer.conn()).unwrap();
ensure_events_schema(writer.conn()).unwrap();
}
let store = SqlEventStore::new_scoped(pool, false, "default");
let attributed = scope_operation_attribution(
OperationAttribution {
op_index: 3,
ref_resolution: RefResolution::Literal,
},
async { make_event("default") },
)
.await;
store.append_event(attributed.clone()).await.unwrap();
assert_eq!(
store.get_event(attributed.id).await.unwrap().unwrap(),
attributed
);
assert_eq!(store.get_event(legacy.id).await.unwrap().unwrap(), legacy);
}
#[test]
fn store_schema_refuses_an_events_table_with_one_attribution_column() {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.execute_batch(PRE_ATTRIBUTION_EVENTS_TABLE).unwrap();
conn.execute_batch("ALTER TABLE events ADD COLUMN op_index INTEGER")
.unwrap();
let error = ensure_events_schema(&conn).unwrap_err();
assert!(
error
.to_string()
.contains("only one operation attribution column"),
"{error}"
);
}
#[tokio::test]
async fn operation_attribution_rejects_unpaired_values_before_append() {
let store = setup_memory_store();
for (op_index, ref_resolution) in [
(Some(0), None),
(None, Some(khive_types::RefResolution::Literal)),
] {
let mut event = make_event("default");
event.op_index = op_index;
event.ref_resolution = ref_resolution;
assert!(store.preflight_event(&event).is_err());
assert!(store.append_event(event).await.is_err());
}
assert_eq!(store.count_events(EventFilter::default()).await.unwrap(), 0);
}
#[tokio::test]
async fn profile_state_version_refuses_overflow_before_any_event_is_persisted() {
let store = setup_memory_store();
let candidate = Uuid::new_v4();
let mut valid = make_event("default").with_profile_state_version(i64::MAX as u64);
valid.payload = json!({
"result_kind": "note",
"candidates": [candidate.to_string()],
});
let overflow = make_event("default").with_profile_state_version(i64::MAX as u64 + 1);
assert_eq!(event_insert_statements(&valid).unwrap().len(), 2);
assert!(event_insert_statements(&overflow).is_err());
assert!(store.preflight_event(&overflow).is_err());
assert!(store.append_event(overflow.clone()).await.is_err());
assert!(store
.append_events(vec![valid.clone(), overflow.clone()])
.await
.is_err());
assert!(observations_for(&store, valid.id).await.is_empty());
assert!(store
.append_events_idempotent(vec![valid.clone(), overflow.clone()])
.await
.is_err());
assert_eq!(store.count_events(EventFilter::default()).await.unwrap(), 0);
assert!(observations_for(&store, valid.id).await.is_empty());
store.append_event(valid.clone()).await.unwrap();
assert_eq!(observations_for(&store, valid.id).await.len(), 1);
assert_eq!(
store.get_event(valid.id).await.unwrap(),
Some(valid.clone())
);
let mut invalid_retry = valid;
invalid_retry.profile_state_version = overflow.profile_state_version;
assert!(store
.append_events_idempotent(vec![invalid_retry])
.await
.is_err());
assert_eq!(store.count_events(EventFilter::default()).await.unwrap(), 1);
assert_eq!(
store
.query_events(
EventFilter::default(),
PageRequest {
limit: 10,
offset: 0
}
)
.await
.unwrap()
.items
.len(),
1
);
}
#[test]
fn profile_state_version_builder_binds_input_value() {
let event = make_event("default").with_profile_state_version(41);
let statements = event_insert_statements(&event).unwrap();
assert_eq!(statements.len(), 1);
match statements[0].params.get(9) {
Some(SqlValue::Integer(value)) => assert_eq!(*value, 41),
other => panic!("expected bound profile_state_version 41, got {other:?}"),
}
}
#[tokio::test]
async fn prepared_event_insert_stores_and_reads_i64_max_profile_state_version() {
let store = setup_memory_store();
let event = make_event("default").with_profile_state_version(i64::MAX as u64);
let statements = event_insert_statements(&event).expect("boundary value is representable");
{
let writer = store.pool.writer().unwrap();
let conn = writer.conn();
for statement in statements {
let mut prepared = conn.prepare(&statement.sql).unwrap();
crate::sql_bridge::bind_params(&mut prepared, &statement.params).unwrap();
assert_eq!(prepared.raw_execute().unwrap(), 1);
}
let stored: i64 = conn
.query_row(
"SELECT profile_state_version FROM events WHERE id = ?1",
[event.id.to_string()],
|row| row.get(0),
)
.unwrap();
assert_eq!(stored, i64::MAX);
}
assert_eq!(store.get_event(event.id).await.unwrap(), Some(event));
}
#[tokio::test]
async fn batch_failure_rolls_back_event_and_observation_rows() {
let store = setup_memory_store();
let candidate = Uuid::new_v4();
let mut first = make_event("default");
first.payload = json!({
"result_kind": "note",
"candidates": [candidate.to_string()],
});
assert_eq!(decode_event_observations(&first).unwrap().len(), 1);
let mut rejected = make_event("default");
rejected.payload = json!({"result_kind": "edge"});
assert!(store
.append_events(vec![first.clone(), rejected.clone()])
.await
.is_err());
assert_eq!(store.count_events(EventFilter::default()).await.unwrap(), 0);
assert!(observations_for(&store, first.id).await.is_empty());
assert!(observations_for(&store, rejected.id).await.is_empty());
}
#[tokio::test]
async fn profile_state_version_audit_finds_preexisting_unreadable_rows() {
let store = setup_memory_store();
let valid_null = make_event("default");
let valid_max = make_event("default").with_profile_state_version(i64::MAX as u64);
let negative = make_event("default");
let text = make_event("default");
let real = make_event("default");
for event in [
valid_null.clone(),
valid_max.clone(),
negative.clone(),
text.clone(),
real.clone(),
] {
store.append_event(event).await.unwrap();
}
{
let writer = store.pool.writer().unwrap();
let conn = writer.conn();
for (id, value) in [
(negative.id, rusqlite::types::Value::Integer(-1)),
(
text.id,
rusqlite::types::Value::Text("not-an-integer".into()),
),
(real.id, rusqlite::types::Value::Real(1.5)),
] {
assert_eq!(
conn.execute(
"UPDATE events SET profile_state_version = ?1 WHERE id = ?2",
rusqlite::params![value, id.to_string()],
)
.unwrap(),
1,
);
}
}
let rows: Vec<(String, String, String, String)> = {
let reader = store.pool.reader().unwrap();
let mut statement = reader
.conn()
.prepare(include_str!(
"../../docs/api/event-profile-state-version-audit.sql"
))
.unwrap();
let found = statement
.query_map([], |row| {
Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?))
})
.unwrap()
.collect::<Result<Vec<_>, _>>()
.unwrap();
found
};
let mut expected = vec![
(
negative.id.to_string(),
"default".to_string(),
"integer".to_string(),
"-1".to_string(),
),
(
text.id.to_string(),
"default".to_string(),
"text".to_string(),
"'not-an-integer'".to_string(),
),
(
real.id.to_string(),
"default".to_string(),
"real".to_string(),
"1.5".to_string(),
),
];
expected.sort_by(|left, right| left.0.cmp(&right.0));
assert_eq!(rows, expected);
assert_eq!(
store.get_event(valid_null.id).await.unwrap(),
Some(valid_null)
);
assert_eq!(
store.get_event(valid_max.id).await.unwrap(),
Some(valid_max)
);
for id in [negative.id, text.id, real.id] {
assert!(store.get_event(id).await.is_err());
}
assert!(store
.query_events(
EventFilter::default(),
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.is_err());
}
#[tokio::test]
async fn test_append_and_get_event() {
let store = setup_memory_store();
let event = make_event("default");
let id = event.id;
store.append_event(event).await.unwrap();
let fetched = store.get_event(id).await.unwrap();
assert!(fetched.is_some());
let fetched = fetched.unwrap();
assert_eq!(fetched.id, id);
assert_eq!(fetched.verb, "search");
assert_eq!(fetched.substrate, SubstrateKind::Note);
assert_eq!(fetched.actor, "agent:test");
assert_eq!(fetched.outcome, EventOutcome::Success);
}
#[tokio::test]
async fn test_append_events_batch() {
let store = setup_memory_store();
let events: Vec<Event> = (0..3).map(|_| make_event("default")).collect();
let summary = store.append_events(events).await.unwrap();
assert_eq!(summary.attempted, 3);
assert_eq!(summary.affected, 3);
assert_eq!(summary.failed, 0);
}
#[tokio::test]
async fn test_count_events() {
let store = setup_memory_store();
for _ in 0..3 {
store.append_event(make_event("default")).await.unwrap();
}
let count = store.count_events(EventFilter::default()).await.unwrap();
assert_eq!(count, 3);
}
#[tokio::test]
async fn test_query_events_filter_by_verb() {
let store = setup_memory_store();
store.append_event(make_event("default")).await.unwrap();
let mut create_event = make_event("default");
create_event.verb = "create".to_string();
store.append_event(create_event).await.unwrap();
let filter = EventFilter {
verbs: vec!["search".to_string()],
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].verb, "search");
}
#[tokio::test]
async fn test_query_events_filter_by_substrate() {
let store = setup_memory_store();
store.append_event(make_event("default")).await.unwrap();
let mut entity_event = make_event("default");
entity_event.substrate = SubstrateKind::Entity;
store.append_event(entity_event).await.unwrap();
let filter = EventFilter {
substrates: vec![SubstrateKind::Entity],
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].substrate, SubstrateKind::Entity);
}
#[tokio::test]
async fn test_outcome_roundtrip() {
let store = setup_memory_store();
let mut denied = make_event("default");
denied.outcome = EventOutcome::Denied;
let denied_id = denied.id;
store.append_event(denied).await.unwrap();
let fetched = store.get_event(denied_id).await.unwrap().unwrap();
assert_eq!(fetched.outcome, EventOutcome::Denied);
}
#[tokio::test]
async fn append_event_writes_observations_atomically() {
let store = setup_memory_store();
let candidate = Uuid::new_v4();
let selected = Uuid::new_v4();
let mut event = make_event("default");
event.kind = EventKind::SearchExecuted;
event.payload = json!({
"result_kind": "note",
"candidates": [candidate.to_string()],
"selected": [selected.to_string()],
"served_by_profile_id": "profile-a"
});
let event_id = event.id;
store.append_event(event).await.unwrap();
let fetched = store.get_event(event_id).await.unwrap();
assert!(fetched.is_some());
let pool = Arc::clone(&store.pool);
let event_id_str = event_id.to_string();
let (candidate_count, selected_count) = tokio::task::spawn_blocking(move || {
let guard = pool.reader().unwrap();
let conn = guard.conn();
let c: i64 = conn
.query_row(
"SELECT COUNT(*) FROM event_observations WHERE event_id = ?1 AND role = 'candidate'",
[&event_id_str],
|r| r.get(0),
)
.unwrap();
let s: i64 = conn
.query_row(
"SELECT COUNT(*) FROM event_observations WHERE event_id = ?1 AND role = 'selected'",
[&event_id_str],
|r| r.get(0),
)
.unwrap();
(c, s)
})
.await
.unwrap();
assert_eq!(candidate_count, 1, "expected one candidate observation row");
assert_eq!(selected_count, 1, "expected one selected observation row");
}
#[tokio::test]
async fn search_executed_rejects_unknown_result_kind() {
let store = setup_memory_store();
let mut event = make_event("default");
event.payload = json!({
"result_kind": "edge",
"candidates": [Uuid::new_v4().to_string()],
"selected": []
});
let event_id = event.id;
let result = store.append_event(event).await;
assert!(result.is_err(), "unknown result_kind must be rejected");
assert!(
store.get_event(event_id).await.unwrap().is_none(),
"invalid event and projection must roll back atomically"
);
}
#[tokio::test]
async fn search_executed_absent_result_kind_projects_historical_note_rows() {
let store = setup_memory_store();
let mut event = make_event("default");
let note_id = Uuid::new_v4();
event.payload = json!({
"candidates": [note_id.to_string()],
"selected": [note_id.to_string()]
});
let event_id = event.id;
store.preflight_event(&event).unwrap();
event_insert_statements(&event).unwrap();
store.append_event(event.clone()).await.unwrap();
let event_id_str = event_id.to_string();
let reader = store.pool.reader().unwrap();
let mut stmt = reader
.conn()
.prepare(
"SELECT entity_id, referent_kind, role FROM event_observations \
WHERE event_id = ?1 ORDER BY role, position",
)
.unwrap();
let rows = stmt
.query_map([&event_id_str], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
))
})
.unwrap()
.collect::<Result<Vec<_>, _>>()
.unwrap();
assert_eq!(
rows,
vec![
(note_id.to_string(), "note".into(), "candidate".into()),
(note_id.to_string(), "note".into(), "selected".into()),
]
);
drop(stmt);
drop(reader);
let replay = store.append_events_idempotent(vec![event]).await.unwrap();
assert_eq!(
replay.rows,
vec![EventAppendDisposition::AlreadyPresentIdentical]
);
}
#[tokio::test]
async fn search_executed_rejects_non_object_payload_root() {
let store = setup_memory_store();
for payload in [json!("not an object"), json!([]), json!(null)] {
let mut event = make_event("default");
event.payload = payload;
let event_id = event.id;
assert!(store.preflight_event(&event).is_err());
assert!(event_insert_statements(&event).is_err());
assert!(store.append_event(event).await.is_err());
assert!(store.get_event(event_id).await.unwrap().is_none());
}
}
#[tokio::test]
async fn search_executed_rejects_present_non_string_result_kind() {
let store = setup_memory_store();
for result_kind in [json!(null), json!(0), json!(["note"])] {
let mut event = make_event("default");
event.payload = json!({"result_kind": result_kind, "candidates": [], "selected": []});
let event_id = event.id;
assert!(store.preflight_event(&event).is_err());
assert!(event_insert_statements(&event).is_err());
assert!(store.append_event(event).await.is_err());
assert!(store.get_event(event_id).await.unwrap().is_none());
}
}
async fn selected_uuids_for(store: &SqlEventStore, event_id: Uuid) -> Vec<String> {
let pool = Arc::clone(&store.pool);
let event_id_str = event_id.to_string();
tokio::task::spawn_blocking(move || {
let guard = pool.reader().unwrap();
let conn = guard.conn();
let mut stmt = conn
.prepare(
"SELECT entity_id FROM event_observations \
WHERE event_id = ?1 AND role = 'selected' ORDER BY position",
)
.unwrap();
stmt.query_map([&event_id_str], |r| r.get(0))
.unwrap()
.collect::<Result<Vec<String>, _>>()
.unwrap()
})
.await
.unwrap()
}
#[tokio::test]
async fn rerank_executed_falls_back_to_reranked_when_final_scores_field_is_absent() {
let store = setup_memory_store();
let candidate = Uuid::new_v4();
let reranked_winner = Uuid::new_v4();
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = json!({
"candidates": [candidate.to_string()],
"reranked": [[reranked_winner.to_string(), [["relevance", 0.9]]]],
});
let event_id = event.id;
store.append_event(event).await.unwrap();
let selected = selected_uuids_for(&store, event_id).await;
assert_eq!(
selected,
vec![reranked_winner.to_string()],
"selected observation must decode the UUID leading the `reranked` tuple \
when `final_scores` is absent"
);
}
#[tokio::test]
async fn rerank_executed_prefers_final_scores_order_over_reranked_when_both_present() {
let store = setup_memory_store();
let a = Uuid::new_v4();
let b = Uuid::new_v4();
let payload = khive_types::RerankExecutedPayload {
served_by_profile_id: Some("profile-a".to_string()),
model_id: khive_types::Id128::from_u128(1),
candidates: vec![
khive_types::Id128::from_bytes(*a.as_bytes()),
khive_types::Id128::from_bytes(*b.as_bytes()),
],
reranked: vec![
(
khive_types::Id128::from_bytes(*b.as_bytes()),
vec![("relevance".to_string(), 0.4)],
),
(
khive_types::Id128::from_bytes(*a.as_bytes()),
vec![("relevance".to_string(), 0.9)],
),
],
final_scores: vec![
(khive_types::Id128::from_bytes(*a.as_bytes()), 0.9),
(khive_types::Id128::from_bytes(*b.as_bytes()), 0.4),
],
latency_us: 1200,
hook_applied: false,
hook_target_match: false,
};
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = serde_json::to_value(&payload).unwrap();
let event_id = event.id;
store.append_event(event).await.unwrap();
let selected = selected_uuids_for(&store, event_id).await;
assert_eq!(
selected,
vec![a.to_string(), b.to_string()],
"selected order must follow `final_scores`, not `reranked`"
);
}
#[tokio::test]
async fn rerank_executed_ignores_stray_selected_field_and_uses_final_scores() {
let store = setup_memory_store();
let a = Uuid::new_v4();
let b = Uuid::new_v4();
let stray_selected_winner = Uuid::new_v4();
let payload = khive_types::RerankExecutedPayload {
served_by_profile_id: Some("profile-a".to_string()),
model_id: khive_types::Id128::from_u128(1),
candidates: vec![
khive_types::Id128::from_bytes(*a.as_bytes()),
khive_types::Id128::from_bytes(*b.as_bytes()),
],
reranked: vec![],
final_scores: vec![
(khive_types::Id128::from_bytes(*a.as_bytes()), 0.9),
(khive_types::Id128::from_bytes(*b.as_bytes()), 0.4),
],
latency_us: 1200,
hook_applied: false,
hook_target_match: false,
};
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
let mut payload_value = serde_json::to_value(&payload).unwrap();
payload_value.as_object_mut().unwrap().insert(
"selected".to_string(),
json!([stray_selected_winner.to_string()]),
);
event.payload = payload_value;
let event_id = event.id;
store.append_event(event).await.unwrap();
let selected = selected_uuids_for(&store, event_id).await;
assert_eq!(
selected,
vec![a.to_string(), b.to_string()],
"stray `selected` field must be ignored for RerankExecuted; \
projection must follow `final_scores` only"
);
}
#[tokio::test]
async fn rerank_executed_uses_final_scores_when_reranked_is_empty() {
let store = setup_memory_store();
let winner = Uuid::new_v4();
let payload = khive_types::RerankExecutedPayload {
served_by_profile_id: None,
model_id: khive_types::Id128::from_u128(1),
candidates: vec![khive_types::Id128::from_bytes(*winner.as_bytes())],
reranked: vec![],
final_scores: vec![(khive_types::Id128::from_bytes(*winner.as_bytes()), 0.75)],
latency_us: 800,
hook_applied: false,
hook_target_match: false,
};
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = serde_json::to_value(&payload).unwrap();
let event_id = event.id;
store.append_event(event).await.unwrap();
let selected = selected_uuids_for(&store, event_id).await;
assert_eq!(
selected,
vec![winner.to_string()],
"final_scores must be used even when reranked is present-but-empty"
);
}
#[tokio::test]
async fn rerank_executed_final_scores_single_element_tuple_rejected() {
let store = setup_memory_store();
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = json!({
"candidates": [],
"final_scores": [[Uuid::new_v4().to_string()]],
});
let event_id = event.id;
let result = store.append_event(event).await;
assert!(result.is_err(), "single-element tuple must be rejected");
let fetched = store.get_event(event_id).await.unwrap();
assert!(fetched.is_none(), "event row must not exist after rollback");
}
#[tokio::test]
async fn rerank_executed_final_scores_extra_element_tuple_rejected() {
let store = setup_memory_store();
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = json!({
"candidates": [],
"final_scores": [[Uuid::new_v4().to_string(), 0.5, "extra"]],
});
let event_id = event.id;
let result = store.append_event(event).await;
assert!(result.is_err(), "extra-element tuple must be rejected");
let fetched = store.get_event(event_id).await.unwrap();
assert!(fetched.is_none(), "event row must not exist after rollback");
}
#[tokio::test]
async fn rerank_executed_final_scores_non_numeric_second_element_rejected() {
let store = setup_memory_store();
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = json!({
"candidates": [],
"final_scores": [[Uuid::new_v4().to_string(), "not-a-number"]],
});
let event_id = event.id;
let result = store.append_event(event).await;
assert!(result.is_err(), "non-numeric score must be rejected");
let fetched = store.get_event(event_id).await.unwrap();
assert!(fetched.is_none(), "event row must not exist after rollback");
}
#[tokio::test]
async fn rerank_executed_reranked_malformed_sub_scores_rejected() {
let store = setup_memory_store();
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = json!({
"candidates": [],
"reranked": [[Uuid::new_v4().to_string(), "not-an-array"]],
});
let event_id = event.id;
let result = store.append_event(event).await;
assert!(
result.is_err(),
"malformed reranked sub-scores must be rejected"
);
let fetched = store.get_event(event_id).await.unwrap();
assert!(fetched.is_none(), "event row must not exist after rollback");
}
#[tokio::test]
async fn feedback_explicit_projects_signal_observation_from_target_id() {
let store = setup_memory_store();
let target = Uuid::new_v4();
let event = Event::new(
"default",
"brain.feedback",
EventKind::FeedbackExplicit,
SubstrateKind::Event,
"agent:test",
)
.with_target(target)
.with_payload(json!({ "signal": "useful" }));
let event_id = event.id;
store.append_event(event).await.unwrap();
let pool = Arc::clone(&store.pool);
let event_id_str = event_id.to_string();
let (signal_count, observed_entity_id): (i64, String) =
tokio::task::spawn_blocking(move || {
let guard = pool.reader().unwrap();
let conn = guard.conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM event_observations WHERE event_id = ?1 AND role = 'signal'",
[&event_id_str],
|r| r.get(0),
)
.unwrap();
let entity_id: String = conn
.query_row(
"SELECT entity_id FROM event_observations WHERE event_id = ?1 AND role = 'signal'",
[&event_id_str],
|r| r.get(0),
)
.unwrap();
(count, entity_id)
})
.await
.unwrap();
assert_eq!(signal_count, 1, "expected one signal observation row");
assert_eq!(
observed_entity_id,
target.to_string(),
"signal observation must carry the feedback's target_id"
);
}
#[tokio::test]
async fn invalid_projection_payload_aborts_event_insert() {
let store = setup_memory_store();
let mut event = make_event("default");
event.kind = EventKind::RerankExecuted;
event.payload = json!({ "candidates": "not-array" });
let event_id = event.id;
let result = store.append_event(event).await;
assert!(result.is_err(), "invalid payload must return Err");
let fetched = store.get_event(event_id).await.unwrap();
assert!(fetched.is_none(), "event row must not exist after rollback");
}
#[tokio::test]
async fn query_events_orders_by_created_at_then_id_desc() {
let store = setup_memory_store();
let ts = 1_700_000_000_000_000_i64;
let id_low = Uuid::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
let id_high = Uuid::parse_str("ffffffff-ffff-ffff-ffff-ffffffffffff").unwrap();
let pool = Arc::clone(&store.pool);
tokio::task::spawn_blocking(move || {
let guard = pool.try_writer().unwrap();
let conn = guard.conn();
conn.execute_batch("BEGIN IMMEDIATE").unwrap();
for id in [id_low, id_high] {
conn.execute(
"INSERT INTO events \
(id, namespace, verb, substrate, actor, kind, outcome, payload, \
payload_schema_version, duration_us, created_at) \
VALUES (?1, 'default', 'search', 'note', 'test', 'audit', 'success', '{}', 1, 0, ?2)",
rusqlite::params![id.to_string(), ts],
)
.unwrap();
}
conn.execute_batch("COMMIT").unwrap();
})
.await
.unwrap();
let page = store
.query_events(
EventFilter::default(),
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 2);
assert_eq!(
page.items[0].id, id_high,
"higher UUID must come first (id DESC tiebreaker)"
);
assert_eq!(page.items[1].id, id_low);
let newer_low = Uuid::from_u128(2);
let newer_high = Uuid::from_u128(3);
let older_low = Uuid::from_u128(4);
let older_high = Uuid::from_u128(5);
let rows = [
(older_high, ts - 10),
(newer_low, ts + 10),
(older_low, ts - 10),
(newer_high, ts + 10),
];
let events = rows
.into_iter()
.map(|(id, created_at)| {
let mut event = make_event("default");
event.id = id;
event.created_at = created_at;
event
})
.collect();
let written = store.append_events(events).await.unwrap();
assert_eq!(written.attempted, 4);
assert_eq!(written.affected, 4);
assert_eq!(written.failed, 0);
let expected = [
newer_high, newer_low, id_high, id_low, older_high, older_low,
];
let whole = store
.query_events(
EventFilter::default(),
PageRequest {
offset: 0,
limit: 10,
},
)
.await
.unwrap();
assert_eq!(
whole.items.iter().map(|event| event.id).collect::<Vec<_>>(),
expected
);
for limit in [1_u32, 2, 3, 4] {
let mut seen = Vec::new();
for offset in (0..expected.len()).step_by(limit as usize) {
let page = store
.query_events(
EventFilter::default(),
PageRequest {
offset: offset as u64,
limit,
},
)
.await
.unwrap();
let ids: Vec<_> = page.items.iter().map(|event| event.id).collect();
let end = (offset + limit as usize).min(expected.len());
assert_eq!(ids, expected[offset..end], "offset={offset} limit={limit}");
seen.extend(ids);
}
assert_eq!(seen, expected, "complete pagination at limit={limit}");
let terminal = store
.query_events(
EventFilter::default(),
PageRequest {
offset: expected.len() as u64,
limit,
},
)
.await
.unwrap();
assert!(terminal.items.is_empty(), "terminal page at limit={limit}");
}
}
#[tokio::test]
async fn query_events_filters_by_kind() {
let store = setup_memory_store();
store.append_event(make_event("default")).await.unwrap();
let mut recall_event = make_event("default");
recall_event.kind = EventKind::RecallExecuted;
store.append_event(recall_event).await.unwrap();
let filter = EventFilter {
kinds: vec![EventKind::RecallExecuted],
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].kind, EventKind::RecallExecuted);
}
#[tokio::test]
async fn query_events_filters_by_session_id() {
let store = setup_memory_store();
let session = Uuid::new_v4();
let mut event = make_event("default");
event.session_id = Some(session);
store.append_event(event).await.unwrap();
store.append_event(make_event("default")).await.unwrap();
let filter = EventFilter {
session_id: Some(session),
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].session_id, Some(session));
}
#[tokio::test]
async fn query_events_filters_by_observed() {
let store = setup_memory_store();
let entity_id = Uuid::new_v4();
let mut event = make_event("default");
event.kind = EventKind::SearchExecuted;
event.payload = json!({
"result_kind": "entity",
"candidates": [entity_id.to_string()],
"selected": []
});
store.append_event(event).await.unwrap();
store.append_event(make_event("default")).await.unwrap();
let filter = EventFilter {
observed: vec![entity_id],
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
}
#[tokio::test]
async fn link_events_project_edge_referents_and_observed_matches_any_role() {
let store = setup_memory_store();
let edge_id = Uuid::new_v4();
let source_id = Uuid::new_v4();
let target_id = Uuid::new_v4();
let event = Event::new(
"default",
"link",
EventKind::LinkCreated,
SubstrateKind::Entity,
"agent:test",
)
.with_target(edge_id)
.with_payload(json!({
"id": edge_id,
"source_id": source_id,
"target_id": target_id,
"source_kind": "entity",
"target_kind": "entity",
"mutation": "created"
}));
let event_id = event.id;
let mut legacy = event.clone();
legacy
.payload
.as_object_mut()
.unwrap()
.remove("source_kind");
legacy
.payload
.as_object_mut()
.unwrap()
.remove("target_kind");
assert_eq!(
decode_link_observations(&event).unwrap(),
decode_link_observations(&legacy).unwrap(),
"explicit entity endpoints must preserve today's observation bytes"
);
store.append_event(event).await.unwrap();
for observed in [source_id, target_id, edge_id] {
let page = store
.query_events(
EventFilter {
observed: vec![observed],
..EventFilter::default()
},
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].id, event_id);
}
let rows = observations_for(&store, event_id).await;
assert_eq!(rows.len(), 3);
for (position, id, kind) in [
(0, source_id, ReferentKind::Entity),
(1, target_id, ReferentKind::Entity),
(2, edge_id, ReferentKind::Edge),
] {
assert_eq!(rows[position].entity_id, id);
assert_eq!(rows[position].referent_kind, kind);
assert_eq!(rows[position].role, ObservationRole::Target);
assert_eq!(rows[position].position, position as u32);
}
}
#[tokio::test]
async fn link_note_endpoints_project_note_referents() {
let store = setup_memory_store();
let source_id = Uuid::new_v4();
let target_id = Uuid::new_v4();
let edge_id = Uuid::new_v4();
let event = Event::new(
"default",
"link",
EventKind::LinkCreated,
SubstrateKind::Entity,
"agent:test",
)
.with_target(edge_id)
.with_payload(json!({
"id": edge_id,
"source_id": source_id,
"target_id": target_id,
"source_kind": "note",
"target_kind": "note",
"relation": "supports"
}));
let event_id = event.id;
store.append_event(event).await.unwrap();
let rows = observations_for(&store, event_id).await;
assert_eq!(rows.len(), 3);
assert_eq!(rows[0].entity_id, source_id);
assert_eq!(rows[0].referent_kind, ReferentKind::Note);
assert_eq!(rows[1].entity_id, target_id);
assert_eq!(rows[1].referent_kind, ReferentKind::Note);
assert_eq!(rows[2].entity_id, edge_id);
assert_eq!(rows[2].referent_kind, ReferentKind::Edge);
}
#[tokio::test]
async fn legacy_link_payload_replays_with_entity_endpoint_fallback() {
let store = setup_memory_store();
let source_id = Uuid::new_v4();
let target_id = Uuid::new_v4();
let edge_id = Uuid::new_v4();
let event = Event::new(
"default",
"link",
EventKind::LinkCreated,
SubstrateKind::Entity,
"agent:test",
)
.with_target(edge_id)
.with_payload(json!({
"id": edge_id,
"source_id": source_id,
"target_id": target_id
}));
let first = store
.append_events_idempotent(vec![event.clone()])
.await
.unwrap();
assert_eq!(first.rows, vec![EventAppendDisposition::Inserted]);
let rows = observations_for(&store, event.id).await;
assert_eq!(rows.len(), 3);
assert_eq!(rows[0].referent_kind, ReferentKind::Entity);
assert_eq!(rows[1].referent_kind, ReferentKind::Entity);
assert_eq!(rows[2].referent_kind, ReferentKind::Edge);
let retry = store.append_events_idempotent(vec![event]).await.unwrap();
assert_eq!(
retry.rows,
vec![EventAppendDisposition::AlreadyPresentIdentical]
);
}
#[tokio::test]
async fn link_to_event_projects_edge_without_event_endpoint() {
let store = setup_memory_store();
let source_id = Uuid::new_v4();
let target_event_id = Uuid::new_v4();
let edge_id = Uuid::new_v4();
let event = Event::new(
"default",
"link",
EventKind::LinkCreated,
SubstrateKind::Entity,
"agent:test",
)
.with_target(edge_id)
.with_payload(json!({
"id": edge_id,
"source_id": source_id,
"target_id": target_event_id,
"source_kind": "note",
"target_kind": "event",
"relation": "annotates"
}));
let event_id = event.id;
store.append_event(event).await.unwrap();
let rows = observations_for(&store, event_id).await;
assert_eq!(rows.len(), 2);
assert_eq!(rows[0].entity_id, source_id);
assert_eq!(rows[0].referent_kind, ReferentKind::Note);
assert_eq!(rows[0].position, 0);
assert_eq!(rows[1].entity_id, edge_id);
assert_eq!(rows[1].referent_kind, ReferentKind::Edge);
assert_eq!(rows[1].position, 2);
assert!(rows.iter().all(|row| row.entity_id != target_event_id));
}
#[tokio::test]
async fn query_events_filters_by_selected() {
let store = setup_memory_store();
let entity_id = Uuid::new_v4();
let mut event = make_event("default");
event.kind = EventKind::SearchExecuted;
event.payload = json!({
"result_kind": "entity",
"candidates": [],
"selected": [entity_id.to_string()]
});
store.append_event(event).await.unwrap();
store.append_event(make_event("default")).await.unwrap();
let filter = EventFilter {
selected: vec![entity_id],
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
}
#[tokio::test]
async fn query_events_filters_by_payload_proposal_id() {
let store = setup_memory_store();
let proposal_id = Uuid::new_v4();
let mut event = make_event("default");
event.kind = EventKind::ProposalCreated;
event.payload = json!({ "proposal_id": proposal_id.to_string() });
store.append_event(event).await.unwrap();
store.append_event(make_event("default")).await.unwrap();
let filter = EventFilter {
payload_proposal_id: Some(proposal_id),
..EventFilter::default()
};
let page = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.items.len(), 1);
}
#[tokio::test]
async fn query_events_observed_filter_missing_projection_returns_clean_error() {
let config = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = Arc::new(ConnectionPool::new(config).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(
"CREATE TABLE IF NOT EXISTS events (\
id TEXT PRIMARY KEY, namespace TEXT NOT NULL, verb TEXT NOT NULL,\
substrate TEXT NOT NULL, actor TEXT NOT NULL, kind TEXT NOT NULL DEFAULT 'audit',\
outcome TEXT NOT NULL, payload TEXT NOT NULL DEFAULT '{}',\
payload_schema_version INTEGER NOT NULL DEFAULT 1,\
duration_us INTEGER NOT NULL DEFAULT 0, created_at INTEGER NOT NULL\
);"
).unwrap();
}
let store = SqlEventStore::new_scoped(pool, false, "default");
let filter = EventFilter {
observed: vec![Uuid::new_v4()],
..EventFilter::default()
};
let result = store
.query_events(
filter,
PageRequest {
limit: 10,
offset: 0,
},
)
.await;
assert!(result.is_err());
let err_msg = result.unwrap_err().to_string();
assert!(
err_msg.contains("event_observations") && err_msg.contains("run migrations"),
"error should mention event_observations and run migrations, got: {err_msg}"
);
}
#[tokio::test]
async fn read_event_with_valid_payload_schema_version() {
let store = setup_memory_store();
let event = make_event("default");
let id = event.id;
store.append_event(event).await.unwrap();
let fetched = store.get_event(id).await.unwrap().unwrap();
assert_eq!(
fetched.payload_schema_version, 1,
"default payload_schema_version must round-trip as 1"
);
}
#[tokio::test]
async fn read_event_rejects_negative_payload_schema_version() {
use crate::pool::PoolConfig;
let config = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = Arc::new(ConnectionPool::new(config).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(EVENTS_DDL).unwrap();
}
let store = SqlEventStore::new_scoped(Arc::clone(&pool), false, "default");
let id = uuid::Uuid::new_v4();
{
let writer = pool.writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO events \
(id, namespace, verb, substrate, actor, kind, outcome, payload, \
payload_schema_version, duration_us, created_at) \
VALUES (?1,'default','test','entity','a','audit','success','{}', -1, 0, 0)",
rusqlite::params![id.to_string()],
)
.unwrap();
}
let result = store.get_event(id).await;
assert!(
result.is_err(),
"negative payload_schema_version must be rejected as a StorageError"
);
}
#[test]
fn observation_position_u32_max_plus_one_is_rejected() {
let overflow_position: usize = u32::MAX as usize + 1;
let result = u32::try_from(overflow_position);
assert!(
result.is_err(),
"u32::MAX + 1 ({overflow_position}) must not fit in u32"
);
}
#[test]
fn observation_position_u32_max_is_accepted() {
let max_position: usize = u32::MAX as usize;
let result = u32::try_from(max_position);
assert!(
result.is_ok(),
"u32::MAX ({max_position}) must be a valid position"
);
assert_eq!(result.unwrap(), u32::MAX);
}
#[tokio::test]
async fn read_event_rejects_payload_schema_version_u32_max_plus_one() {
let config = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = Arc::new(ConnectionPool::new(config).unwrap());
{
let writer = pool.writer().unwrap();
writer.conn().execute_batch(EVENTS_DDL).unwrap();
}
let store = SqlEventStore::new_scoped(Arc::clone(&pool), false, "default");
let id = uuid::Uuid::new_v4();
let overflow_version: i64 = i64::from(u32::MAX) + 1;
{
let writer = pool.writer().unwrap();
writer
.conn()
.execute(
"INSERT INTO events \
(id, namespace, verb, substrate, actor, kind, outcome, payload, \
payload_schema_version, duration_us, created_at) \
VALUES (?1,'default','test','entity','a','audit','success','{}', ?2, 0, 0)",
rusqlite::params![id.to_string(), overflow_version],
)
.unwrap();
}
let result = store.get_event(id).await;
assert!(
result.is_err(),
"payload_schema_version = u32::MAX + 1 ({overflow_version}) must be rejected"
);
}
#[tokio::test]
async fn page_offset_over_i64max_rejected() {
let store = setup_memory_store();
store.append_event(make_event("default")).await.unwrap();
let result = store
.query_events(
EventFilter::default(),
PageRequest {
offset: (i64::MAX as u64) + 1,
limit: 10,
},
)
.await;
assert!(
matches!(result, Err(StorageError::InvalidInput { .. })),
"expected InvalidInput, got {result:?}"
);
}
#[tokio::test]
async fn append_events_routes_through_writer_task_when_flag_enabled() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("write_queue_events.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(EVENTS_DDL).unwrap();
}
let store = SqlEventStore::new_scoped(Arc::clone(&pool), true, "default");
let e1 = make_event("default");
let e2 = make_event("default");
let id1 = e1.id;
let id2 = e2.id;
let summary = store.append_events(vec![e1, e2]).await.unwrap();
assert_eq!(summary.attempted, 2);
assert_eq!(summary.affected, 2);
assert_eq!(summary.failed, 0);
assert!(store.get_event(id1).await.unwrap().is_some());
assert!(store.get_event(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"
);
}
fn malformed_recall_event(namespace: &str) -> Event {
Event::new(
namespace,
"recall",
EventKind::RecallExecuted,
SubstrateKind::Note,
"agent:test",
)
.with_payload(json!({ "candidates": "not-an-array" }))
}
#[tokio::test]
async fn preflight_rejects_malformed_observation_before_enqueue() {
let store = setup_memory_store();
let bad = malformed_recall_event("default");
let bad_id = bad.id;
let result = store.preflight_event(&bad);
assert!(
result.is_err(),
"preflight must reject a payload the observation decoder cannot build statements for"
);
assert!(store.get_event(bad_id).await.unwrap().is_none());
let good = make_event("default");
assert!(store.preflight_event(&good).is_ok());
}
#[tokio::test]
async fn idempotent_batch_uses_one_writer_acquisition() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("idempotent_acquisition.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(EVENTS_DDL).unwrap();
}
let store = SqlEventStore::new_scoped(Arc::clone(&pool), true, "default");
let before = pool.writer_acquisition_snapshot().writer_task_acquisitions;
let events: Vec<Event> = (0..3).map(|_| make_event("default")).collect();
let ids: Vec<Uuid> = events.iter().map(|e| e.id).collect();
let result = store.append_events_idempotent(events).await.unwrap();
assert_eq!(result.rows.len(), 3);
assert!(result
.rows
.iter()
.all(|d| *d == EventAppendDisposition::Inserted));
for id in ids {
assert!(store.get_event(id).await.unwrap().is_some());
}
let after = pool.writer_acquisition_snapshot().writer_task_acquisitions;
assert_eq!(
after - before,
1,
"a 3-row idempotent batch must acquire the writer exactly once, not once per row"
);
}
#[tokio::test]
async fn idempotent_retry_accepts_only_full_equality() {
let store = setup_memory_store();
let event = make_event("default");
let id = event.id;
let first = store
.append_events_idempotent(vec![event.clone()])
.await
.unwrap();
assert_eq!(first.rows, vec![EventAppendDisposition::Inserted]);
let retry = store
.append_events_idempotent(vec![event.clone()])
.await
.unwrap();
assert_eq!(
retry.rows,
vec![EventAppendDisposition::AlreadyPresentIdentical]
);
assert_eq!(store.count_events(EventFilter::default()).await.unwrap(), 1);
let mut mismatched = event.clone();
mismatched.verb = "different-verb".to_string();
let conflict = store
.append_events_idempotent(vec![mismatched])
.await
.unwrap();
assert_eq!(
conflict.rows,
vec![EventAppendDisposition::IdentityConflict]
);
let stored = store.get_event(id).await.unwrap().unwrap();
assert_eq!(stored.verb, "search");
}
#[tokio::test]
async fn identity_collision_is_per_row_and_preserves_unrelated_rows() {
let store = setup_memory_store();
let original = make_event("default");
let original_id = original.id;
store
.append_events_idempotent(vec![original.clone()])
.await
.unwrap();
let mut conflicting = original.clone();
conflicting.verb = "changed".to_string();
let fresh = make_event("default");
let fresh_id = fresh.id;
let result = store
.append_events_idempotent(vec![conflicting, fresh])
.await
.unwrap();
assert_eq!(
result.rows,
vec![
EventAppendDisposition::IdentityConflict,
EventAppendDisposition::Inserted,
],
"one row's identity conflict must not affect an unrelated row in the same batch"
);
assert!(store.get_event(fresh_id).await.unwrap().is_some());
let stored_original = store.get_event(original_id).await.unwrap().unwrap();
assert_eq!(stored_original.verb, "search");
}
#[tokio::test]
async fn mixed_batch_store_failure_rolls_back_insertions() {
let store = setup_memory_store();
let good = make_event("default");
let good_id = good.id;
let bad = malformed_recall_event("default");
let err = store
.append_events_idempotent(vec![good, bad])
.await
.unwrap_err();
let _ = err;
assert!(
store.get_event(good_id).await.unwrap().is_none(),
"a genuine store failure later in the batch must roll back earlier insertions \
in the same transaction, unlike a per-row identity conflict"
);
assert_eq!(store.count_events(EventFilter::default()).await.unwrap(), 0);
}
#[tokio::test]
async fn reply_loss_after_commit_retries_exactly_once() {
let store = setup_memory_store();
let event = make_event("default");
let id = event.id;
let first = store
.append_events_idempotent(vec![event.clone()])
.await
.unwrap();
assert_eq!(first.rows, vec![EventAppendDisposition::Inserted]);
let retry = store.append_events_idempotent(vec![event]).await.unwrap();
assert_eq!(
retry.rows,
vec![EventAppendDisposition::AlreadyPresentIdentical]
);
assert_eq!(
store.count_events(EventFilter::default()).await.unwrap(),
1,
"an ambiguous-ack retry must land exactly once, never a duplicate row"
);
assert!(store.get_event(id).await.unwrap().is_some());
}
#[tokio::test]
async fn query_events_declines_to_compute_a_total() {
let store = setup_memory_store();
for _ in 0..3 {
store.append_event(make_event("default")).await.unwrap();
}
let page = store
.query_events(
EventFilter::default(),
PageRequest {
limit: 2,
offset: 0,
},
)
.await
.unwrap();
assert_eq!(page.total, None);
assert_eq!(page.items.len(), 2);
let counted = store.count_events(EventFilter::default()).await.unwrap();
assert_eq!(counted, 3);
}
#[tokio::test]
async fn refusal_target_filter_keeps_query_count_namespace_and_observations_separate() {
let pool = Arc::new(
ConnectionPool::new(PoolConfig {
path: None,
write_queue_enabled: Some(false),
write_routing_strict: false,
..PoolConfig::default()
})
.unwrap(),
);
pool.writer()
.unwrap()
.conn()
.execute_batch(EVENTS_DDL)
.unwrap();
let store = SqlEventStore::new_scoped(Arc::clone(&pool), false, "default");
let other = SqlEventStore::new_scoped(pool, false, "other");
let subject = Uuid::new_v4();
let refusal = |namespace: &str, target| {
Event::new(
namespace,
"knowledge.upsert_atoms",
EventKind::Refusal,
SubstrateKind::Event,
"actor:refusal-test",
)
.with_target(target)
.with_outcome(EventOutcome::Denied)
.with_payload(json!({"subject_kind": "knowledge_atom"}))
};
let first = refusal("default", subject);
let second = refusal("default", subject);
store
.append_events(vec![
first.clone(),
second.clone(),
refusal("default", Uuid::new_v4()),
Event::new(
"default",
"knowledge.upsert_atoms",
EventKind::Refusal,
SubstrateKind::Event,
"actor:refusal-test",
),
Event::new(
"default",
"test",
EventKind::Audit,
SubstrateKind::Event,
"actor:refusal-test",
)
.with_target(subject),
])
.await
.unwrap();
other.append_event(refusal("other", subject)).await.unwrap();
let filter = EventFilter {
target_id: Some(subject),
kinds: vec![EventKind::Refusal],
..EventFilter::default()
};
assert_eq!(store.count_events(filter.clone()).await.unwrap(), 2);
let mut found = Vec::new();
for offset in 0..2 {
let page = store
.query_events(filter.clone(), PageRequest { limit: 1, offset })
.await
.unwrap();
assert_eq!(page.items.len(), 1);
assert_eq!(page.items[0].target_id, Some(subject));
assert_eq!(page.items[0].namespace, "default");
found.push(page.items[0].id);
}
found.sort();
let mut expected = vec![first.id, second.id];
expected.sort();
assert_eq!(found, expected);
assert_eq!(other.count_events(filter.clone()).await.unwrap(), 1);
assert_eq!(
store
.count_events(EventFilter {
target_id: Some(subject),
..EventFilter::default()
})
.await
.unwrap(),
3
);
let observations = EventFilter {
observed: vec![subject],
..filter
};
assert_eq!(store.count_events(observations.clone()).await.unwrap(), 0);
assert!(
store
.query_events(
observations,
PageRequest {
limit: 10,
offset: 0
}
)
.await
.unwrap()
.items
.is_empty(),
"knowledge subjects must not become graph observations"
);
}