diff --git a/crates/khive-db/src/stores/event.rs b/crates/khive-db/src/stores/event.rs index f1a85337b..4e36990d4 100644 --- a/crates/khive-db/src/stores/event.rs +++ b/crates/khive-db/src/stores/event.rs @@ -438,7 +438,8 @@ pub async fn append_event_on_writer( fn decode_event_observations(event: &Event) -> Result, rusqlite::Error> { match event.kind { EventKind::RerankExecuted => decode_rerank_observations(event), - EventKind::RecallExecuted | EventKind::SearchExecuted => decode_recall_observations(event), + EventKind::RecallExecuted => decode_recall_observations(event), + EventKind::SearchExecuted => decode_search_observations(event), EventKind::LinkCreated => decode_link_observations(event), EventKind::EntityCreated | EventKind::EntityUpdated @@ -552,7 +553,10 @@ fn payload_uuid(event: &Event, field: &'static str) -> Result, rusq .map_err(|e| invalid_payload(event.kind, field, e)) } -fn decode_candidate_observations(event: &Event) -> Result, rusqlite::Error> { +fn decode_candidate_observations( + event: &Event, + referent_kind: ReferentKind, +) -> Result, rusqlite::Error> { let mut rows = Vec::new(); for (position, entity_id) in payload_uuid_array(event, "candidates")? @@ -569,7 +573,7 @@ fn decode_candidate_observations(event: &Event) -> Result, rows.push(EventObservation { event_id: event.id, entity_id, - referent_kind: ReferentKind::Note, + referent_kind, role: ObservationRole::Candidate, position: position_u32, }); @@ -581,6 +585,7 @@ fn decode_candidate_observations(event: &Event) -> Result, fn push_selected_observations( event: &Event, selected: Vec, + referent_kind: ReferentKind, rows: &mut Vec, ) -> Result<(), rusqlite::Error> { for (position, entity_id) in selected.into_iter().enumerate() { @@ -594,7 +599,7 @@ fn push_selected_observations( rows.push(EventObservation { event_id: event.id, entity_id, - referent_kind: ReferentKind::Note, + referent_kind, role: ObservationRole::Selected, position: position_u32, }); @@ -602,13 +607,43 @@ fn push_selected_observations( Ok(()) } -/// `RecallExecuted`/`SearchExecuted` payloads carry a flat `selected: Vec` -/// field (ADR-041 §"Projection rules"). These payloads are untyped JSON; the -/// ADR-041 projection contract makes `selected` the only field consulted here. +/// `RecallExecuted` payloads carry flat candidate and selected note UUID lists. fn decode_recall_observations(event: &Event) -> Result, rusqlite::Error> { - let mut rows = decode_candidate_observations(event)?; + let mut rows = decode_candidate_observations(event, ReferentKind::Note)?; + let selected = payload_uuid_array_opt(event, "selected")?.unwrap_or_default(); + push_selected_observations(event, selected, ReferentKind::Note, &mut rows)?; + Ok(rows) +} + +/// `SearchExecuted.result_kind` identifies which substrate owns every UUID in +/// the candidate and selected lists. Rejecting missing or unknown values keeps +/// the append-only projection from persisting an untyped reference. +fn decode_search_observations(event: &Event) -> Result, rusqlite::Error> { + let referent_kind = match event + .payload + .get("result_kind") + .and_then(|value| value.as_str()) + { + Some("entity") => ReferentKind::Entity, + Some("note") => ReferentKind::Note, + Some(_) => { + return Err(invalid_payload( + event.kind, + "result_kind", + "expected \"entity\" or \"note\"", + )); + } + None => { + return Err(invalid_payload( + event.kind, + "result_kind", + "expected string \"entity\" or \"note\"", + )); + } + }; + let mut rows = decode_candidate_observations(event, referent_kind)?; let selected = payload_uuid_array_opt(event, "selected")?.unwrap_or_default(); - push_selected_observations(event, selected, &mut rows)?; + push_selected_observations(event, selected, referent_kind, &mut rows)?; Ok(rows) } @@ -622,11 +657,11 @@ fn decode_recall_observations(event: &Event) -> Result, ru /// absent (`None`) — a present-but-malformed `final_scores` errors /// immediately instead of masking the problem by falling through. fn decode_rerank_observations(event: &Event) -> Result, rusqlite::Error> { - let mut rows = decode_candidate_observations(event)?; + let mut rows = decode_candidate_observations(event, ReferentKind::Note)?; let selected = payload_final_scores_uuid_array_opt(event, "final_scores")? .or(payload_reranked_uuid_array_opt(event, "reranked")?) .unwrap_or_default(); - push_selected_observations(event, selected, &mut rows)?; + push_selected_observations(event, selected, ReferentKind::Note, &mut rows)?; Ok(rows) } diff --git a/crates/khive-db/src/stores/event_tests.rs b/crates/khive-db/src/stores/event_tests.rs index 02bc5f833..ede18555f 100644 --- a/crates/khive-db/src/stores/event_tests.rs +++ b/crates/khive-db/src/stores/event_tests.rs @@ -25,6 +25,7 @@ fn make_event(namespace: &str) -> Event { SubstrateKind::Note, "agent:test", ) + .with_payload(json!({ "result_kind": "note" })) } #[tokio::test] @@ -146,6 +147,7 @@ async fn append_event_writes_observations_atomically() { 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" @@ -187,6 +189,43 @@ async fn append_event_writes_observations_atomically() { 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_rejects_absent_result_kind() { + let store = setup_memory_store(); + let mut event = make_event("default"); + event.payload = json!({ + "candidates": [Uuid::new_v4().to_string()], + "selected": [] + }); + let event_id = event.id; + + let result = store.append_event(event).await; + assert!(result.is_err(), "absent result_kind must be rejected"); + assert!( + store.get_event(event_id).await.unwrap().is_none(), + "invalid event and projection must roll back atomically" + ); +} + async fn selected_uuids_for(store: &SqlEventStore, event_id: Uuid) -> Vec { let pool = Arc::clone(&store.pool); let event_id_str = event_id.to_string(); @@ -647,6 +686,7 @@ async fn query_events_filters_by_observed() { let mut event = make_event("default"); event.kind = EventKind::SearchExecuted; event.payload = json!({ + "result_kind": "entity", "candidates": [entity_id.to_string()], "selected": [] }); @@ -677,6 +717,7 @@ async fn query_events_filters_by_selected() { let mut event = make_event("default"); event.kind = EventKind::SearchExecuted; event.payload = json!({ + "result_kind": "entity", "candidates": [], "selected": [entity_id.to_string()] }); diff --git a/crates/khive-pack-brain/src/tests.rs b/crates/khive-pack-brain/src/tests.rs index 42fd3e523..819bd9b3a 100644 --- a/crates/khive-pack-brain/src/tests.rs +++ b/crates/khive-pack-brain/src/tests.rs @@ -5896,6 +5896,9 @@ mod event_counts_tests { ); event.created_at = created_at; event.payload = payload; + if kind == EventKind::SearchExecuted && event.payload.get("result_kind").is_none() { + event.payload["result_kind"] = json!("note"); + } rt.events(token) .expect("event store") .append_event(event) diff --git a/crates/khive-pack-kg/src/handlers/search.rs b/crates/khive-pack-kg/src/handlers/search.rs index 8dda7bbe2..fa47716ef 100644 --- a/crates/khive-pack-kg/src/handlers/search.rs +++ b/crates/khive-pack-kg/src/handlers/search.rs @@ -6,10 +6,12 @@ use std::collections::HashMap; /// See `docs/api/scan-cliff.md`. const FILTERED_SCAN_CAP: u32 = 500; -use serde_json::Value; +use std::time::Instant; + +use serde_json::{json, Value}; use uuid::Uuid; -use khive_runtime::{NamespaceToken, RuntimeError, VerbRegistry}; +use khive_runtime::{KhiveRuntime, NamespaceToken, RuntimeError, VerbRegistry}; use khive_storage::types::PageRequest; use khive_storage::EntityFilter; @@ -26,6 +28,7 @@ impl KgPack { params: Value, registry: &VerbRegistry, ) -> Result { + let search_start = Instant::now(); let p: SearchParams = deser(params)?; let limit = p.limit.unwrap_or(10).min(100); let spec = resolve_kind_spec(&p.kind, registry)?; @@ -141,6 +144,13 @@ impl KgPack { }) }) .collect(); + self.track_search_serve( + token, + &p.query, + "entity", + &result, + search_start.elapsed().as_micros() as i64, + ); to_json(&result) } KindSpec::Note { specific } => { @@ -242,6 +252,13 @@ impl KgPack { }) }) .collect(); + self.track_search_serve( + token, + &p.query, + "note", + &result, + search_start.elapsed().as_micros() as i64, + ); to_json(&result) } KindSpec::Edge => Err(RuntimeError::InvalidInput( @@ -255,4 +272,92 @@ impl KgPack { )), } } + + /// Fire-and-forget `search_executed` telemetry (ADR-103 event plane), + /// mirroring `memory.recall`'s `track_recall_serve` seam (#866): the + /// event append runs off the response path via `track_background_task` + /// so a slow or failing event store never affects a served search. + fn track_search_serve( + &self, + token: &NamespaceToken, + query_raw: &str, + result_kind: &'static str, + results: &[Value], + latency_us: i64, + ) { + let selected: Vec = results + .iter() + .filter_map(|r| r.get("id").and_then(Value::as_str).map(str::to_string)) + .collect(); + let result_count = selected.len(); + let query = query_raw.to_string(); + let actor = format!("{}:{}", token.actor().kind, token.actor().id); + let runtime = self.runtime.clone(); + let token = token.clone(); + + khive_runtime::track_background_task(async move { + emit_search_executed_event( + &runtime, + &token, + actor, + query, + result_kind, + selected, + result_count, + latency_us, + ) + .await; + }); + } +} + +/// Append best-effort search telemetry without affecting the search response. +#[allow(clippy::too_many_arguments)] +async fn emit_search_executed_event( + rt: &KhiveRuntime, + token: &NamespaceToken, + actor: String, + query: String, + result_kind: &'static str, + selected: Vec, + result_count: usize, + latency_us: i64, +) { + let store = match rt.events(token) { + Ok(store) => store, + Err(err) => { + tracing::warn!( + error = %err, + namespace = token.namespace().as_str(), + event_kind = "search_executed", + "search_executed event store acquisition failed; search result is unaffected" + ); + return; + } + }; + let payload = json!({ + "actor": actor, + "served_by_profile_id": Value::Null, + "query": query, + "result_kind": result_kind, + "result_count": result_count, + "candidates": selected, + "selected": selected, + "latency_us": latency_us, + }); + let event = khive_storage::Event::new( + token.namespace().as_str(), + "search", + khive_types::EventKind::SearchExecuted, + khive_types::SubstrateKind::Event, + actor, + ) + .with_payload(payload) + .with_duration_us(latency_us); + if let Err(err) = store.append_event(event).await { + tracing::warn!( + error = %err, + "search_executed event append failed; search result is unaffected" + ); + } } diff --git a/crates/khive-pack-kg/tests/integration.rs b/crates/khive-pack-kg/tests/integration.rs index ed6c2e012..acace9dfc 100644 --- a/crates/khive-pack-kg/tests/integration.rs +++ b/crates/khive-pack-kg/tests/integration.rs @@ -11,7 +11,7 @@ use khive_runtime::{ EntityCreateSpec, KhiveRuntime, Namespace, NamespaceToken, ParamDef, RuntimeError, VerbCategory, VerbRegistry, VerbRegistryBuilder, VerifiedActor, Visibility, }; -use khive_storage::{Edge, EdgeRelation, Note}; +use khive_storage::{Edge, EdgeRelation, Note, SqlStatement, SqlValue}; use khive_types::Pack; use serde_json::{json, Value}; @@ -11505,3 +11505,262 @@ async fn list_proposal_limit_over_cap_reports_effective_limit() { assert_eq!(response["effective_limit"], 500); assert_eq!(response["limit_clamped"], true); } + +// ── #806: `search_executed` event-plane emission ──────────────────────────── +// +// Mirrors memory.recall's `#866` `recall_executed` regression +// (khive-pack-memory/src/handlers/recall.rs's +// `recall_emits_exactly_one_recall_executed_event`): `search` must append +// exactly one `SearchExecuted` event per served search, off the response +// path via `track_background_task`, so poll briefly instead of assuming it +// has landed by the time `search` returns. + +async fn poll_search_executed_events( + store: &std::sync::Arc, +) -> Vec { + for _ in 0..100 { + let page = store + .query_events( + khive_storage::EventFilter { + kinds: vec![khive_types::EventKind::SearchExecuted], + ..Default::default() + }, + khive_storage::types::PageRequest { + limit: 50, + offset: 0, + }, + ) + .await + .expect("query_events"); + if !page.items.is_empty() { + return page.items; + } + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + } + Vec::new() +} + +async fn assert_search_projection( + rt: &KhiveRuntime, + store: &std::sync::Arc, + event: &khive_storage::Event, + hits: &[Value], + referent_kind: &str, +) { + let expected_ids: Vec = hits + .iter() + .map(|hit| { + hit["id"] + .as_str() + .expect("search hit id must be a UUID string") + .to_string() + }) + .collect(); + assert_eq!(event.payload["candidates"], json!(expected_ids)); + assert_eq!(event.payload["selected"], json!(expected_ids)); + + let mut reader = rt.sql().reader().await.expect("sql reader must open"); + let rows = reader + .query_all(SqlStatement { + sql: "SELECT entity_id, referent_kind, role, position \ + FROM event_observations WHERE event_id = ?1 \ + ORDER BY CASE role WHEN 'candidate' THEN 0 ELSE 1 END, position" + .into(), + params: vec![SqlValue::Text(event.id.to_string())], + label: Some("search_observation_projection".into()), + }) + .await + .expect("projection query must succeed"); + + let projected: Vec<(String, String, String, i64)> = rows + .iter() + .map(|row| { + let text = |column| match row.get(column) { + Some(SqlValue::Text(value)) => value.clone(), + other => panic!("{column} must be text, got {other:?}"), + }; + let position = match row.get("position") { + Some(SqlValue::Integer(value)) => *value, + other => panic!("position must be integer, got {other:?}"), + }; + ( + text("entity_id"), + text("referent_kind"), + text("role"), + position, + ) + }) + .collect(); + let expected_projection: Vec<(String, String, String, i64)> = ["candidate", "selected"] + .into_iter() + .flat_map(|role| { + expected_ids.iter().enumerate().map(move |(position, id)| { + ( + id.clone(), + referent_kind.to_string(), + role.to_string(), + position as i64, + ) + }) + }) + .collect(); + assert_eq!(projected, expected_projection); + + let first_id = uuid::Uuid::parse_str(&expected_ids[0]).expect("search hit id must parse"); + for (filter_name, filter) in [ + ( + "observed", + khive_storage::EventFilter { + kinds: vec![khive_types::EventKind::SearchExecuted], + observed: vec![first_id], + ..Default::default() + }, + ), + ( + "selected", + khive_storage::EventFilter { + kinds: vec![khive_types::EventKind::SearchExecuted], + selected: vec![first_id], + ..Default::default() + }, + ), + ] { + let page = store + .query_events( + filter, + khive_storage::types::PageRequest { + limit: 10, + offset: 0, + }, + ) + .await + .unwrap_or_else(|error| panic!("{filter_name} filter failed: {error}")); + assert_eq!( + page.items.iter().map(|item| item.id).collect::>(), + vec![event.id], + "{filter_name} must find the search event through its projection" + ); + } +} + +#[tokio::test] +async fn search_entity_emits_exactly_one_search_executed_event() { + let rt = KhiveRuntime::memory().expect("in-memory runtime must succeed"); + let ns = Namespace::local(); + let token = rt.authorize(ns.clone()).expect("authorize local"); + + rt.create_entity( + &token, + "concept", + None, + "806 search executed entity", + None, + None, + vec![], + ) + .await + .expect("create entity"); + + let mut builder = VerbRegistryBuilder::new(); + builder.register(KgPack::new(rt.clone())); + let registry = builder.build().expect("registry builds"); + + let result = registry + .dispatch( + "search", + json!({"kind": "entity", "query": "806 search executed entity"}), + ) + .await + .expect("search must succeed"); + let hits = result.as_array().expect("bare array result"); + assert!(!hits.is_empty(), "must find the seeded entity"); + + let store = rt.events(&token).expect("event store for local namespace"); + let search_events = poll_search_executed_events(&store).await; + + assert_eq!( + search_events.len(), + 1, + "exactly one search_executed event per served search, got: {search_events:?}" + ); + let event = &search_events[0]; + assert_eq!(event.kind, khive_types::EventKind::SearchExecuted); + assert_eq!(event.verb, "search"); + assert_eq!(event.payload["result_kind"], json!("entity")); + assert_eq!(event.payload["result_count"], json!(hits.len())); + assert_eq!(event.payload["query"], json!("806 search executed entity")); + assert_eq!( + event.payload["served_by_profile_id"], + Value::Null, + "kg search has no profile resolution — the field must stay present but null" + ); + assert!( + event.payload["latency_us"] + .as_i64() + .is_some_and(|us| us >= 0), + "latency_us must be a non-negative measured duration, got: {:?}", + event.payload["latency_us"] + ); + assert!( + event.payload["actor"] + .as_str() + .is_some_and(|s| !s.is_empty()), + "actor must be stamped, got: {:?}", + event.payload["actor"] + ); + let selected = event.payload["selected"] + .as_array() + .expect("selected must be a UUID array"); + assert_eq!( + selected.len(), + hits.len(), + "selected must carry the full served result-ID list" + ); + assert_search_projection(&rt, &store, event, hits, "entity").await; +} + +#[tokio::test] +async fn search_note_emits_exactly_one_search_executed_event_with_note_result_kind() { + let rt = KhiveRuntime::memory().expect("in-memory runtime must succeed"); + let ns = Namespace::local(); + let token = rt.authorize(ns.clone()).expect("authorize local"); + + rt.create_note( + &token, + "observation", + None, + "806 search executed note unique_marker_9931", + None, + None, + vec![], + ) + .await + .expect("create note"); + + let mut builder = VerbRegistryBuilder::new(); + builder.register(KgPack::new(rt.clone())); + let registry = builder.build().expect("registry builds"); + + let result = registry + .dispatch( + "search", + json!({"kind": "note", "query": "unique_marker_9931"}), + ) + .await + .expect("search must succeed"); + let hits = result.as_array().expect("bare array result"); + assert!(!hits.is_empty(), "must find the seeded note"); + + let store = rt.events(&token).expect("event store for local namespace"); + let search_events = poll_search_executed_events(&store).await; + + assert_eq!( + search_events.len(), + 1, + "exactly one search_executed event per served search, got: {search_events:?}" + ); + let event = &search_events[0]; + assert_eq!(event.payload["result_kind"], json!("note")); + assert_eq!(event.payload["result_count"], json!(hits.len())); + assert_search_projection(&rt, &store, event, hits, "note").await; +} diff --git a/crates/khive-runtime/tests/integration.rs b/crates/khive-runtime/tests/integration.rs index 58bee874a..c69f94820 100644 --- a/crates/khive-runtime/tests/integration.rs +++ b/crates/khive-runtime/tests/integration.rs @@ -1001,11 +1001,13 @@ async fn synthetic_edge_observed_as_selected_returns_memory_note() { // Step 2: create an event of kind SearchExecuted with a payload that // includes `selected: [memory_id]`. The storage layer's `append_event` - // implementation calls `decode_recall_observations`, which reads - // `payload["selected"]` and inserts a row into `event_observations` with - // role="selected" and entity_id=memory_id. (`selected` is part of - // `SearchExecuted`/`RecallExecuted`'s ADR-041 projection contract; - // their payloads are untyped JSON — + // routes `SearchExecuted` to `decode_search_observations`, which requires + // the `result_kind` discriminator (`entity` | `note`) to identify which + // substrate owns every UUID in `candidates` and `selected`, then inserts a + // row into `event_observations` with role="selected" and + // entity_id=memory_id. (`selected` is part of `SearchExecuted`/ + // `RecallExecuted`'s ADR-041 projection contract; `RecallExecuted` still + // decodes via `decode_recall_observations` from flat note UUID lists, // unlike `RerankExecuted`, which projects `selected` rows from // `final_scores`/`reranked` instead, since its typed payload has no // `selected` field.) @@ -1018,6 +1020,7 @@ async fn synthetic_edge_observed_as_selected_returns_memory_note() { "agent:test", ); event.payload = serde_json::json!({ + "result_kind": "note", "candidates": [], "selected": [memory_id.to_string()] });