diff --git a/docs/README.md b/docs/README.md index e824f54..506917e 100644 --- a/docs/README.md +++ b/docs/README.md @@ -23,6 +23,8 @@ some are older research reports kept for rationale. architecture decisions. main follow-up: `klbr-u0q`. - [`memory-implementation-status.md`](./memory-implementation-status.md) - current implementation map, verification commands, and known open work. +- [`memory-storage-ownership.md`](./memory-storage-ownership.md) - canonical + source tables versus rebuildable retrieval projections and migration rules. - [`folgezettel-memory.md`](./folgezettel-memory.md) - current agent-facing memory tool ontology: sourced zettel tree, titles as api handles, bodies as attractors, and trail-aware retrieval. diff --git a/docs/memory-implementation-status.md b/docs/memory-implementation-status.md index 95e3af4..bb0fdab 100644 --- a/docs/memory-implementation-status.md +++ b/docs/memory-implementation-status.md @@ -28,9 +28,21 @@ going next. - unicode-token fts rows in `promptable_text_fts` - trigram fts rows in `promptable_text_trigram` for no-space, substring, identifier, and exact-ish candidate generation -- `MemoryGarden` can write markdown note files and sync them back into sqlite. -- legacy `memory_edges` now mirror into canonical `edges`, `sync_reference_indexes` backfills old rows, and memory-id provenance APIs read canonical ref edges. -- memory status changes now project onto canonical refs and rebuild fts visibility. +- `MemoryGarden` reconciles files with sqlite using persisted file/store hashes. + create, update, rename, delete, and conflict outcomes are explicit; deletion + tombstones an unchanged store note and two-sided edits overwrite neither side. +- session ingestion is resumable and idempotent through `session_ingests`. + turn, note, and memory projection writes use local transactions, and + provenance failures are no longer discarded. +- [`memory-storage-ownership.md`](./memory-storage-ownership.md) defines source + content tables and canonical refs/edges versus rebuildable projections. + episode cards are note-owned instead of dual-written to legacy L1 memory. +- legacy `memory_edges` mirror into canonical `edges`; explicit projection + repair backfills old rows, while normal retrieval performs no global repair. +- memory status changes update only the affected canonical refs and FTS rows. +- async pipeline storage jobs use a bounded blocking gate. retrieval opens + independent query-only WAL connections while writes retain one ordered + transactional connection. - note, episode, and chunk refs can resolve benchmark session ids from ref metadata/frontmatter. - dense retrieval now searches canonical refs through `embedding_items`, with bounded local embedding text. @@ -85,6 +97,9 @@ going next. order, and update-resolution questions. packet planning widens candidate and packet limits, traces fact-group coverage/gap state, and favors distinct session/ref fact groups instead of letting one session monopolize context. + collect mode now consumes a bounded second pass, follows query-aligned target + gaps, merges only new refs, and reports explicit sufficiency/degradation stop + reasons. - the structured query planner parser is tolerant of common model label shapes (`op`, `operation`, `query_op`, `count`, `latest`, `average`, etc.) while preserving the typed `QueryOp` interface. longmemeval category hints can @@ -101,6 +116,9 @@ going next. `duration`, `date`, `time`, `ordinal`) so aggregate synthesis does not have to trust the first number in a body. rendered `` include `` and `` blocks before bodies. + rows also carry query-derived target identity, extraction/polarity provenance, + and separate `observed_at`/`valid_at` axes. reducers require a consistent + target and defer ambiguous preference rows to the grounded reader. - `MemoryPipeline::answer` has an initial planned synthesis path for aggregate count/sum/avg, chronological ordering, previous/latest update resolution, and preference recommendations with personal-support guards. @@ -139,14 +157,13 @@ rtk cargo test -p klbr-bench rtk cargo run -p klbr-bench -- continuous-loop /tmp/klbr-continuous-loop ``` -current result: +current local result on aarch64-darwin through `nix develop`: ```text -cargo check -p klbr-core: 0 errors, 1 pre-existing warning -cargo check -p klbr-bench: 0 errors, 1 pre-existing warning +cargo check --workspace: passed focused folgezettel tests: mk surface 2 passed; trail-support packet test 1 passed -klbr-core tests: 164 passed, 1 ignored -klbr-bench tests: 19 passed +klbr-core tests: 183 passed; 1 local-embedder integration test skipped +klbr-bench tests: 21 passed continuous-loop smoke: passed with the pre-existing agent.rs dead-code warning session-first/fts-only trace smoke: evaluated 1, CandidateSessionRecallAny@5 1.0000 ``` diff --git a/docs/memory-storage-ownership.md b/docs/memory-storage-ownership.md new file mode 100644 index 0000000..fe69686 --- /dev/null +++ b/docs/memory-storage-ownership.md @@ -0,0 +1,81 @@ +# memory storage ownership + +status: current architecture decision +decided: 2026-07-10 +tracked by: `klbr-1jn.3` + +## decision + +klbr has one canonical identity and lifecycle substrate: `refs`, aliases, and +`edges`. content remains canonical in the source table for its entity family. +search indexes and rendered retrieval material are projections and must never +be repaired as a side effect of reading. + +| concern | canonical owner | derived / compatibility structures | +|---|---|---| +| immutable conversation source | `turns`, `turn_chunks` | `promptable_text`, fts, embeddings | +| markdown, zettel, profile, procedural, and episode-card content | `markdown_notes`, `markdown_note_chunks` | `promptable_text`, fts, embeddings | +| legacy L1 memory and versions | `memories`, `memory_versions` | `vec_memories`, `promptable_text`, fts | +| stable identity, aliases, lifecycle | `refs`, `ref_aliases` | none | +| retrieval lane/kind annotations | `ref_metadata` | lane fields copied into fts rows | +| provenance and topology | `edges` | `memory_edges` is migration compatibility only | +| dense vectors for canonical refs | `embedding_items` | `vec_memories` is legacy-memory compatibility only | +| lexical search | none | `promptable_text`, `promptable_text_fts`, `promptable_text_trigram` | + +episode event cards are markdown-note entities. production ingest must not also +create a semantically identical `memories` row. dense retrieval embeds the +episode note ref through `embedding_items` like every other canonical ref. + +## invariants + +- a logical artifact has one source content row and one canonical ref. +- mutations update the source row, ref lifecycle, and affected projections in + one local transaction. +- readers never call global backfill or index rebuild functions. +- projection rows may be deleted and rebuilt from canonical owners. +- lifecycle filtering comes from `refs.status`; projections cannot revive an + inactive ref. +- `edges` is the only graph written by new code. legacy graph rows may be read + or migrated, but are not parallel authority. + +## migration + +1. stop dual-writing episode cards to `memories`. +2. incrementally maintain fts and embedding projections on mutation. +3. add an explicit integrity audit and repair command for old databases. +4. backfill canonical episode-note embeddings, then retire episode-shaped L1 + duplicates after verifying their source refs and edges were preserved. +5. remove legacy graph/vector compatibility only after deployed databases have + crossed the migration and rollback window. + +rollback keeps the source tables and refs untouched. projections can be rebuilt; +compatibility tables remain readable until their later removal issue lands. + +## garden reconciliation + +garden sync records the last file hash and store hash for every path/ref pair. +one-sided file changes import normally. renames preserve the note ref. deletion +tombstones the note only when the store still matches the last sync. if both +sides changed, sync reports a conflict and overwrites neither side. + +automatic store-to-file export is intentionally unsupported for now: agent/db +changes are protected as conflicts until a person exports or reconciles them. +this is preferable to silently replacing a manually edited note or reviving a +deleted one. + +## sqlite execution + +the canonical store keeps one ordered writer connection for transactions. +async memory-pipeline operations acquire a bounded permit and run sqlite work +on Tokio's blocking pool. retrieval jobs open independent query-only WAL +connections, so concurrent reads do not wait on or share the mutable writer. +dropping a caller cancels interest in the result, not an already-running sqlite +call; the bounded gate supplies backpressure and prevents unbounded blocking +jobs. transaction commit remains the cancellation boundary for writes. + +## rejected alternatives + +keeping every table authoritative preserves the current dual-write and repair +failure modes. making `memories` own notes would erase the file-backed note +identity and folgezettel semantics. making `promptable_text` canonical would +turn a tokenizer-specific cache into durable user data. diff --git a/flake.nix b/flake.nix index 8e7406d..125ee6d 100644 --- a/flake.nix +++ b/flake.nix @@ -6,7 +6,10 @@ outputs = inp: inp.parts.lib.mkFlake { inputs = inp; } { - systems = [ "x86_64-linux" ]; + systems = [ + "x86_64-linux" + "aarch64-darwin" + ]; imports = [ inp.nci.flakeModule ]; perSystem = { @@ -23,10 +26,10 @@ ++ (with pkgs; [ cargo-outdated clang - wild bun typescript-language-server - ]); + ]) + ++ pkgs.lib.optionals pkgs.stdenv.isLinux [ pkgs.wild ]; }); }; }; diff --git a/klbr-core/src/evidence.rs b/klbr-core/src/evidence.rs index e27d58a..c45441a 100644 --- a/klbr-core/src/evidence.rs +++ b/klbr-core/src/evidence.rs @@ -2,7 +2,7 @@ use std::collections::{HashMap, HashSet}; use std::sync::OnceLock; use anyhow::Result; -use chrono::{SecondsFormat, TimeZone, Utc}; +use chrono::{NaiveDate, SecondsFormat, TimeZone, Utc}; use regex::Regex; use crate::memory::{to_base36, MemoryLane, MemoryStore, ResolvedRef}; @@ -257,7 +257,21 @@ pub struct EvidenceFactRow { pub entity_type: String, pub session_id: Option, pub timestamp: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub observed_at: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub valid_at: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub valid_until: Option, pub ordinal: usize, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub target_key: Option, + #[serde(default = "unspecified_polarity")] + pub polarity: String, + #[serde(default)] + pub extraction_confidence: f32, + #[serde(default = "default_extraction_method")] + pub extraction_method: String, pub value_text: Option, pub value_number: Option, #[serde(default)] @@ -265,11 +279,23 @@ pub struct EvidenceFactRow { pub text: String, } +fn unspecified_polarity() -> String { + "unspecified".to_string() +} + +fn default_extraction_method() -> String { + "legacy_unspecified".to_string() +} + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct EvidenceTimelineEvent { pub ref_id: String, pub session_id: Option, pub timestamp: i64, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub observed_at: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub valid_at: Option, pub ordinal: usize, pub text: String, } @@ -303,7 +329,7 @@ impl EvidencePlanner { pub fn build_packets_with_collect_mode( &self, - _query: &str, + query: &str, candidates: &[EvidenceAtom], limit: usize, packet_token_budget: usize, @@ -315,7 +341,7 @@ impl EvidencePlanner { for (index, candidate) in candidates.iter().enumerate() { let expanded = - self.expand_candidate_packet(candidate, index + 1, packet_token_budget)?; + self.expand_candidate_packet(query, candidate, index + 1, packet_token_budget)?; let packet_id = expanded.packet.packet_id.clone(); omissions.extend(expanded.omissions); if let Some(position) = packet_positions.get(&packet_id).copied() { @@ -343,6 +369,7 @@ impl EvidencePlanner { fn expand_candidate_packet( &self, + query: &str, candidate: &EvidenceAtom, rank: usize, packet_token_budget: usize, @@ -436,17 +463,21 @@ impl EvidencePlanner { ); let packet_id = format!("pkt_{}", stable_base36(&packet_seed)); let fact_rows = - self.fact_rows_for_packet(&refs, &bodies, candidate.session_id.as_deref())?; + self.fact_rows_for_packet(query, &refs, &bodies, candidate.session_id.as_deref())?; let timeline = fact_rows .iter() .filter_map(|row| { - row.timestamp.map(|timestamp| EvidenceTimelineEvent { - ref_id: row.ref_id.clone(), - session_id: row.session_id.clone(), - timestamp, - ordinal: row.ordinal, - text: row.text.clone(), - }) + row.valid_at + .or(row.observed_at) + .map(|timestamp| EvidenceTimelineEvent { + ref_id: row.ref_id.clone(), + session_id: row.session_id.clone(), + timestamp, + observed_at: row.observed_at, + valid_at: row.valid_at, + ordinal: row.ordinal, + text: row.text.clone(), + }) }) .collect(); let omissions = pending_omissions @@ -483,6 +514,7 @@ impl EvidencePlanner { fn fact_rows_for_packet( &self, + query: &str, refs: &[String], bodies: &[String], fallback_session_id: Option<&str>, @@ -503,6 +535,8 @@ impl EvidencePlanner { let timestamp = self.memory.ref_timestamp(ref_id)?; let values = numeric_value_mentions(body); let (value_text, value_number) = primary_numeric_value(&values); + let (target_key, extraction_confidence) = query_target_key(query, body); + let valid_at = explicit_valid_at(&values); let text = truncate_chars(body, 180); let row_seed = format!("{ref_id}:{ordinal}:{}", text); Ok(EvidenceFactRow { @@ -511,7 +545,14 @@ impl EvidencePlanner { entity_type, session_id, timestamp, + observed_at: timestamp, + valid_at, + valid_until: None, ordinal, + target_key, + polarity: "unspecified".to_string(), + extraction_confidence, + extraction_method: "query_token_overlap_v1".to_string(), value_text, value_number, values, @@ -542,6 +583,10 @@ pub(crate) struct RankedPacketOutcome { #[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)] pub struct EvidenceCoverageTrace { pub collect_mode: bool, + pub passes: usize, + pub stop_reason: String, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub followup_queries: Vec, pub available_fact_groups: usize, pub selected_fact_groups: usize, pub fact_row_count: usize, @@ -622,12 +667,18 @@ pub fn render_evidence_packet_with_refs( xml_escape(&fact.ref_id), xml_escape(&fact.entity_type), fact.ordinal, - xml_escape(fact.session_id.as_deref().unwrap_or("")) + xml_escape(fact.session_id.as_deref().unwrap_or("")), ); - if let Some(timestamp) = fact.timestamp { - attrs.push_str(&format!(" t=\"{}\"", timestamp)); - if let Some(date) = format_timestamp_utc(timestamp) { - attrs.push_str(&format!(" date=\"{}\"", xml_escape(&date))); + if let Some(target_key) = &fact.target_key { + attrs.push_str(&format!(" target=\"{}\"", xml_escape(target_key))); + } + if let Some(observed_at) = fact.observed_at { + attrs.push_str(&format!(" observed_at=\"{}\"", observed_at)); + } + if let Some(valid_at) = fact.valid_at { + attrs.push_str(&format!(" valid_at=\"{}\"", valid_at)); + if let Some(date) = format_timestamp_utc(valid_at) { + attrs.push_str(&format!(" valid_date=\"{}\"", xml_escape(&date))); } } if let Some(value) = &fact.value_text { @@ -675,12 +726,20 @@ pub fn render_evidence_packet_with_refs( .filter_map(|index| packet.timeline.get(index)) { let line = format!( - " {}\n", + " {}\n", xml_escape(&event.ref_id), event.timestamp, format_timestamp_utc(event.timestamp) .map(|date| format!(" date=\"{}\"", xml_escape(&date))) .unwrap_or_default(), + event + .observed_at + .map(|timestamp| format!(" observed_at=\"{timestamp}\"")) + .unwrap_or_default(), + event + .valid_at + .map(|timestamp| format!(" valid_at=\"{timestamp}\"")) + .unwrap_or_default(), event.ordinal, xml_escape(event.session_id.as_deref().unwrap_or("")), xml_escape(&truncate_chars(&event.text, 96)), @@ -787,6 +846,13 @@ fn evidence_coverage_trace( .sum(); EvidenceCoverageTrace { collect_mode, + passes: 1, + stop_reason: if collect_mode { + "initial_pass".to_string() + } else { + "single_pass".to_string() + }, + followup_queries: vec![], available_fact_groups: available.len(), selected_fact_groups: selected_count, fact_row_count, @@ -1177,6 +1243,42 @@ fn primary_numeric_value(values: &[EvidenceValueMention]) -> (Option, Op .unwrap_or((None, None)) } +fn query_target_key(query: &str, body: &str) -> (Option, f32) { + let query_terms = normalized_terms(query); + if query_terms.is_empty() { + return (None, 0.0); + } + let body_terms = normalized_terms(body).into_iter().collect::>(); + let matching = query_terms + .iter() + .filter(|term| body_terms.contains(*term)) + .collect::>(); + let confidence = matching.len() as f32 / query_terms.len() as f32; + (matching.last().map(|term| (*term).clone()), confidence) +} + +fn explicit_valid_at(values: &[EvidenceValueMention]) -> Option { + values + .iter() + .filter(|value| value.kind == EvidenceValueKind::Date) + .find_map(|value| { + ["%Y-%m-%d", "%Y/%m/%d", "%m/%d/%Y"] + .into_iter() + .find_map(|format| NaiveDate::parse_from_str(&value.text, format).ok()) + }) + .and_then(|date| date.and_hms_opt(0, 0, 0)) + .map(|datetime| datetime.and_utc().timestamp()) +} + +fn normalized_terms(text: &str) -> Vec { + text.split(|ch: char| !ch.is_alphanumeric() && ch != '_') + .map(str::trim) + .filter(|term| term.chars().count() >= 3) + .filter(|term| !term.chars().all(|ch| ch.is_numeric())) + .map(|term| term.to_lowercase()) + .collect() +} + fn numeric_value_mentions(text: &str) -> Vec { numeric_token_regex() .find_iter(text) @@ -1561,7 +1663,14 @@ mod tests { entity_type: "turn_chunk".to_string(), session_id: Some("s1".to_string()), timestamp: Some(42), + observed_at: Some(42), + valid_at: None, + valid_until: None, ordinal: 0, + target_key: Some("personal_best".to_string()), + polarity: "asserted".to_string(), + extraction_confidence: 1.0, + extraction_method: "test_fixture".to_string(), value_text: Some("27:45".to_string()), value_number: None, values: vec![EvidenceValueMention { @@ -1576,6 +1685,8 @@ mod tests { ref_id: "ref_fact".to_string(), session_id: Some("s1".to_string()), timestamp: 42, + observed_at: Some(42), + valid_at: None, ordinal: 0, text: "personal best was 27:45".to_string(), }]; @@ -1585,7 +1696,7 @@ mod tests { assert!(rendered.contains("")); assert!(rendered.contains("")); assert!(rendered.contains("ref=\"ref_fact\"")); - assert!(rendered.contains("date=\"1970-01-01T00:00:42Z\"")); + assert!(rendered.contains("observed_at=\"42\"")); assert!(rendered.contains("value=\"27:45\"")); assert!(rendered.contains("values=\"time:27:45\"")); assert!(rendered.contains(", +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct GardenSyncReport { + pub outcomes: Vec, +} + +impl GardenSyncReport { + pub fn has_conflicts(&self) -> bool { + self.outcomes + .iter() + .any(|outcome| outcome.action == "conflict") + } +} + #[derive(Debug, Clone)] pub struct MemoryGarden { root: PathBuf, @@ -49,12 +72,25 @@ impl MemoryGarden { Ok(path) } - pub fn sync_to_store(&self, memory: &MemoryStore) -> Result> { + pub fn sync_to_store(&self, memory: &MemoryStore) -> Result { let mut paths = Vec::new(); collect_markdown_files(&self.root, &mut paths)?; paths.sort(); - let mut records = Vec::new(); + let states = memory.garden_sync_states()?; + let by_path = states + .iter() + .cloned() + .map(|state| (state.path.clone(), state)) + .collect::>(); + let by_ref = states + .iter() + .cloned() + .map(|state| (state.note_ref.clone(), state)) + .collect::>(); + let mut seen_paths = HashSet::new(); + let mut seen_refs = HashSet::new(); + let mut report = GardenSyncReport::default(); for path in paths { let contents = fs::read_to_string(&path) .with_context(|| format!("failed to read {}", path.display()))?; @@ -64,9 +100,97 @@ impl MemoryGarden { .to_string_lossy() .replace('\\', "/"); let input = parse_note_file(&contents, relative)?; - records.push(memory.upsert_markdown_note(&input)?); + let relative = input.path.clone().unwrap_or_default(); + let note_ref = input.note_ref.clone().unwrap_or_default(); + let file_hash = simple_hash(&contents); + let previous = by_path.get(&relative).or_else(|| by_ref.get(¬e_ref)); + let store_hash = memory.markdown_note_hash(¬e_ref)?; + let file_changed = previous.is_none_or(|state| state.file_hash != file_hash); + let store_changed = previous.is_some_and(|state| { + store_hash + .as_deref() + .is_some_and(|hash| hash != state.store_hash) + }); + seen_paths.insert(relative.clone()); + seen_refs.insert(note_ref.clone()); + + if file_changed && store_changed { + report.outcomes.push(GardenSyncOutcome { + path: relative, + note_ref, + action: "conflict".to_string(), + message: Some("file and store both changed since the last sync".to_string()), + }); + continue; + } + + let action = if previous.is_none() { + "created" + } else if previous.is_some_and(|state| state.path != relative) { + "renamed" + } else if file_changed { + "updated" + } else if store_changed { + report.outcomes.push(GardenSyncOutcome { + path: relative, + note_ref, + action: "conflict".to_string(), + message: Some( + "store changed; store-to-file export is intentionally manual".to_string(), + ), + }); + continue; + } else { + "unchanged" + }; + let record = memory.upsert_markdown_note(&input)?; + let store_hash = memory + .markdown_note_hash(&record.note_ref)? + .ok_or_else(|| { + anyhow::anyhow!("synced note '{}' has no store hash", record.note_ref) + })?; + if let Some(old) = previous.filter(|state| state.path != relative) { + memory.remove_garden_sync(&old.path)?; + } + memory.record_garden_sync(&relative, &record.note_ref, &file_hash, &store_hash)?; + report.outcomes.push(GardenSyncOutcome { + path: relative, + note_ref: record.note_ref, + action: action.to_string(), + message: None, + }); + } + + for GardenSyncState { + path, + note_ref, + store_hash, + .. + } in states + { + if seen_paths.contains(&path) || seen_refs.contains(¬e_ref) { + continue; + } + let current_store_hash = memory.markdown_note_hash(¬e_ref)?; + if current_store_hash.as_deref() != Some(store_hash.as_str()) { + report.outcomes.push(GardenSyncOutcome { + path, + note_ref, + action: "conflict".to_string(), + message: Some("file was deleted after the store changed".to_string()), + }); + continue; + } + memory.tombstone_markdown_note(¬e_ref)?; + memory.remove_garden_sync(&path)?; + report.outcomes.push(GardenSyncOutcome { + path, + note_ref, + action: "deleted".to_string(), + message: None, + }); } - Ok(records) + Ok(report) } } @@ -271,3 +395,83 @@ fn sanitize_file_stem(value: &str) -> String { }) .collect() } + +#[cfg(test)] +mod tests { + use super::*; + use tempfile::{NamedTempFile, TempDir}; + + fn note(note_ref: &str, body: &str) -> MarkdownNoteInput { + MarkdownNoteInput { + note_ref: Some(note_ref.to_string()), + lane: MemoryLane::Semantic, + kind: "zettel_note".to_string(), + title: "garden synchronization is explicit".to_string(), + path: Some(format!("semantic/{note_ref}.md")), + body: body.to_string(), + sources: vec![], + follow: None, + entities: vec![], + status: "active".to_string(), + frontmatter: serde_json::json!({}), + } + } + + #[test] + fn garden_sync_tracks_rename_and_tombstones_deleted_files() -> Result<()> { + let root = TempDir::new()?; + let db = NamedTempFile::new()?; + let store = MemoryStore::open(db.path().to_str().unwrap(), 4)?; + let garden = MemoryGarden::new(root.path()); + let original = garden.write_note(note("zgarden", "initial body"))?; + + let created = garden.sync_to_store(&store)?; + assert_eq!(created.outcomes[0].action, "created"); + let renamed = root.path().join("semantic/renamed.md"); + fs::rename(&original, &renamed)?; + let rename_report = garden.sync_to_store(&store)?; + assert!(rename_report + .outcomes + .iter() + .any(|outcome| outcome.action == "renamed")); + + fs::remove_file(&renamed)?; + let delete_report = garden.sync_to_store(&store)?; + assert!(delete_report + .outcomes + .iter() + .any(|outcome| outcome.action == "deleted")); + assert!(store + .search_refs_fts("initial body", &[MemoryLane::Semantic], 5)? + .is_empty()); + Ok(()) + } + + #[test] + fn garden_sync_reports_two_sided_edit_without_overwriting_store() -> Result<()> { + let root = TempDir::new()?; + let db = NamedTempFile::new()?; + let store = MemoryStore::open(db.path().to_str().unwrap(), 4)?; + let garden = MemoryGarden::new(root.path()); + let path = garden.write_note(note("zconflict", "initial body"))?; + garden.sync_to_store(&store)?; + + store.upsert_markdown_note(¬e("zconflict", "store edit"))?; + fs::write( + &path, + render_note_file(¬e("zconflict", "file edit"), "zconflict"), + )?; + let report = garden.sync_to_store(&store)?; + + assert!(report.has_conflicts()); + assert!(store + .search_refs_fts("store edit", &[MemoryLane::Semantic], 5)? + .iter() + .any(|entry| entry.alias.as_deref() == Some("zconflict"))); + assert!(!store + .search_refs_fts("file edit", &[MemoryLane::Semantic], 5)? + .iter() + .any(|entry| entry.body.contains("file edit"))); + Ok(()) + } +} diff --git a/klbr-core/src/memory.rs b/klbr-core/src/memory.rs index 4182ad7..1ef3cd7 100644 --- a/klbr-core/src/memory.rs +++ b/klbr-core/src/memory.rs @@ -24,6 +24,16 @@ pub struct HistoryEntry { pub tool_call_id: Option, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct SessionIngestState { + pub session_id: String, + pub status: String, + pub expected_turns: usize, + pub completed_turns: usize, + pub attempt_count: usize, + pub last_error: Option, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct MemorySummary { pub id: i64, @@ -182,6 +192,14 @@ pub struct MarkdownNoteRecord { pub chunk_refs: Vec, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct GardenSyncState { + pub path: String, + pub note_ref: String, + pub file_hash: String, + pub store_hash: String, +} + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct MemoryStoreStats { pub memories: i64, @@ -195,10 +213,32 @@ pub struct MemoryStoreStats { pub db_bytes: i64, } +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub struct ProjectionAudit { + pub missing_fts: i64, + pub stale_fts: i64, + pub missing_trigram: i64, + pub stale_trigram: i64, + pub missing_ref_metadata: i64, + pub orphan_embedding_items: i64, +} + +impl ProjectionAudit { + pub fn is_healthy(&self) -> bool { + self.missing_fts == 0 + && self.stale_fts == 0 + && self.missing_trigram == 0 + && self.stale_trigram == 0 + && self.missing_ref_metadata == 0 + && self.orphan_embedding_items == 0 + } +} + /// sqlite-backed episodic memory store using sqlite-vec for cosine ANN. #[derive(Clone)] pub struct MemoryStore { conn: Arc>, + db_path: Arc, embed_dim: usize, } @@ -216,6 +256,7 @@ impl MemoryStore { } let store = Self { conn: Arc::new(Mutex::new(conn)), + db_path: Arc::new(path.to_string()), embed_dim, }; store.init_schema()?; @@ -228,6 +269,18 @@ impl MemoryStore { &self.conn } + pub fn read_handle(&self) -> Result { + let conn = Connection::open(self.db_path.as_str())?; + conn.execute("PRAGMA foreign_keys = ON;", [])?; + conn.busy_timeout(std::time::Duration::from_millis(5000))?; + conn.execute("PRAGMA query_only = ON;", [])?; + Ok(Self { + conn: Arc::new(Mutex::new(conn)), + db_path: self.db_path.clone(), + embed_dim: self.embed_dim, + }) + } + fn init_schema(&self) -> Result<()> { let conn = self.conn.lock().unwrap(); conn.execute_batch(&format!( @@ -269,6 +322,16 @@ impl MemoryStore { messages TEXT NOT NULL, ts INTEGER NOT NULL DEFAULT (unixepoch()) ); + CREATE TABLE IF NOT EXISTS session_ingests ( + session_id TEXT PRIMARY KEY, + status TEXT NOT NULL CHECK (status IN ('running', 'failed', 'complete')), + expected_turns INTEGER NOT NULL CHECK (expected_turns >= 0), + completed_turns INTEGER NOT NULL DEFAULT 0 CHECK (completed_turns >= 0), + attempt_count INTEGER NOT NULL DEFAULT 1 CHECK (attempt_count > 0), + last_error TEXT, + created_at INTEGER NOT NULL DEFAULT (unixepoch()), + updated_at INTEGER NOT NULL DEFAULT (unixepoch()) + ) STRICT; CREATE TABLE IF NOT EXISTS memory_edges ( id INTEGER PRIMARY KEY, from_memory_id INTEGER NOT NULL, @@ -407,6 +470,13 @@ impl MemoryStore { created_at INTEGER NOT NULL DEFAULT (unixepoch()), UNIQUE(note_id, ord) ) STRICT; + CREATE TABLE IF NOT EXISTS garden_sync_state ( + path TEXT PRIMARY KEY, + note_ref TEXT NOT NULL, + file_hash TEXT NOT NULL, + store_hash TEXT NOT NULL, + synced_at INTEGER NOT NULL DEFAULT (unixepoch()) + ) STRICT; CREATE VIRTUAL TABLE IF NOT EXISTS promptable_text_fts USING fts5( ref_id UNINDEXED, @@ -510,6 +580,7 @@ impl MemoryStore { DROP TABLE IF EXISTS promptable_text_trigram; DROP TABLE IF EXISTS promptable_text_fts; DROP TABLE IF EXISTS markdown_note_chunks; + DROP TABLE IF EXISTS garden_sync_state; DROP TABLE IF EXISTS markdown_notes; DROP TABLE IF EXISTS ref_metadata; DROP TABLE IF EXISTS promptable_text; @@ -523,6 +594,7 @@ impl MemoryStore { DROP TABLE IF EXISTS memory_tombstones; DROP TABLE IF EXISTS memories; DROP TABLE IF EXISTS turns; + DROP TABLE IF EXISTS session_ingests; DROP TABLE IF EXISTS context_snapshots; PRAGMA foreign_keys = ON;", )?; @@ -640,6 +712,13 @@ impl MemoryStore { created_at INTEGER NOT NULL DEFAULT (unixepoch()), UNIQUE(note_id, ord) ) STRICT; + CREATE TABLE IF NOT EXISTS garden_sync_state ( + path TEXT PRIMARY KEY, + note_ref TEXT NOT NULL, + file_hash TEXT NOT NULL, + store_hash TEXT NOT NULL, + synced_at INTEGER NOT NULL DEFAULT (unixepoch()) + ) STRICT; CREATE INDEX IF NOT EXISTS ref_metadata_lane ON ref_metadata(lane, kind); CREATE INDEX IF NOT EXISTS markdown_notes_lane ON markdown_notes(lane, kind, status); CREATE INDEX IF NOT EXISTS markdown_note_chunks_note ON markdown_note_chunks(note_id, ord);", @@ -735,15 +814,16 @@ impl MemoryStore { } let tags_json = serde_json::to_string(&input.tags).unwrap_or_else(|_| "[]".into()); - let conn = self.conn.lock().unwrap(); + let mut conn = self.conn.lock().unwrap(); + let tx = conn.transaction()?; if input.memory_id.is_none() && input.source_ref.as_deref() != Some(ANCHOR_SOURCE_REF) { if let Some((id, existing_tags, existing_pinned)) = - find_duplicate_memory(&conn, input.namespace.as_str(), input.text.as_str())? + find_duplicate_memory(&tx, input.namespace.as_str(), input.text.as_str())? { let merged_tags = merge_tags(existing_tags, input.tags.clone()); let pinned = existing_pinned || input.pinned; - conn.execute( + tx.execute( "UPDATE memories SET tags = ?1, pinned = ?2, ingest_time = ?3, ts = ?3 WHERE id = ?4", @@ -754,12 +834,13 @@ impl MemoryStore { id ], )?; + tx.commit()?; return Ok(id); } } if let Some(id) = input.memory_id { - conn.execute( + tx.execute( "INSERT INTO memories ( id, namespace, layer, content, pinned, tags, event_time, ingest_time, embedding_model, embedding_dim, embedding_version, status, source_ref @@ -780,22 +861,23 @@ impl MemoryStore { input.source_ref.as_deref(), ], )?; - conn.execute( + tx.execute( "INSERT INTO vec_memories (rowid, embedding) VALUES (?1, ?2)", params![id, f32s_to_bytes(&input.embedding)], )?; self.register_memory_ref( - &conn, + &tx, id, input.text.as_str(), &tags_json, &input.status, input.source_ref.as_deref(), )?; + tx.commit()?; return Ok(id); } - conn.execute( + tx.execute( "INSERT INTO memories ( namespace, layer, content, pinned, tags, event_time, ingest_time, embedding_model, embedding_dim, embedding_version, status, source_ref @@ -815,19 +897,20 @@ impl MemoryStore { input.source_ref.as_deref(), ], )?; - let id = conn.last_insert_rowid(); - conn.execute( + let id = tx.last_insert_rowid(); + tx.execute( "INSERT INTO vec_memories (rowid, embedding) VALUES (?1, ?2)", params![id, f32s_to_bytes(&input.embedding)], )?; self.register_memory_ref( - &conn, + &tx, id, input.text.as_str(), &tags_json, &input.status, input.source_ref.as_deref(), )?; + tx.commit()?; Ok(id) } @@ -892,20 +975,10 @@ impl MemoryStore { params![&ver_alias, &ref_ver], )?; - // 7. Insert promptable_text for memory and version 1 - let tokens = content.chars().count() / 4; - conn.execute( - "INSERT INTO promptable_text (ref_id, body, token_count, body_hash, tokenizer_id, computed_at) - VALUES (?1, ?2, ?3, ?4, 'char_count_div_4', ?5)", - params![&ref_mem, content, tokens as i64, &hash, timestamp], - )?; upsert_ref_metadata(conn, &ref_mem, lane, "memory")?; upsert_ref_metadata(conn, &ref_ver, lane, "memory_version")?; - conn.execute( - "INSERT INTO promptable_text (ref_id, body, token_count, body_hash, tokenizer_id, computed_at) - VALUES (?1, ?2, ?3, ?4, 'char_count_div_4', ?5)", - params![&ref_ver, content, tokens as i64, &hash, timestamp], - )?; + upsert_promptable_text(conn, &ref_mem, content, &hash, timestamp)?; + upsert_promptable_text(conn, &ref_ver, content, &hash, timestamp)?; Ok(()) } @@ -1034,22 +1107,21 @@ impl MemoryStore { } } - // Update promptable_text (cache) for the parent memory ref with the new body - let tokens = content.chars().count() / 4; - conn.execute( - "INSERT INTO promptable_text (ref_id, body, token_count, body_hash, tokenizer_id, computed_at) - VALUES (?1, ?2, ?3, ?4, 'char_count_div_4', ?5) - ON CONFLICT(ref_id) DO UPDATE SET body = excluded.body, token_count = excluded.token_count, body_hash = excluded.body_hash, computed_at = excluded.computed_at", - params![&ref_mem, content, tokens as i64, &hash, now], - )?; - - // Insert promptable_text for the new version ref - conn.execute( - "INSERT INTO promptable_text (ref_id, body, token_count, body_hash, tokenizer_id, computed_at) - VALUES (?1, ?2, ?3, ?4, 'char_count_div_4', ?5) - ON CONFLICT(ref_id) DO UPDATE SET body = excluded.body, token_count = excluded.token_count, body_hash = excluded.body_hash, computed_at = excluded.computed_at", - params![&ref_ver, content, tokens as i64, &hash, now], - )?; + upsert_promptable_text(&conn, &ref_mem, content, &hash, now)?; + upsert_promptable_text(&conn, &ref_ver, content, &hash, now)?; + if let Some(previous_version_id) = prev_version_id { + if let Some(previous_ref) = conn + .query_row( + "SELECT ref_id FROM refs + WHERE entity_type = 'memory_version' AND entity_id = ?1", + params![previous_version_id], + |row| row.get::<_, String>(0), + ) + .optional()? + { + refresh_promptable_fts_ref(&conn, &previous_ref)?; + } + } Ok(()) } @@ -1087,7 +1159,7 @@ impl MemoryStore { params![status_to_str(&status), unix_timestamp(), id], )?; set_memory_ref_status(&conn, id, &status)?; - rebuild_promptable_fts(&conn)?; + refresh_memory_projection_refs(&conn, id)?; Ok(()) } @@ -1119,17 +1191,21 @@ impl MemoryStore { params![id, reason, now], )?; - for descendant in provenance_descendants(&conn, id)? { + let descendants = provenance_descendants(&conn, id)?; + for descendant in &descendants { conn.execute( "UPDATE memories SET status = 'suppressed', pinned = 0, ingest_time = ?1, ts = ?1 WHERE id = ?2 AND status != 'tombstoned'", params![now, descendant], )?; - set_memory_ref_status(&conn, descendant, &MemoryStatus::Suppressed)?; + set_memory_ref_status(&conn, *descendant, &MemoryStatus::Suppressed)?; } set_memory_ref_status(&conn, id, &MemoryStatus::Tombstoned)?; - rebuild_promptable_fts(&conn)?; + refresh_memory_projection_refs(&conn, id)?; + for descendant in descendants { + refresh_memory_projection_refs(&conn, descendant)?; + } Ok(()) } @@ -1269,6 +1345,10 @@ impl MemoryStore { } pub fn sync_reference_indexes(&self) -> Result<()> { + self.repair_reference_indexes() + } + + pub fn repair_reference_indexes(&self) -> Result<()> { let conn = self.conn.lock().unwrap(); backfill_memory_refs(&conn)?; backfill_turn_refs(&conn)?; @@ -1279,13 +1359,150 @@ impl MemoryStore { Ok(()) } - pub fn upsert_markdown_note(&self, input: &MarkdownNoteInput) -> Result { + pub fn audit_reference_indexes(&self) -> Result { let conn = self.conn.lock().unwrap(); - let record = upsert_markdown_note_record(&conn, input)?; - rebuild_promptable_fts(&conn)?; + let count = |sql: &str| conn.query_row(sql, [], |row| row.get::<_, i64>(0)); + Ok(ProjectionAudit { + missing_fts: count( + "SELECT COUNT(*) FROM promptable_text p + JOIN refs r ON r.ref_id = p.ref_id AND r.status = 'active' + LEFT JOIN promptable_text_fts f ON f.ref_id = p.ref_id + WHERE f.ref_id IS NULL", + )?, + stale_fts: count( + "SELECT COUNT(*) FROM promptable_text_fts f + LEFT JOIN refs r ON r.ref_id = f.ref_id + WHERE r.ref_id IS NULL OR r.status != 'active'", + )?, + missing_trigram: count( + "SELECT COUNT(*) FROM promptable_text p + JOIN refs r ON r.ref_id = p.ref_id AND r.status = 'active' + LEFT JOIN promptable_text_trigram f ON f.ref_id = p.ref_id + WHERE f.ref_id IS NULL", + )?, + stale_trigram: count( + "SELECT COUNT(*) FROM promptable_text_trigram f + LEFT JOIN refs r ON r.ref_id = f.ref_id + WHERE r.ref_id IS NULL OR r.status != 'active'", + )?, + missing_ref_metadata: count( + "SELECT COUNT(*) FROM refs r + LEFT JOIN ref_metadata m ON m.ref_id = r.ref_id + WHERE m.ref_id IS NULL", + )?, + orphan_embedding_items: count( + "SELECT COUNT(*) FROM embedding_items e + LEFT JOIN refs r ON r.ref_id = e.ref_id + WHERE r.ref_id IS NULL", + )?, + }) + } + + pub fn upsert_markdown_note(&self, input: &MarkdownNoteInput) -> Result { + let mut conn = self.conn.lock().unwrap(); + let tx = conn.transaction()?; + let record = upsert_markdown_note_record(&tx, input)?; + tx.commit()?; Ok(record) } + pub fn markdown_note_hash(&self, note_ref: &str) -> Result> { + let conn = self.conn.lock().unwrap(); + conn.query_row( + "SELECT body_hash FROM markdown_notes WHERE note_ref = ?1", + params![note_ref], + |row| row.get(0), + ) + .optional() + .map_err(Into::into) + } + + pub fn garden_sync_states(&self) -> Result> { + let conn = self.conn.lock().unwrap(); + let mut stmt = conn.prepare( + "SELECT path, note_ref, file_hash, store_hash + FROM garden_sync_state ORDER BY path ASC", + )?; + let states = stmt + .query_map([], |row| { + Ok(GardenSyncState { + path: row.get(0)?, + note_ref: row.get(1)?, + file_hash: row.get(2)?, + store_hash: row.get(3)?, + }) + })? + .collect::>>()?; + Ok(states) + } + + pub fn record_garden_sync( + &self, + path: &str, + note_ref: &str, + file_hash: &str, + store_hash: &str, + ) -> Result<()> { + self.conn.lock().unwrap().execute( + "INSERT INTO garden_sync_state (path, note_ref, file_hash, store_hash, synced_at) + VALUES (?1, ?2, ?3, ?4, unixepoch()) + ON CONFLICT(path) DO UPDATE SET + note_ref = excluded.note_ref, + file_hash = excluded.file_hash, + store_hash = excluded.store_hash, + synced_at = excluded.synced_at", + params![path, note_ref, file_hash, store_hash], + )?; + Ok(()) + } + + pub fn remove_garden_sync(&self, path: &str) -> Result<()> { + self.conn.lock().unwrap().execute( + "DELETE FROM garden_sync_state WHERE path = ?1", + params![path], + )?; + Ok(()) + } + + pub fn tombstone_markdown_note(&self, note_ref: &str) -> Result<()> { + let mut conn = self.conn.lock().unwrap(); + let tx = conn.transaction()?; + let note_id = tx + .query_row( + "SELECT note_id FROM markdown_notes WHERE note_ref = ?1", + params![note_ref], + |row| row.get::<_, i64>(0), + ) + .optional()? + .ok_or_else(|| anyhow::anyhow!("markdown note '{note_ref}' not found"))?; + tx.execute( + "UPDATE markdown_notes SET status = 'tombstoned', updated_at = unixepoch() + WHERE note_id = ?1", + params![note_id], + )?; + let refs = { + let mut stmt = tx.prepare( + "SELECT r.ref_id FROM refs r + WHERE (r.entity_type IN ('episode', 'semantic_note', 'profile_note', 'procedural_note') + AND r.entity_id = ?1)", + )?; + let refs = stmt + .query_map(params![note_id], |row| row.get::<_, String>(0))? + .collect::>>()?; + refs + }; + for ref_id in refs { + tx.execute( + "UPDATE refs SET status = 'tombstoned', deleted_at = unixepoch() + WHERE ref_id = ?1", + params![&ref_id], + )?; + refresh_promptable_fts_ref(&tx, &ref_id)?; + } + tx.commit()?; + Ok(()) + } + pub fn active_promptable_refs( &self, lanes: &[MemoryLane], @@ -2501,10 +2718,12 @@ impl MemoryStore { pub fn tombstone_ref(&self, ref_id: &str) -> Result<()> { let conn = self.conn.lock().unwrap(); + let canonical = resolve_ref_alias(&conn, ref_id)?.unwrap_or_else(|| ref_id.to_string()); conn.execute( "UPDATE refs SET status = 'tombstoned', deleted_at = unixepoch() WHERE ref_id = ?1", - params![ref_id], + params![&canonical], )?; + refresh_promptable_fts_ref(&conn, &canonical)?; Ok(()) } @@ -2533,26 +2752,30 @@ impl MemoryStore { anyhow::bail!("superseded ref '{}' not found in refs table", old_ref); } let edge_id = insert_ref_edge(&conn, &new_ref_id, &old_ref_id, "supersedes", "{}")?; - rebuild_promptable_fts(&conn)?; + refresh_promptable_fts_ref(&conn, &old_ref_id)?; + refresh_promptable_fts_ref(&conn, &new_ref_id)?; Ok(edge_id) } pub fn suppress_ref(&self, ref_id: &str) -> Result<()> { let conn = self.conn.lock().unwrap(); + let canonical = resolve_ref_alias(&conn, ref_id)?.unwrap_or_else(|| ref_id.to_string()); conn.execute( "UPDATE refs SET status = 'suppressed' WHERE ref_id = ?1", - params![ref_id], + params![&canonical], )?; + refresh_promptable_fts_ref(&conn, &canonical)?; Ok(()) } pub fn purge_ref(&self, ref_id: &str) -> Result<()> { let conn = self.conn.lock().unwrap(); + let canonical = resolve_ref_alias(&conn, ref_id)?.unwrap_or_else(|| ref_id.to_string()); let info: Option<(String, i64)> = conn .query_row( "SELECT entity_type, entity_id FROM refs WHERE ref_id = ?1", - params![ref_id], + params![&canonical], |row| Ok((row.get(0)?, row.get(1)?)), ) .optional()?; @@ -2561,7 +2784,7 @@ impl MemoryStore { if entity_type == "turn_chunk" { conn.execute( "UPDATE turn_chunks SET raw_text = '' WHERE ref_id = ?1", - params![ref_id], + params![&canonical], )?; conn.execute( "UPDATE turns SET content = '' WHERE id = ?1", @@ -2582,13 +2805,14 @@ impl MemoryStore { conn.execute( "DELETE FROM promptable_text WHERE ref_id = ?1", - params![ref_id], + params![&canonical], )?; conn.execute( "UPDATE refs SET status = 'purged', deleted_at = unixepoch() WHERE ref_id = ?1", - params![ref_id], + params![&canonical], )?; + refresh_promptable_fts_ref(&conn, &canonical)?; Ok(()) } @@ -3031,6 +3255,104 @@ impl MemoryStore { ) } + pub fn begin_session_ingest(&self, session_id: &str, expected_turns: usize) -> Result<()> { + self.conn.lock().unwrap().execute( + "INSERT INTO session_ingests ( + session_id, status, expected_turns, completed_turns, attempt_count, + last_error, created_at, updated_at + ) VALUES (?1, 'running', ?2, 0, 1, NULL, unixepoch(), unixepoch()) + ON CONFLICT(session_id) DO UPDATE SET + status = 'running', + expected_turns = excluded.expected_turns, + attempt_count = session_ingests.attempt_count + 1, + last_error = NULL, + updated_at = unixepoch()", + params![session_id, expected_turns as i64], + )?; + Ok(()) + } + + pub fn session_turn( + &self, + session_id: &str, + turn_index: usize, + ) -> Result> { + let conn = self.conn.lock().unwrap(); + conn.query_row( + "SELECT id, role, content, thinking, tool_calls, tool_call_id, ts + FROM turns + WHERE json_extract(metadata, '$.session_id') = ?1 + AND json_extract(metadata, '$.turn_index') = ?2 + ORDER BY id ASC + LIMIT 1", + params![session_id, turn_index as i64], + row_to_history_entry, + ) + .optional() + .map_err(Into::into) + } + + pub fn update_session_ingest_progress( + &self, + session_id: &str, + completed_turns: usize, + ) -> Result<()> { + self.conn.lock().unwrap().execute( + "UPDATE session_ingests + SET completed_turns = MAX(completed_turns, ?2), updated_at = unixepoch() + WHERE session_id = ?1", + params![session_id, completed_turns as i64], + )?; + Ok(()) + } + + pub fn fail_session_ingest(&self, session_id: &str, error: &str) -> Result<()> { + self.conn.lock().unwrap().execute( + "UPDATE session_ingests + SET status = 'failed', last_error = ?2, updated_at = unixepoch() + WHERE session_id = ?1", + params![session_id, truncate_error(error)], + )?; + Ok(()) + } + + pub fn complete_session_ingest(&self, session_id: &str) -> Result<()> { + let conn = self.conn.lock().unwrap(); + let changed = conn.execute( + "UPDATE session_ingests + SET status = 'complete', completed_turns = expected_turns, + last_error = NULL, updated_at = unixepoch() + WHERE session_id = ?1 AND completed_turns >= expected_turns", + params![session_id], + )?; + if changed == 0 { + anyhow::bail!("session ingest '{session_id}' is missing turns or was not started"); + } + Ok(()) + } + + pub fn session_ingest_state(&self, session_id: &str) -> Result> { + let conn = self.conn.lock().unwrap(); + conn.query_row( + "SELECT session_id, status, expected_turns, completed_turns, + attempt_count, last_error + FROM session_ingests WHERE session_id = ?1", + params![session_id], + |row| { + Ok(SessionIngestState { + session_id: row.get(0)?, + status: row.get(1)?, + expected_turns: row.get::<_, i64>(2)? as usize, + completed_turns: row.get::<_, i64>(3)? as usize, + attempt_count: row.get::<_, i64>(4)? as usize, + last_error: row.get(5)?, + }) + }, + ) + .optional() + .map_err(Into::into) + } + pub fn log_turn_with_tools( &self, role: &str, @@ -3062,8 +3384,9 @@ impl MemoryStore { ) -> Result { let tool_calls_json = tool_calls.map(serde_json::to_string).transpose()?; let lane = infer_turn_lane(metadata_json); - let conn = self.conn.lock().unwrap(); - conn.execute( + let mut conn = self.conn.lock().unwrap(); + let tx = conn.transaction()?; + tx.execute( "INSERT INTO turns (role, content, thinking, tool_calls, tool_call_id, ts, metadata) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)", params![ @@ -3077,23 +3400,23 @@ impl MemoryStore { ], )?; - let turn_id = conn.last_insert_rowid(); + let turn_id = tx.last_insert_rowid(); let turn_base36 = to_base36(turn_id as u64); let turn_alias = format!("d{}", turn_base36); let turn_ref = generate_canonical_id(); // 1. Register the parent turn ref and display alias - conn.execute( + tx.execute( "INSERT INTO refs (ref_id, entity_type, entity_id, status) VALUES (?1, 'synthetic', ?2, 'active')", params![&turn_ref, turn_id], )?; - conn.execute( + tx.execute( "INSERT INTO ref_aliases (alias, ref_id, alias_kind, status) VALUES (?1, ?2, 'display', 'active')", params![&turn_alias, &turn_ref], )?; - upsert_ref_metadata(&conn, &turn_ref, lane, "turn")?; + upsert_ref_metadata(&tx, &turn_ref, lane, "turn")?; if role == "user" || role == "assistant" { let chunks = extract_markdown_chunks_with_offsets(content); @@ -3105,21 +3428,21 @@ impl MemoryStore { let hash = simple_hash(&chunk); // 1. Insert into refs - conn.execute( + tx.execute( "INSERT INTO refs (ref_id, entity_type, entity_id, status, content_hash) VALUES (?1, 'turn_chunk', ?2, 'active', ?3)", params![&chunk_ref, turn_id, &hash], )?; // 2. Insert into ref_aliases - conn.execute( + tx.execute( "INSERT INTO ref_aliases (alias, ref_id, alias_kind, status) VALUES (?1, ?2, 'display', 'active')", params![&chunk_alias, &chunk_ref], )?; // 3. Insert into turn_chunks - conn.execute( + tx.execute( "INSERT INTO turn_chunks (ref_id, turn_id, ord, raw_text, raw_hash, byte_start, byte_end, parser_name, parser_version, parser_options) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 'pulldown-cmark', ?8, '{}')", params![ @@ -3134,14 +3457,8 @@ impl MemoryStore { ], )?; - // 4. Insert into promptable_text (cache) - let tokens = chunk.chars().count() / 4; - conn.execute( - "INSERT INTO promptable_text (ref_id, body, token_count, body_hash, tokenizer_id, computed_at) - VALUES (?1, ?2, ?3, ?4, 'char_count_div_4', ?5)", - params![&chunk_ref, &chunk, tokens as i64, &hash, timestamp], - )?; - upsert_ref_metadata(&conn, &chunk_ref, lane, "turn_chunk")?; + upsert_ref_metadata(&tx, &chunk_ref, lane, "turn_chunk")?; + upsert_promptable_text(&tx, &chunk_ref, &chunk, &hash, timestamp)?; } } else { // For other roles, register as a single chunk 'a' @@ -3150,33 +3467,29 @@ impl MemoryStore { let chunk_ref = generate_canonical_id(); let hash = simple_hash(content); - conn.execute( + tx.execute( "INSERT INTO refs (ref_id, entity_type, entity_id, status, content_hash) VALUES (?1, 'turn_chunk', ?2, 'active', ?3)", params![&chunk_ref, turn_id, &hash], )?; - conn.execute( + tx.execute( "INSERT INTO ref_aliases (alias, ref_id, alias_kind, status) VALUES (?1, ?2, 'display', 'active')", params![&chunk_alias, &chunk_ref], )?; - conn.execute( + tx.execute( "INSERT INTO turn_chunks (ref_id, turn_id, ord, raw_text, raw_hash, byte_start, byte_end, parser_name, parser_version, parser_options) VALUES (?1, ?2, 0, ?3, ?4, 0, ?5, 'fallback', ?6, '{}')", params![&chunk_ref, turn_id, content, &hash, content.len() as i64, env!("CARGO_PKG_VERSION")], )?; - let tokens = content.chars().count() / 4; - conn.execute( - "INSERT INTO promptable_text (ref_id, body, token_count, body_hash, tokenizer_id, computed_at) - VALUES (?1, ?2, ?3, ?4, 'char_count_div_4', ?5)", - params![&chunk_ref, content, tokens as i64, &hash, timestamp], - )?; - upsert_ref_metadata(&conn, &chunk_ref, lane, "turn_chunk")?; + upsert_ref_metadata(&tx, &chunk_ref, lane, "turn_chunk")?; + upsert_promptable_text(&tx, &chunk_ref, content, &hash, timestamp)?; } + tx.commit()?; Ok(HistoryEntry { id: turn_id, timestamp, @@ -3710,12 +4023,25 @@ fn upsert_markdown_note_record( )", params![note_id, chunks.len() as i64], )?; + let note_chunk_refs = { + let mut stmt = conn.prepare( + "SELECT ref_id FROM markdown_note_chunks WHERE note_id = ?1 ORDER BY ord ASC", + )?; + let refs = stmt + .query_map(params![note_id], |row| row.get::<_, String>(0))? + .collect::>>()?; + refs + }; + refresh_promptable_fts_ref(conn, ¬e_ref_id)?; + for chunk_ref in note_chunk_refs { + refresh_promptable_fts_ref(conn, &chunk_ref)?; + } for source in &input.sources { - let _ = insert_ref_edge(conn, ¬e_ref, source, "derived_from", "{}"); + insert_ref_edge(conn, ¬e_ref, source, "derived_from", "{}")?; } if let Some(follow) = &input.follow { - let _ = insert_ref_edge(conn, ¬e_ref, follow, "continues", "{}"); + insert_ref_edge(conn, ¬e_ref, follow, "continues", "{}")?; } Ok(MarkdownNoteRecord { @@ -3959,6 +4285,55 @@ fn upsert_promptable_text( computed_at = excluded.computed_at", params![ref_id, body, tokens as i64, hash, timestamp], )?; + refresh_promptable_fts_ref(conn, ref_id)?; + Ok(()) +} + +fn refresh_promptable_fts_ref(conn: &Connection, ref_id: &str) -> Result<()> { + conn.execute( + "DELETE FROM promptable_text_fts WHERE ref_id = ?1", + params![ref_id], + )?; + conn.execute( + "DELETE FROM promptable_text_trigram WHERE ref_id = ?1", + params![ref_id], + )?; + conn.execute( + "INSERT INTO promptable_text_fts (ref_id, body, lane, entity_type) + SELECT p.ref_id, p.body, COALESCE(m.lane, 'semantic'), r.entity_type + FROM promptable_text p + JOIN refs r ON r.ref_id = p.ref_id + LEFT JOIN ref_metadata m ON m.ref_id = p.ref_id + WHERE p.ref_id = ?1 AND r.status = 'active'", + params![ref_id], + )?; + conn.execute( + "INSERT INTO promptable_text_trigram (ref_id, body, lane, entity_type) + SELECT p.ref_id, p.body, COALESCE(m.lane, 'semantic'), r.entity_type + FROM promptable_text p + JOIN refs r ON r.ref_id = p.ref_id + LEFT JOIN ref_metadata m ON m.ref_id = p.ref_id + WHERE p.ref_id = ?1 AND r.status = 'active'", + params![ref_id], + )?; + Ok(()) +} + +fn refresh_memory_projection_refs(conn: &Connection, memory_id: i64) -> Result<()> { + let mut stmt = conn.prepare( + "SELECT ref_id FROM refs + WHERE (entity_type = 'memory' AND entity_id = ?1) + OR (entity_type = 'memory_version' AND entity_id IN ( + SELECT version_id FROM memory_versions WHERE memory_id = ?1 + ))", + )?; + let refs = stmt + .query_map(params![memory_id], |row| row.get::<_, String>(0))? + .collect::>>()?; + drop(stmt); + for ref_id in refs { + refresh_promptable_fts_ref(conn, &ref_id)?; + } Ok(()) } @@ -4484,6 +4859,10 @@ fn unix_timestamp() -> i64 { .unwrap_or_default() } +fn truncate_error(error: &str) -> String { + error.chars().take(1_000).collect() +} + #[cfg(test)] mod tests { use super::*; @@ -5308,6 +5687,116 @@ mod tests { Ok(()) } + #[test] + fn turn_ingest_rolls_back_every_projection_when_a_late_write_fails() -> Result<()> { + let tmp = NamedTempFile::new()?; + let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; + store.conn().lock().unwrap().execute_batch( + "CREATE TRIGGER reject_promptable_turn + BEFORE INSERT ON promptable_text + BEGIN + SELECT RAISE(FAIL, 'injected promptable failure'); + END;", + )?; + + assert!(store + .log_turn("user", "this write must roll back", None) + .is_err()); + let conn = store.conn().lock().unwrap(); + for table in [ + "turns", + "refs", + "ref_aliases", + "turn_chunks", + "ref_metadata", + ] { + let count: i64 = + conn.query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| { + row.get(0) + })?; + assert_eq!(count, 0, "{table} retained a partial turn write"); + } + Ok(()) + } + + #[test] + fn markdown_note_rolls_back_when_a_source_ref_is_invalid() -> Result<()> { + let tmp = NamedTempFile::new()?; + let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; + + let result = store.upsert_markdown_note(&MarkdownNoteInput { + note_ref: Some("ninvalid_source".to_string()), + lane: MemoryLane::Semantic, + kind: "zettel_note".to_string(), + title: "invalid source must not partially persist".to_string(), + path: None, + body: "the source does not exist".to_string(), + sources: vec!["missing_ref".to_string()], + follow: None, + entities: vec![], + status: "active".to_string(), + frontmatter: serde_json::json!({}), + }); + + assert!(result.is_err()); + let conn = store.conn().lock().unwrap(); + let note_count: i64 = + conn.query_row("SELECT COUNT(*) FROM markdown_notes", [], |row| row.get(0))?; + assert_eq!(note_count, 0); + let alias_count: i64 = conn.query_row( + "SELECT COUNT(*) FROM ref_aliases WHERE alias = 'ninvalid_source'", + [], + |row| row.get(0), + )?; + assert_eq!(alias_count, 0); + Ok(()) + } + + #[test] + fn session_ingest_state_supports_idempotent_turn_reuse() -> Result<()> { + let tmp = NamedTempFile::new()?; + let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; + store.begin_session_ingest("session-a", 1)?; + let metadata = serde_json::json!({ + "session_id": "session-a", + "turn_index": 0, + "benchmark_session": true, + }); + let inserted = store.log_turn_with_metadata_at("user", "hello", None, &metadata, 42)?; + store.update_session_ingest_progress("session-a", 1)?; + store.complete_session_ingest("session-a")?; + + let reused = store.session_turn("session-a", 0)?.unwrap(); + assert_eq!(reused.id, inserted.id); + store.begin_session_ingest("session-a", 1)?; + let state = store.session_ingest_state("session-a")?.unwrap(); + assert_eq!(state.status, "running"); + assert_eq!(state.completed_turns, 1); + assert_eq!(state.attempt_count, 2); + Ok(()) + } + + #[test] + fn projection_audit_detects_and_explicit_repair_restores_fts_drift() -> Result<()> { + let tmp = NamedTempFile::new()?; + let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; + store.log_turn("user", "incremental projection sentinel", None)?; + assert!(store.audit_reference_indexes()?.is_healthy()); + + store.conn().lock().unwrap().execute( + "DELETE FROM promptable_text_fts + WHERE ref_id = (SELECT ref_id FROM promptable_text LIMIT 1)", + [], + )?; + let drift = store.audit_reference_indexes()?; + assert_eq!(drift.missing_fts, 1); + assert!(!drift.is_healthy()); + + store.repair_reference_indexes()?; + assert!(store.audit_reference_indexes()?.is_healthy()); + Ok(()) + } + #[test] fn test_markdown_note_upsert_indexes_refs_chunks_and_sources() -> Result<()> { let tmp = NamedTempFile::new()?; diff --git a/klbr-core/src/pipeline.rs b/klbr-core/src/pipeline.rs index a36c1dc..3bc5849 100644 --- a/klbr-core/src/pipeline.rs +++ b/klbr-core/src/pipeline.rs @@ -1,8 +1,10 @@ use std::collections::{HashMap, HashSet, VecDeque}; +use std::sync::Arc; use std::time::Duration; use anyhow::Result; use chrono::{SecondsFormat, TimeZone, Utc}; +use tokio::sync::Semaphore; use crate::{ config::MemoryConfig, @@ -212,7 +214,7 @@ impl PipelineProfile { || lower.contains("hybrid-packet")); Self { name: name.to_string(), - write_episode_memories: !raw_turns && !fts_only, + write_episode_memories: lower.contains("legacy-episode-memory"), write_episode_notes: !raw_turns, lexical: !dense_only, dense: !raw_turns && !fts_only, @@ -228,6 +230,7 @@ pub struct MemoryPipeline { pub llm: LlmClient, pub config: MemoryConfig, pub profile: PipelineProfile, + storage_gate: Arc, } impl MemoryPipeline { @@ -237,6 +240,7 @@ impl MemoryPipeline { llm, config, profile: PipelineProfile::named("klbr-full"), + storage_gate: Arc::new(Semaphore::new(8)), } } @@ -245,6 +249,46 @@ impl MemoryPipeline { self } + async fn read_db(&self, operation: F) -> Result + where + T: Send + 'static, + F: FnOnce(MemoryStore) -> Result + Send + 'static, + { + let permit = self + .storage_gate + .clone() + .acquire_owned() + .await + .map_err(|_| anyhow::anyhow!("storage read gate closed"))?; + let memory = self.memory.clone(); + tokio::task::spawn_blocking(move || { + let _permit = permit; + operation(memory.read_handle()?) + }) + .await + .map_err(|error| anyhow::anyhow!("storage read task failed: {error}"))? + } + + async fn write_db(&self, operation: F) -> Result + where + T: Send + 'static, + F: FnOnce(MemoryStore) -> Result + Send + 'static, + { + let permit = self + .storage_gate + .clone() + .acquire_owned() + .await + .map_err(|_| anyhow::anyhow!("storage write gate closed"))?; + let memory = self.memory.clone(); + tokio::task::spawn_blocking(move || { + let _permit = permit; + operation(memory) + }) + .await + .map_err(|error| anyhow::anyhow!("storage write task failed: {error}"))? + } + async fn plan_query_op(&self, query: &BenchQuery) -> OpPlan { if let Some(op) = query.op_hint { return OpPlan::for_op(op, "query_op_hint"); @@ -270,17 +314,20 @@ impl MemoryPipeline { match completion { Ok(Ok((reply, _))) => parse_structured_op_plan(&reply).unwrap_or_else(|| OpPlan { + status: "parse_failed".to_string(), reason: "planner_parse_failed_default_lookup".to_string(), ..OpPlan::default() }), Ok(Err(error)) => { tracing::warn!("structured query planner failed: {error:?}"); OpPlan { + status: "request_failed".to_string(), reason: "planner_request_failed_default_lookup".to_string(), ..OpPlan::default() } } Err(_) => OpPlan { + status: "timeout".to_string(), reason: "planner_timeout_default_lookup".to_string(), ..OpPlan::default() }, @@ -288,12 +335,52 @@ impl MemoryPipeline { } pub async fn reset(&self, _run: BenchRun) -> Result<()> { - self.memory.reset()?; - self.memory.sync_reference_indexes()?; + self.write_db(|memory| memory.reset()).await?; Ok(()) } pub async fn observe_session(&self, session: BenchSession) -> Result { + let session_id = session.session_id.clone(); + let expected_turns = session.turns.len(); + self.write_db(move |memory| memory.begin_session_ingest(&session_id, expected_turns)) + .await?; + match self.observe_session_inner(&session).await { + Ok(trace) => { + let session_id = session.session_id.clone(); + if let Err(error) = self + .write_db(move |memory| memory.complete_session_ingest(&session_id)) + .await + { + let failed_session = session.session_id.clone(); + let message = error.to_string(); + let _ = self + .write_db(move |memory| { + memory.fail_session_ingest(&failed_session, &message) + }) + .await; + return Err(error); + } + Ok(trace) + } + Err(error) => { + let failed_session = session.session_id.clone(); + let message = error.to_string(); + if let Err(state_error) = self + .write_db(move |memory| memory.fail_session_ingest(&failed_session, &message)) + .await + { + tracing::error!( + session_id = %session.session_id, + error = ?state_error, + "failed to persist session ingest failure state" + ); + } + Err(error) + } + } + } + + async fn observe_session_inner(&self, session: &BenchSession) -> Result { let timestamp = session.timestamp.unwrap_or_else(unix_timestamp); let mut turn_ids = Vec::new(); let mut turn_refs = Vec::new(); @@ -303,18 +390,38 @@ impl MemoryPipeline { let role = normalized_turn_role(&turn.role); let metadata = serde_json::json!({ "lane": "episodic", - "session_id": session.session_id, + "session_id": session.session_id.clone(), "benchmark_session": true, "turn_index": turn_idx, }); - let entry = self.memory.log_turn_with_metadata_at( - role, - &turn.content, - None, - &metadata, - turn.timestamp.unwrap_or(timestamp), - )?; + let lookup_session = session.session_id.clone(); + let existing = self + .read_db(move |memory| memory.session_turn(&lookup_session, turn_idx)) + .await?; + let entry = match existing { + Some(entry) => entry, + None => { + let role = role.to_string(); + let content = turn.content.clone(); + let turn_timestamp = turn.timestamp.unwrap_or(timestamp); + self.write_db(move |memory| { + memory.log_turn_with_metadata_at( + &role, + &content, + None, + &metadata, + turn_timestamp, + ) + }) + .await? + } + }; turn_ids.push(entry.id); + let progress_session = session.session_id.clone(); + self.write_db(move |memory| { + memory.update_session_ingest_progress(&progress_session, turn_idx + 1) + }) + .await?; let turn_alias = format!("d{}", to_base36(entry.id as u64)); turn_refs.push(turn_alias.clone()); @@ -325,9 +432,9 @@ impl MemoryPipeline { } } - let episode_body = render_episode_card(&session, timestamp, &source_refs); + let episode_body = render_episode_card(session, timestamp, &source_refs); if self.profile.write_episode_notes && !episode_body.trim().is_empty() { - let _ = self.memory.upsert_markdown_note(&MarkdownNoteInput { + let note = MarkdownNoteInput { note_ref: Some(format!("e{}", stable_base36(&session.session_id))), lane: MemoryLane::Episodic, kind: "episode_event_card".to_string(), @@ -341,8 +448,10 @@ impl MemoryPipeline { follow: None, entities: vec![format!("session:{}", session.session_id)], status: "active".to_string(), - frontmatter: episode_card_frontmatter(&session, timestamp, &source_refs), - }); + frontmatter: episode_card_frontmatter(session, timestamp, &source_refs), + }; + self.write_db(move |memory| memory.upsert_markdown_note(¬e)) + .await?; } let episode_memory_id = @@ -350,7 +459,7 @@ impl MemoryPipeline { None } else { let emb = self.llm.embed(&dense_embedding_text(&episode_body)).await?; - let id = self.memory.store_with_metadata(&MemoryRecordInput { + let record = MemoryRecordInput { memory_id: None, namespace: "default".to_string(), layer: MemoryLayer::L1, @@ -368,20 +477,26 @@ impl MemoryPipeline { ], pinned: false, embedding: emb, - })?; + }; + let id = self + .write_db(move |memory| memory.store_with_metadata(&record)) + .await?; let episode_alias = format!("m{}", to_base36(id as u64)); for source_ref in &source_refs { - let _ = - self.memory - .add_reflink_edge(&episode_alias, source_ref, "derived_from"); + let source_ref = source_ref.clone(); + let episode_alias = episode_alias.clone(); + self.write_db(move |memory| { + memory.add_reflink_edge(&episode_alias, &source_ref, "derived_from") + }) + .await?; } Some(id) }; let episode_ref = episode_memory_id.map(|id| format!("m{}", to_base36(id as u64))); Ok(WriteTrace { - session_id: session.session_id, + session_id: session.session_id.clone(), turn_ids, turn_refs, episode_memory_id, @@ -395,7 +510,6 @@ impl MemoryPipeline { query: BenchQuery, budget: ContextBudget, ) -> Result { - self.memory.sync_reference_indexes()?; let exact_refs = extract_ref_codes(&query.text) .into_iter() .collect::>(); @@ -407,7 +521,7 @@ impl MemoryPipeline { let packet_token_budget = op_plan.packet_token_budget(budget.max_tokens); let mut candidates = Vec::new(); - candidates.extend(self.resolve_exact_refs(&exact_refs)?); + candidates.extend(self.resolve_exact_refs(&exact_refs).await?); if route.archival_allowed && self.profile.lexical { candidates.extend( self.search_sparse(&query.text, &routed_lanes, candidate_limit) @@ -427,25 +541,88 @@ impl MemoryPipeline { } if route.archival_allowed && self.profile.graph && budget.graph_depth > 0 { - let seed_refs = canonical_refs(&self.memory, &exact_refs)?; - candidates.extend(self.expand_graph(&seed_refs, candidate_limit)?); + candidates.extend(self.expand_graph(&exact_refs, candidate_limit).await?); } let ranked_limit = packet_limit.saturating_mul(4).max(candidate_limit); - let (candidates, stage_one) = if self.profile.session_first { + let (mut candidates, mut stage_one) = if self.profile.session_first { rank_session_first_candidates(candidates, ranked_limit, query.reference_time) } else { let candidates = rank_seed_candidates(candidates, ranked_limit); let stage_one = trace_from_ranked_candidates("source_interleave", &candidates); (candidates, stage_one) }; - let packet_plan = self.evidence_planner().build_packets_with_collect_mode( - &query.text, - &candidates, - packet_limit, - packet_token_budget, - op_plan.collect_all, - )?; + let mut packet_plan = self + .build_packets_async( + query.text.clone(), + candidates.clone(), + packet_limit, + packet_token_budget, + op_plan.collect_all, + ) + .await?; + if op_plan.collect_all && op_plan.max_passes > 1 && route.archival_allowed { + let followup_queries = collect_followup_queries(&query.text, &packet_plan.packets); + let seen_refs = candidates + .iter() + .map(|candidate| candidate.ref_id.clone()) + .collect::>(); + let mut followup = Vec::new(); + for followup_query in &followup_queries { + if self.profile.lexical { + followup.extend( + self.search_sparse(followup_query, &routed_lanes, candidate_limit) + .await?, + ); + } + if self.profile.dense { + followup.extend( + self.search_dense( + followup_query, + query.reference_time, + &routed_lanes, + candidate_limit, + ) + .await?, + ); + } + } + followup.retain(|candidate| !seen_refs.contains(&candidate.ref_id)); + if followup.is_empty() { + packet_plan.coverage.passes = 2; + packet_plan.coverage.stop_reason = "no_new_evidence".to_string(); + packet_plan.coverage.followup_queries = followup_queries; + } else { + candidates.extend(followup); + (candidates, stage_one) = if self.profile.session_first { + rank_session_first_candidates(candidates, ranked_limit, query.reference_time) + } else { + let ranked = rank_seed_candidates(candidates, ranked_limit); + let trace = trace_from_ranked_candidates("source_interleave", &ranked); + (ranked, trace) + }; + packet_plan = self + .build_packets_async( + query.text.clone(), + candidates.clone(), + packet_limit, + packet_token_budget, + true, + ) + .await?; + packet_plan.coverage.passes = 2; + packet_plan.coverage.followup_queries = followup_queries; + packet_plan.coverage.stop_reason = if packet_plan.coverage.gap_count == 0 + && packet_plan.coverage.selected_fact_groups >= 2 + { + "sufficient_after_gap_repair".to_string() + } else { + "budget_exhausted".to_string() + }; + } + } else if op_plan.status != "ok" { + packet_plan.coverage.stop_reason = "planner_degraded".to_string(); + } let (packets, packet_rerank) = self .rerank_packets_if_enabled(&query.text, packet_plan.packets) .await; @@ -480,14 +657,14 @@ impl MemoryPipeline { && retrieved.packet_omissions.is_empty() { fallback_packets = self - .evidence_planner() - .build_packets_with_collect_mode( - &query.text, - &retrieved.candidates, + .build_packets_async( + query.text.clone(), + retrieved.candidates.clone(), retrieved.op_plan.packet_limit(budget.top_k), retrieved.op_plan.packet_token_budget(budget.max_tokens), retrieved.op_plan.collect_all, - )? + ) + .await? .packets; fallback_packets.as_slice() } else { @@ -533,30 +710,34 @@ impl MemoryPipeline { query: BenchQuery, budget: ContextBudget, ) -> Result { - self.memory.sync_reference_indexes()?; let route = route_memory_query(&query.text, false); let op_plan = self.plan_query_op(&query).await; let lanes = route.lanes.clone(); let candidate_limit = op_plan.candidate_limit(budget.top_k); let packet_limit = op_plan.packet_limit(budget.top_k); let packet_token_budget = op_plan.packet_token_budget(budget.max_tokens); - let entries = if route.archival_allowed { - self.memory - .active_promptable_refs(&lanes, candidate_limit.saturating_mul(8).max(24))? + let candidates = if route.archival_allowed { + let read_lanes = lanes.clone(); + self.read_db(move |memory| { + memory + .active_promptable_refs(&read_lanes, candidate_limit.saturating_mul(8).max(24))? + .into_iter() + .map(|entry| ref_search_entry_to_atom(&memory, entry)) + .collect() + }) + .await? } else { vec![] }; - let candidates = entries - .into_iter() - .map(|entry| self.ref_search_entry_to_atom(entry)) - .collect::>>()?; - let packet_plan = self.evidence_planner().build_packets_with_collect_mode( - &query.text, - &candidates, - packet_limit, - packet_token_budget, - op_plan.collect_all, - )?; + let packet_plan = self + .build_packets_async( + query.text.clone(), + candidates.clone(), + packet_limit, + packet_token_budget, + op_plan.collect_all, + ) + .await?; let (packets, packet_rerank) = self .rerank_packets_if_enabled(&query.text, packet_plan.packets) .await; @@ -603,43 +784,49 @@ impl MemoryPipeline { }) } - fn resolve_exact_refs(&self, refs: &[String]) -> Result> { - canonical_refs(&self.memory, refs)? - .into_iter() - .filter_map(|ref_id| self.ref_entry(&ref_id, "exact", 0.0).transpose()) - .collect() + async fn resolve_exact_refs(&self, refs: &[String]) -> Result> { + let refs = refs.to_vec(); + self.read_db(move |memory| { + canonical_refs(&memory, &refs)? + .into_iter() + .filter_map(|ref_id| ref_entry(&memory, &ref_id, "exact", 0.0).transpose()) + .collect() + }) + .await } - fn search_lexical( + async fn search_lexical( &self, query: &str, lanes: &[MemoryLane], limit: usize, ) -> Result> { let channel_limit = limit.saturating_mul(4).max(limit).max(16); - let mut entries = self.memory.search_refs_fts(query, lanes, channel_limit)?; - if entries.len() < limit { - let mut seen = entries - .iter() - .map(|entry| entry.ref_id.clone()) - .collect::>(); - for entry in self - .memory - .search_refs_trigram(query, lanes, channel_limit)? - { - if seen.insert(entry.ref_id.clone()) { - entries.push(entry); - } - if entries.len() >= limit { - break; + let query = query.to_string(); + let lanes = lanes.to_vec(); + self.read_db(move |memory| { + let mut entries = memory.search_refs_fts(&query, &lanes, channel_limit)?; + if entries.len() < limit { + let mut seen = entries + .iter() + .map(|entry| entry.ref_id.clone()) + .collect::>(); + for entry in memory.search_refs_trigram(&query, &lanes, channel_limit)? { + if seen.insert(entry.ref_id.clone()) { + entries.push(entry); + } + if entries.len() >= limit { + break; + } } } - } - entries - .into_iter() - .take(limit) - .map(|entry| self.ref_search_entry_to_atom(entry)) - .collect() + entries + .into_iter() + .take(limit) + .map(|entry| ref_search_entry_to_atom(&memory, entry)) + .collect() + }) + .await } async fn search_sparse( @@ -650,15 +837,17 @@ impl MemoryPipeline { ) -> Result> { let sparse_weights = self.llm.embed_sparse(query).await?; if sparse_weights.is_empty() { - return self.search_lexical(query, lanes, limit); + return self.search_lexical(query, lanes, limit).await; } - let entries = self - .memory - .search_refs_sparse(&sparse_weights, lanes, limit)?; - entries - .into_iter() - .map(|entry| self.ref_search_entry_to_atom(entry)) - .collect() + let lanes = lanes.to_vec(); + self.read_db(move |memory| { + memory + .search_refs_sparse(&sparse_weights, &lanes, limit)? + .into_iter() + .map(|entry| ref_search_entry_to_atom(&memory, entry)) + .collect() + }) + .await } async fn search_dense( @@ -670,101 +859,72 @@ impl MemoryPipeline { ) -> Result> { let emb = self.llm.embed(query).await?; let embedding_model = self.llm.config.embedder.model.clone(); - for source in self - .memory - .promptable_refs_missing_embeddings(lanes, &embedding_model)? - { + let missing_lanes = lanes.to_vec(); + let missing_model = embedding_model.clone(); + let missing = self + .read_db(move |memory| { + memory.promptable_refs_missing_embeddings(&missing_lanes, &missing_model) + }) + .await?; + for source in missing { let item_embedding = self.llm.embed(&dense_embedding_text(&source.body)).await?; - self.memory.upsert_embedding_item( - &source.ref_id, - &embedding_model, - &source.body_hash, - &item_embedding, - )?; + let model = embedding_model.clone(); + self.write_db(move |memory| { + memory.upsert_embedding_item( + &source.ref_id, + &model, + &source.body_hash, + &item_embedding, + ) + }) + .await?; } - self.memory - .search_refs_dense(&emb, lanes, &embedding_model, limit)? - .into_iter() - .filter(|entry| entry.score < self.config.sim_threshold) - .map(|entry| self.ref_search_entry_to_atom(entry)) - .collect() + let lanes = lanes.to_vec(); + let threshold = self.config.sim_threshold; + self.read_db(move |memory| { + memory + .search_refs_dense(&emb, &lanes, &embedding_model, limit)? + .into_iter() + .filter(|entry| entry.score < threshold) + .map(|entry| ref_search_entry_to_atom(&memory, entry)) + .collect() + }) + .await } - fn expand_graph(&self, seed_refs: &[String], limit: usize) -> Result> { + async fn expand_graph(&self, seed_refs: &[String], limit: usize) -> Result> { if seed_refs.is_empty() { return Ok(vec![]); } - Ok(self - .memory - .expand_edges(seed_refs, limit)? - .into_iter() - .filter_map(|resolved| resolved_to_entry(resolved, "graph")) - .collect()) - } - - fn ref_entry(&self, ref_id: &str, source: &str, score: f32) -> Result> { - let Some(data) = self.memory.get_resolved_ref(ref_id)? else { - return Ok(None); - }; - let Some(body) = data.body else { - return Ok(None); - }; - let lane = self - .memory - .search_refs_fts(&body, &[], 1) - .ok() - .and_then(|entries| { - entries - .into_iter() - .find(|entry| entry.ref_id == ref_id) - .map(|entry| entry.lane) - }) - .unwrap_or(MemoryLane::Semantic); - Ok(Some(EvidenceAtom { - ref_id: ref_id.to_string(), - alias: None, - session_id: self.memory.ref_session_id(ref_id).ok().flatten(), - entity_type: EntityType::from_str(&data.entity_type), - lane, - source: RetrievalSource::from_str(source), - rank_in_source: 1, - score, - token_count: data.token_count.unwrap_or_else(|| body.chars().count() / 4), - timestamp: self.memory.ref_timestamp(ref_id).ok().flatten(), - body, - anchor_kind: if source == "exact" { - AnchorKind::ExplicitRef - } else { - AnchorKind::QueryMatch - }, - })) - } - - fn ref_search_entry_to_atom(&self, entry: RefSearchEntry) -> Result { - let timestamp = self - .memory - .ref_timestamp(&entry.ref_id) - .ok() - .flatten() - .or(entry.timestamp); - Ok(EvidenceAtom { - session_id: self.memory.ref_session_id(&entry.ref_id).ok().flatten(), - ref_id: entry.ref_id, - alias: entry.alias, - entity_type: EntityType::from_str(&entry.entity_type), - lane: entry.lane, - source: RetrievalSource::from_str(&entry.source), - rank_in_source: 1, - score: entry.score, - token_count: entry.token_count, - timestamp, - body: entry.body, - anchor_kind: AnchorKind::QueryMatch, + let seed_refs = seed_refs.to_vec(); + self.read_db(move |memory| { + Ok(memory + .expand_edges(&seed_refs, limit)? + .into_iter() + .filter_map(|resolved| resolved_to_entry(resolved, "graph")) + .collect()) }) + .await } - pub fn evidence_planner(&self) -> EvidencePlanner { - EvidencePlanner::new(self.memory.clone()) + async fn build_packets_async( + &self, + query: String, + candidates: Vec, + limit: usize, + packet_token_budget: usize, + collect_mode: bool, + ) -> Result { + self.read_db(move |memory| { + EvidencePlanner::new(memory).build_packets_with_collect_mode( + &query, + &candidates, + limit, + packet_token_budget, + collect_mode, + ) + }) + .await } async fn rerank_packets_if_enabled( @@ -851,6 +1011,70 @@ fn canonical_refs(memory: &MemoryStore, refs: &[String]) -> Result> Ok(mapped.into_iter().map(|(_, ref_id)| ref_id).collect()) } +fn ref_entry( + memory: &MemoryStore, + ref_id: &str, + source: &str, + score: f32, +) -> Result> { + let Some(data) = memory.get_resolved_ref(ref_id)? else { + return Ok(None); + }; + let Some(body) = data.body else { + return Ok(None); + }; + let lane = memory + .search_refs_fts(&body, &[], 1) + .ok() + .and_then(|entries| { + entries + .into_iter() + .find(|entry| entry.ref_id == ref_id) + .map(|entry| entry.lane) + }) + .unwrap_or(MemoryLane::Semantic); + Ok(Some(EvidenceAtom { + ref_id: ref_id.to_string(), + alias: None, + session_id: memory.ref_session_id(ref_id).ok().flatten(), + entity_type: EntityType::from_str(&data.entity_type), + lane, + source: RetrievalSource::from_str(source), + rank_in_source: 1, + score, + token_count: data.token_count.unwrap_or_else(|| body.chars().count() / 4), + timestamp: memory.ref_timestamp(ref_id).ok().flatten(), + body, + anchor_kind: if source == "exact" { + AnchorKind::ExplicitRef + } else { + AnchorKind::QueryMatch + }, + })) +} + +fn ref_search_entry_to_atom(memory: &MemoryStore, entry: RefSearchEntry) -> Result { + let timestamp = memory + .ref_timestamp(&entry.ref_id) + .ok() + .flatten() + .or(entry.timestamp); + Ok(EvidenceAtom { + session_id: memory.ref_session_id(&entry.ref_id).ok().flatten(), + ref_id: entry.ref_id, + alias: entry.alias, + entity_type: EntityType::from_str(&entry.entity_type), + lane: entry.lane, + source: RetrievalSource::from_str(&entry.source), + rank_in_source: 1, + score: entry.score, + token_count: entry.token_count, + timestamp, + body: entry.body, + anchor_kind: AnchorKind::QueryMatch, + }) +} + fn resolved_to_entry(resolved: ResolvedRef, source: &str) -> Option { match resolved { ResolvedRef::Active { @@ -1249,23 +1473,30 @@ fn trace_from_ranked_candidates(strategy: &str, candidates: &[EvidenceAtom]) -> } fn synthesize_planned_answer(plan: &OpPlan, packets: &[EvidencePacket]) -> Option { - let facts = packet_fact_rows(packets); + let extracted = packet_fact_rows(packets); + let facts = query_aligned_facts(&extracted); if facts.is_empty() { return None; } match plan.op { QueryOp::Lookup | QueryOp::AbstainOrFalsePremiseCheck => None, - QueryOp::AggregateCount => Some(distinct_fact_count(&facts).to_string()), + QueryOp::AggregateCount => { + consistent_target(&facts)?; + Some(distinct_fact_count(&facts).to_string()) + } QueryOp::AggregateSum => { + consistent_target(&facts)?; let values = numeric_fact_values(&facts); (!values.is_empty()).then(|| format_number(values.iter().sum::())) } QueryOp::AggregateAvg => { + consistent_target(&facts)?; let values = numeric_fact_values(&facts); (!values.is_empty()) .then(|| format_number(values.iter().sum::() / values.len().max(1) as f64)) } QueryOp::OrderOrRank => { + consistent_target(&facts)?; let ordered = ordered_fact_rows(&facts); (!ordered.is_empty()).then(|| { ordered @@ -1276,6 +1507,7 @@ fn synthesize_planned_answer(plan: &OpPlan, packets: &[EvidencePacket]) -> Optio }) } QueryOp::UpdateResolution => { + consistent_target(&facts)?; let ordered = ordered_fact_rows(&facts); match ordered.as_slice() { [] => None, @@ -1295,15 +1527,28 @@ fn synthesize_planned_answer(plan: &OpPlan, packets: &[EvidencePacket]) -> Optio if !has_personal_support(packets) { return Some("I don't know.".to_string()); } - facts - .iter() - .find(|fact| fact.session_id.is_some()) - .or_else(|| facts.first()) - .map(|fact| render_fact_brief(fact)) + // Generic extraction does not yet recover preference polarity or + // contradictions safely. Keep the rows as cited reader evidence, + // but do not turn the first episodic row into a recommendation. + None } } } +fn query_aligned_facts(facts: &[EvidenceFactRow]) -> Vec { + facts + .iter() + .filter(|fact| fact.target_key.is_some() && fact.extraction_confidence > 0.0) + .cloned() + .collect() +} + +fn consistent_target(facts: &[EvidenceFactRow]) -> Option<&str> { + let mut targets = facts.iter().filter_map(|fact| fact.target_key.as_deref()); + let first = targets.next()?; + targets.all(|target| target == first).then_some(first) +} + fn packet_fact_rows(packets: &[EvidencePacket]) -> Vec { let mut rows = Vec::new(); let mut seen = HashSet::new(); @@ -1322,14 +1567,41 @@ fn packet_fact_rows(packets: &[EvidencePacket]) -> Vec { rows } +fn collect_followup_queries(query: &str, packets: &[EvidencePacket]) -> Vec { + let mut seen = HashSet::new(); + let mut targets = packets + .iter() + .flat_map(|packet| packet.fact_rows.iter()) + .filter_map(|fact| fact.target_key.clone()) + .filter(|target| seen.insert(target.clone())) + .take(3) + .collect::>(); + if targets.is_empty() { + targets.push(query.to_string()); + } + targets +} + fn distinct_fact_count(facts: &[EvidenceFactRow]) -> usize { - facts + let turn_chunks = facts .iter() + .filter(|fact| fact.entity_type == "turn_chunk") + .collect::>(); + let selected = if turn_chunks.is_empty() { + facts.iter().collect::>() + } else { + turn_chunks + }; + selected + .into_iter() .map(|fact| { - fact.session_id - .as_ref() - .map(|session_id| format!("session:{session_id}")) - .unwrap_or_else(|| format!("ref:{}", fact.ref_id)) + format!( + "{}:{}:{}:{}", + fact.target_key.as_deref().unwrap_or(""), + fact.session_id.as_deref().unwrap_or(""), + fact.value_text.as_deref().unwrap_or(""), + fact.text + ) }) .collect::>() .len() @@ -1393,15 +1665,21 @@ fn aggregate_values_for_fact(fact: &EvidenceFactRow) -> Vec { fn ordered_fact_rows(facts: &[EvidenceFactRow]) -> Vec<&EvidenceFactRow> { let mut ordered = facts.iter().collect::>(); ordered.sort_by(|left, right| { - left.timestamp - .unwrap_or(i64::MAX) - .cmp(&right.timestamp.unwrap_or(i64::MAX)) + fact_sort_time(left) + .cmp(&fact_sort_time(right)) .then_with(|| left.ordinal.cmp(&right.ordinal)) .then_with(|| left.ref_id.cmp(&right.ref_id)) }); ordered } +fn fact_sort_time(fact: &EvidenceFactRow) -> i64 { + fact.valid_at + .or(fact.observed_at) + .or(fact.timestamp) + .unwrap_or(i64::MAX) +} + fn render_fact_brief(fact: &EvidenceFactRow) -> String { fact.value_text .clone() @@ -2607,6 +2885,74 @@ mod tests { assert_eq!(retrieval.op_plan.op, QueryOp::UpdateResolution); assert!(retrieval.op_plan.collect_all); assert_eq!(retrieval.op_plan.reason, "query_op_hint"); + assert_eq!(retrieval.coverage.passes, 2); + assert_eq!(retrieval.coverage.followup_queries, vec!["value"]); + assert!(matches!( + retrieval.coverage.stop_reason.as_str(), + "no_new_evidence" | "sufficient_after_gap_repair" | "budget_exhausted" + )); + Ok(()) + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn concurrent_ingest_and_retrieval_use_bounded_blocking_storage_jobs() -> Result<()> { + let tmp = NamedTempFile::new()?; + let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; + let llm = LlmClient::new(crate::models::ModelsConfig::default()); + let pipeline = MemoryPipeline::new(store, llm, MemoryConfig::default()) + .with_profile("raw-turns/fts-only"); + pipeline + .observe_session(BenchSession { + session_id: "seed-session".to_string(), + timestamp: Some(1), + turns: vec![BenchTurn { + role: "user".to_string(), + content: "the storage sentinel is blue quartz".to_string(), + timestamp: Some(1), + }], + }) + .await?; + let budget = ContextBudget { + max_tokens: 800, + top_k: 3, + graph_depth: 0, + }; + let query = BenchQuery { + query_id: "q-concurrent".to_string(), + text: "what color is the storage sentinel?".to_string(), + reference_time: Some(3), + op_hint: None, + }; + let ingest_pipeline = pipeline.clone(); + let retrieval_a = pipeline.clone(); + let retrieval_b = pipeline.clone(); + + let (ingest, first, second) = tokio::time::timeout(Duration::from_secs(5), async move { + tokio::join!( + ingest_pipeline.observe_session(BenchSession { + session_id: "concurrent-session".to_string(), + timestamp: Some(2), + turns: vec![BenchTurn { + role: "user".to_string(), + content: "a concurrent write stays ordered".to_string(), + timestamp: Some(2), + }], + }), + retrieval_a.retrieve_evidence(query.clone(), budget), + retrieval_b.retrieve_evidence(query, budget), + ) + }) + .await?; + + assert!(ingest.is_ok()); + assert!(first? + .candidates + .iter() + .any(|candidate| candidate.body.contains("blue quartz"))); + assert!(second? + .candidates + .iter() + .any(|candidate| candidate.body.contains("blue quartz"))); Ok(()) } @@ -2903,6 +3249,84 @@ mod tests { assert_eq!(answer, None); } + #[test] + fn planned_synthesis_rejects_same_unit_values_for_different_targets() { + let mut books = fact_packet( + "pkt_books", + "ref_books", + Some("s1"), + Some(10), + Some("$40"), + Some(40.0), + "spent $40 on books", + MemoryLane::Episodic, + ); + books.fact_rows[0].target_key = Some("books".to_string()); + books.fact_rows[0].values[0].kind = EvidenceValueKind::Currency; + books.fact_rows[0].values[0].unit = Some("$".to_string()); + let mut rent = fact_packet( + "pkt_rent", + "ref_rent", + Some("s2"), + Some(20), + Some("$900"), + Some(900.0), + "paid $900 rent", + MemoryLane::Episodic, + ); + rent.fact_rows[0].target_key = Some("rent".to_string()); + rent.fact_rows[0].values[0].kind = EvidenceValueKind::Currency; + rent.fact_rows[0].values[0].unit = Some("$".to_string()); + + let answer = synthesize_planned_answer( + &OpPlan::for_op(QueryOp::AggregateSum, "test"), + &[books, rent], + ); + + assert_eq!(answer, None); + } + + #[test] + fn collect_followup_queries_use_query_aligned_fact_targets() { + let mut books = fact_packet( + "pkt_books", + "ref_books", + Some("s1"), + Some(10), + Some("3"), + Some(3.0), + "read 3 books", + MemoryLane::Episodic, + ); + books.fact_rows[0].target_key = Some("books".to_string()); + + assert_eq!( + collect_followup_queries("how many books", &[books]), + vec!["books"] + ); + } + + #[test] + fn preference_synthesis_defers_when_polarity_is_not_structured() { + let packet = fact_packet( + "pkt_preference", + "ref_preference", + Some("s1"), + Some(10), + Some("jazz"), + None, + "dawn likes jazz but dislikes crowds", + MemoryLane::Profile, + ); + + let answer = synthesize_planned_answer( + &OpPlan::for_op(QueryOp::PreferenceRecommendation, "test"), + &[packet], + ); + + assert_eq!(answer, None); + } + #[test] fn planned_synthesis_orders_fact_rows_by_timestamp() { let packets = vec![ @@ -2976,6 +3400,43 @@ mod tests { assert!(answer.contains("latest: 26:30")); } + #[test] + fn update_resolution_orders_retroactive_reports_by_valid_time() { + let mut older_event_reported_late = fact_packet( + "pkt_old_event", + "ref_old_event", + Some("s2"), + Some(200), + Some("27:45"), + None, + "the 2019 personal best was 27:45", + MemoryLane::Episodic, + ); + older_event_reported_late.fact_rows[0].observed_at = Some(200); + older_event_reported_late.fact_rows[0].valid_at = Some(100); + let mut newer_event_reported_early = fact_packet( + "pkt_new_event", + "ref_new_event", + Some("s1"), + Some(150), + Some("26:30"), + None, + "the 2020 personal best was 26:30", + MemoryLane::Episodic, + ); + newer_event_reported_early.fact_rows[0].observed_at = Some(150); + newer_event_reported_early.fact_rows[0].valid_at = Some(300); + + let answer = synthesize_planned_answer( + &OpPlan::for_op(QueryOp::UpdateResolution, "test"), + &[newer_event_reported_early, older_event_reported_late], + ) + .unwrap(); + + assert!(answer.contains("previous: 27:45")); + assert!(answer.contains("latest: 26:30")); + } + #[test] fn preference_synthesis_abstains_without_personal_support() { let packets = vec![fact_packet( @@ -3187,7 +3648,14 @@ mod tests { entity_type: "turn_chunk".to_string(), session_id: session_id.map(str::to_string), timestamp, + observed_at: timestamp, + valid_at: None, + valid_until: None, ordinal: 0, + target_key: Some("target".to_string()), + polarity: "unspecified".to_string(), + extraction_confidence: 1.0, + extraction_method: "test_fixture".to_string(), value_text: value_text.map(str::to_string), value_number, values: value_text @@ -3209,6 +3677,8 @@ mod tests { ref_id: ref_id.to_string(), session_id: session_id.map(str::to_string), timestamp, + observed_at: Some(timestamp), + valid_at: None, ordinal: 0, text: text.to_string(), }] diff --git a/klbr-core/src/planner.rs b/klbr-core/src/planner.rs index b5297ac..d74ca6b 100644 --- a/klbr-core/src/planner.rs +++ b/klbr-core/src/planner.rs @@ -34,6 +34,7 @@ pub struct OpPlan { pub collect_all: bool, pub require_personal_support: bool, pub max_passes: usize, + pub status: String, pub reason: String, } @@ -44,6 +45,7 @@ impl Default for OpPlan { collect_all: false, require_personal_support: false, max_passes: 1, + status: "unavailable".to_string(), reason: "planner_unavailable_default_lookup".to_string(), } } @@ -65,6 +67,7 @@ impl OpPlan { collect_all, require_personal_support, max_passes: if collect_all { 2 } else { 1 }, + status: "ok".to_string(), reason: reason.into(), } } @@ -185,6 +188,7 @@ mod tests { assert_eq!(plan.op, QueryOp::Lookup); assert!(!plan.collect_all); + assert_eq!(plan.status, "unavailable"); assert_eq!(plan.candidate_limit(5), 5); assert_eq!(plan.packet_limit(5), 5); } @@ -194,6 +198,7 @@ mod tests { let plan = OpPlan::for_op(QueryOp::AggregateCount, "model_structured"); assert!(plan.collect_all); + assert_eq!(plan.status, "ok"); assert_eq!(plan.candidate_limit(5), 15); assert_eq!(plan.packet_limit(5), 10); }