From e9871e9e442e9aff56e3d08352be82cc0ae795a1 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Mon, 29 Jun 2026 01:21:39 +0300 Subject: [PATCH] remove hardcoded memory query cues --- .beads/issues.jsonl | 2 +- docs/long-term-memory-arch.md | 24 +- docs/memory-implementation-status.md | 18 +- klbr-bench/src/main.rs | 25 +- klbr-core/src/agent.rs | 8 +- klbr-core/src/evidence.rs | 80 +--- klbr-core/src/pipeline.rs | 604 ++++----------------------- klbr-core/src/planner.rs | 319 ++++---------- 8 files changed, 191 insertions(+), 889 deletions(-) diff --git a/.beads/issues.jsonl b/.beads/issues.jsonl index 2b6b202..6fd72bd 100644 --- a/.beads/issues.jsonl +++ b/.beads/issues.jsonl @@ -1,4 +1,4 @@ -{"_type":"issue","id":"klbr-u0q","title":"Implement session-first lexical-contained memory retrieval","description":"Build the next retrieval architecture from docs/long-term-memory-arch.md: stage-one session/event candidate generation, lexical containment, and language-agnostic operation planning experiments without adding english cue-word hacks.","design":"Keep EvidencePacket as the rank object. Use exact refs, dense episode/session cards, learned sparse or fts side channels, and graph expansion from strong seeds only. Keep QueryOp shape but replace cue-list policy with a structured classifier plus deterministic fallback.","acceptance_criteria":"A bench profile can retrieve from session/episode cards first and use chunk fts/sparse hits as packet enrichment; lexical-only policy paths are removed or contained behind candidate channels; traces expose candidate session recall and packet/rendered evidence metrics; multilingual fixtures cover non-English cue-free retrieval cases.","status":"open","priority":1,"issue_type":"feature","owner":"90008@klbr.net","created_at":"2026-06-28T22:06:32Z","created_by":"dawn","updated_at":"2026-06-28T22:06:32Z","dependency_count":0,"dependent_count":0,"comment_count":0} +{"_type":"issue","id":"klbr-u0q","title":"Implement session-first lexical-contained memory retrieval","description":"Build the next retrieval architecture from docs/long-term-memory-arch.md: stage-one session/event candidate generation, lexical containment, and language-agnostic operation planning experiments without adding english cue-word hacks.","design":"Keep EvidencePacket as the rank object. Use exact refs, dense episode/session cards, learned sparse or fts side channels, and graph expansion from strong seeds only. Keep QueryOp shape but replace cue-list policy with a structured classifier plus deterministic fallback.","acceptance_criteria":"A bench profile can retrieve from session/episode cards first and use chunk fts/sparse hits as packet enrichment; lexical-only policy paths are removed or contained behind candidate channels; traces expose candidate session recall and packet/rendered evidence metrics; multilingual fixtures cover non-English cue-free retrieval cases.","notes":"2026-06-29 partial: removed hardcoded english cue-list planner, english lane routing, english negation/entity-boundary packet filters, and wh/pronoun neighbor-expansion heuristics. OpPlan now comes from structured model JSON when an llm endpoint is configured, otherwise conservative lookup. Core and bench tests pass.","status":"in_progress","priority":1,"issue_type":"feature","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-28T22:06:32Z","created_by":"dawn","updated_at":"2026-06-28T22:20:45Z","started_at":"2026-06-28T22:10:12Z","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-jwn","title":"Organize evolving memory architecture docs","description":"Make the memory architecture docs distinguish current truth, implementation status, proposals, research inputs, and archived rationale so future agents do not treat older reports as canonical.","acceptance_criteria":"Docs have a clear index and status taxonomy; current memory architecture direction points at docs/long-term-memory-arch.md; older research reports are marked as historical or supporting; AGENTS.md routes future agents through the doc index.","status":"closed","priority":1,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-28T22:01:20Z","created_by":"dawn","updated_at":"2026-06-28T22:07:09Z","started_at":"2026-06-28T22:01:23Z","closed_at":"2026-06-28T22:07:09Z","close_reason":"Completed docs index, current architecture review rewrite, status banners, AGENTS routing, and follow-up implementation issue klbr-u0q.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-h4l","title":"Cache bench embeddings locally and run 25-sample QA bench","description":"Make LongMemEval bench embeddings use a project-local gitignored sqlite cache so reruns do not pay for the same embeddings twice. Then run a 25-sample stratified QA bench, inspect failures, try bounded non-overfit tweaks if evidence supports them, and prepare a report packet if results remain weak or tuning would be overfit.","design":"Prefer a global project cache path such as benchmarks/cache/*.db over temp/run-local caches. Keep benchmark polling sparse: start the run, wait for artifacts or process completion, then inspect once.","acceptance_criteria":"Embedding cache lives under the project and is ignored by git; cache hits are reused across bench runs; a 25-sample bench artifact exists with failure analysis; any code changes are tested, committed, and pushed.","notes":"Implemented project-local embedding cache at benchmarks/cache/embeddings.db and wired LongMemEval bench LlmClient construction through it. Ran stratified 25-sample QA on seed 4937553249516211047: current-code run benchmarks/runs/longmemeval-s/klbr-full/2026-06-28_210327.693783Z scored official accuracy 0.6800. Report packet written to report_packet.md in that run dir. High-budget targeted reruns fixed only dd2973ad and a1cc6108; remaining failures point to fact/timeline synthesis rather than safe cue tuning.","status":"closed","priority":1,"issue_type":"task","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-28T20:08:32Z","created_by":"dawn","updated_at":"2026-06-28T21:21:27Z","started_at":"2026-06-28T20:08:35Z","closed_at":"2026-06-28T21:21:27Z","close_reason":"Implemented project-local dense/sparse embedding cache, ran 25-sample QA bench, inspected failures, tried bounded planner/high-budget experiments, and wrote report packet.","dependency_count":0,"dependent_count":0,"comment_count":0} {"_type":"issue","id":"klbr-f0q.1","title":"Add operation-aware query planning","description":"Introduce a small OpPlan/QueryOp classifier for memory QA queries: lookup, aggregate_count, aggregate_sum, aggregate_avg, order_or_rank, update_resolution, preference_recommendation, and abstain_or_false_premise_check. Use it to switch only high-value classes away from plain lookup behavior.","design":"Keep the label set intentionally small and rule-based first; avoid growing English lexical hacks into the main ranking decision boundary.","acceptance_criteria":"Production memory retrieval traces expose the selected query op and plan; existing lookup behavior remains unchanged by default; aggregate/order/update/preference queries can be detected deterministically in tests.","status":"closed","priority":1,"issue_type":"feature","assignee":"dawn","owner":"90008@klbr.net","created_at":"2026-06-28T19:44:46Z","created_by":"dawn","updated_at":"2026-06-28T19:54:15Z","started_at":"2026-06-28T19:45:04Z","closed_at":"2026-06-28T19:54:15Z","close_reason":"Implemented OpPlan/QueryOp classification, serialized op_plan into retrieval traces, preserved lookup defaults, added deterministic planner and pipeline trace tests.","dependencies":[{"issue_id":"klbr-f0q.1","depends_on_id":"klbr-f0q","type":"parent-child","created_at":"2026-06-28T22:44:46Z","created_by":"dawn","metadata":"{}"}],"dependency_count":0,"dependent_count":1,"comment_count":0} diff --git a/docs/long-term-memory-arch.md b/docs/long-term-memory-arch.md index df513a8..79e3f9d 100644 --- a/docs/long-term-memory-arch.md +++ b/docs/long-term-memory-arch.md @@ -16,11 +16,10 @@ decides what memory means. the important correction is scope. klbr already has many pieces from the report: canonical refs, promptable text, dense and sparse embeddings, fts fallback, -episode cards, evidence packets, packet rrf, packet reranking, and a small -operation planner. the remaining work is not "add packets." it is to make the -first retrieval horizon more session/event based, remove hard-coded english cue -policy, and make aggregate/update/order synthesis less dependent on reader -freestyle. +episode cards, evidence packets, packet rrf, packet reranking, and a structured +operation planner interface. the remaining work is not "add packets." it is to +make the first retrieval horizon more session/event based and make +aggregate/update/order synthesis less dependent on reader freestyle. ## current diagnosis @@ -87,11 +86,10 @@ not allowed as the main answer: - additive lexical bonuses that suppress strong dense/session evidence. - "fix this miss" keyword patches. -current code still has a cue-list operation planner in `klbr-core/src/planner.rs`. -that planner is useful scaffolding, but it is technical debt. the replacement -should emit the same small `QueryOp` shape using a language-agnostic classifier -or a structured model call, then keep deterministic fallback behavior for -offline tests. +current code no longer has the cue-list operation planner or english query +constraint filters. `klbr-core/src/planner.rs` keeps the `QueryOp` interface, +parses structured model replies, and falls back to conservative `lookup` when a +planner model is unavailable. ## operation planning @@ -162,9 +160,9 @@ signals, and estimated tokens. decide the evidence horizon. chunk fts/sparse hits should enrich those session packets. -2. replace cue-list planner policy. - keep `QueryOp`, but move classification behind a structured classifier with - deterministic tests and a conservative fallback. +2. harden structured operation planning. + keep `QueryOp`, improve structured planner reliability, and preserve the + conservative offline fallback. do not reintroduce query-word lists. 3. finish collect-mode fact coverage. track fact groups for aggregate/order/update questions. expose sufficiency diff --git a/docs/memory-implementation-status.md b/docs/memory-implementation-status.md index f54f4c7..4434973 100644 --- a/docs/memory-implementation-status.md +++ b/docs/memory-implementation-status.md @@ -35,10 +35,6 @@ going next. - `EvidencePlanner` builds `exact_ref`, `turn_window`, `episode_bridge`, and `graph_bridge` packets; packet fusion replaces raw same-session candidate dedupe with same-session ref-overlap suppression. -- retrieval applies lightweight query constraints after packet planning: - negated query terms suppress fuzzy packets that mention the excluded target, - and narrow entity-boundary conflicts block sibling-entity bleed such as - film-vs-camera recall. exact-ref packets bypass these filters. - episodic notes written by the pipeline are now `episode_event_card` artifacts with structured event metadata, source-grounded timeline snippets, stated facts/decisions, attachment markers, and source refs. @@ -51,6 +47,13 @@ going next. packet metrics, and optional packet-rerank pre/post order and scores. - optional packet reranking scores bounded packet documents, protects exact-ref packets, and falls back with traceable status/message when the reranker fails. +- operation planning keeps the `QueryOp`/`OpPlan` interface, but no longer uses + hard-coded english cue lists. runtime planning calls a structured model + classifier when an llm endpoint is configured and otherwise falls back to + conservative `lookup`. +- query-time english negation/entity-boundary packet filters were removed. + lexical infrastructure remains candidate generation through sparse embeddings + or sqlite fts, not handwritten query policy. - `klbr-bench -- run` now uses the production pipeline instead of manually calling low-level retrieval: - longmemeval-s/m style data - flexible locomo adapter with fixture coverage and protocol-note artifacts @@ -62,10 +65,9 @@ going next. ## open direction -- `klbr-u0q`: implement session/event-first retrieval and lexical containment. -- replace the cue-list operation planner in `klbr-core/src/planner.rs` with a - language-agnostic structured classifier while keeping deterministic fallback - tests. +- `klbr-u0q`: implement session/event-first retrieval and finish lexical + containment beyond the first cue-list removal. +- harden the structured planner path and add model-backed planner fixtures. - add a session/event-first retrieval profile where dense episode cards seed sessions and chunk hits enrich evidence packets. - finish collect-mode fact-group sufficiency and gap diagnosis for aggregate, diff --git a/klbr-bench/src/main.rs b/klbr-bench/src/main.rs index a90b574..b400137 100644 --- a/klbr-bench/src/main.rs +++ b/klbr-bench/src/main.rs @@ -680,12 +680,12 @@ async fn run_continuous_loop_queries( require_candidates: true, }, ContinuousLoopQueryCase { - query_id: "negated-project-switch".to_string(), - text: "working on the dotfiles syncer today, not touching klbr".to_string(), - expected_terms: Vec::new(), - forbidden_terms: vec!["rust cli harness".to_string()], - expected_omission: Some("negated_query_term:klbr".to_string()), - expect_packets: false, + query_id: "project-syncer".to_string(), + text: "dotfiles syncer rust cli harness".to_string(), + expected_terms: vec!["rust cli harness".to_string()], + forbidden_terms: Vec::new(), + expected_omission: None, + expect_packets: true, require_candidates: true, }, ]; @@ -786,9 +786,9 @@ fn continuous_loop_checks( }) .collect::>>()?; let query_pass_count = queries.iter().filter(|query| query.passed).count(); - let negated_ok = queries + let project_ok = queries .iter() - .find(|query| query.query_id == "negated-project-switch") + .find(|query| query.query_id == "project-syncer") .is_some_and(|query| query.passed); let exact_ok = queries .iter() @@ -836,9 +836,10 @@ fn continuous_loop_checks( details: format!("{} / {} queries passed", query_pass_count, queries.len()), }, ContinuousLoopCheck { - name: "negated_passive_recall_suppressed".to_string(), - passed: negated_ok, - details: "dotfiles query keeps klbr packet out of assembled context".to_string(), + name: "project_passive_recall".to_string(), + passed: project_ok, + details: "project query retrieves the matching packet without query-word suppression" + .to_string(), }, ContinuousLoopCheck { name: "exact_ref_resolution".to_string(), @@ -6488,7 +6489,7 @@ mod tests { assert!(report .queries .iter() - .any(|query| query.query_id == "negated-project-switch" && query.passed)); + .any(|query| query.query_id == "project-syncer" && query.passed)); assert!(report .queries .iter() diff --git a/klbr-core/src/agent.rs b/klbr-core/src/agent.rs index 67048d4..e2b65b4 100644 --- a/klbr-core/src/agent.rs +++ b/klbr-core/src/agent.rs @@ -2203,8 +2203,7 @@ mod tests { } #[tokio::test] - async fn runtime_packet_recall_skips_archival_search_for_discourse_local_prompts() -> Result<()> - { + async fn runtime_packet_recall_does_not_use_discourse_cue_filter() -> Result<()> { let tmp = NamedTempFile::new()?; let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; let llm = LlmClient::new(ModelsConfig::default()); @@ -2239,7 +2238,10 @@ mod tests { ) .await?; - assert!(packets.is_empty()); + assert!(packets + .iter() + .flat_map(|packet| &packet.bodies) + .any(|body| body.contains("old archival phrase"))); Ok(()) } diff --git a/klbr-core/src/evidence.rs b/klbr-core/src/evidence.rs index 972c72a..0a75ee7 100644 --- a/klbr-core/src/evidence.rs +++ b/klbr-core/src/evidence.rs @@ -186,51 +186,6 @@ pub struct EvidenceOmittedRef { pub reason: String, } -fn query_is_wh_like(query: &str) -> bool { - let lower = query.to_lowercase(); - ["what", "who", "where", "why", "how", "when", "which"] - .iter() - .any(|word| lower.starts_with(word) || lower.contains(&format!(" {word}"))) -} - -fn looks_incomplete(body: &str) -> bool { - let trimmed = body.trim(); - trimmed.len() < 30 - || trimmed.ends_with(':') - || trimmed.ends_with(',') - || ["he", "she", "they", "it", "that", "this"] - .iter() - .any(|pronoun| trimmed.to_lowercase().contains(pronoun)) -} - -fn neighbor_gain(neighbor_body: &str, anchor_body: &str, query: &str) -> f32 { - let lower_neighbor = neighbor_body.to_lowercase(); - let lower_query = query.to_lowercase(); - - let mut fts_score = 0.0; - for word in lower_query.split_whitespace() { - if word.len() >= 4 && lower_neighbor.contains(word) { - fts_score += 1.0; - } - } - - let mut entity_score = 0.0; - for word in anchor_body.split_whitespace() { - if let Some(first_char) = word.chars().next() { - if first_char.is_uppercase() && word.len() >= 3 { - let clean_word = word - .trim_matches(|c: char| !c.is_alphabetic()) - .to_lowercase(); - if !clean_word.is_empty() && lower_neighbor.contains(&clean_word) { - entity_score += 1.5; - } - } - } - } - - fts_score + entity_score -} - #[derive(Clone)] pub struct EvidencePlanner { memory: MemoryStore, @@ -243,17 +198,17 @@ impl EvidencePlanner { pub fn build_packets( &self, - query: &str, + _query: &str, candidates: &[EvidenceAtom], limit: usize, packet_token_budget: usize, ) -> Result { - self.build_packets_with_collect_mode(query, candidates, limit, packet_token_budget, false) + self.build_packets_with_collect_mode(_query, candidates, limit, packet_token_budget, false) } pub fn build_packets_with_collect_mode( &self, - query: &str, + _query: &str, candidates: &[EvidenceAtom], limit: usize, packet_token_budget: usize, @@ -265,7 +220,7 @@ impl EvidencePlanner { for (index, candidate) in candidates.iter().enumerate() { let expanded = - self.expand_candidate_packet(query, candidate, index + 1, packet_token_budget)?; + self.expand_candidate_packet(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() { @@ -290,7 +245,6 @@ impl EvidencePlanner { fn expand_candidate_packet( &self, - query: &str, candidate: &EvidenceAtom, rank: usize, packet_token_budget: usize, @@ -313,30 +267,8 @@ impl EvidencePlanner { let expansion_refs = match packet_kind { EvidencePacketKind::TurnWindow => { - let tight_refs = self - .memory - .turn_window_ref_ids(&candidate.ref_id, 1, 1, 10)?; - let mut final_refs = tight_refs; - if query_is_wh_like(query) || looks_incomplete(&candidate.body) { - let wider_refs = - self.memory - .turn_window_ref_ids(&candidate.ref_id, 2, 2, 18)?; - let extra_refs: Vec = wider_refs - .into_iter() - .filter(|r| !final_refs.contains(r)) - .collect(); - if !extra_refs.is_empty() { - let extra_entries = self - .memory - .promptable_entries_for_refs(&extra_refs, "turn_window")?; - for entry in extra_entries { - if neighbor_gain(&entry.body, &candidate.body, query) >= 1.0 { - final_refs.push(entry.ref_id); - } - } - } - } - final_refs + self.memory + .turn_window_ref_ids(&candidate.ref_id, 1, 1, 10)? } EvidencePacketKind::ExactRef | EvidencePacketKind::EpisodeBridge diff --git a/klbr-core/src/pipeline.rs b/klbr-core/src/pipeline.rs index eb199aa..84d40ea 100644 --- a/klbr-core/src/pipeline.rs +++ b/klbr-core/src/pipeline.rs @@ -1,4 +1,4 @@ -use std::collections::{HashMap, HashSet, VecDeque}; +use std::collections::{HashMap, VecDeque}; use std::time::Duration; use anyhow::Result; @@ -16,7 +16,7 @@ use crate::{ }, models::{LlmClient, Message, RerankResult}, mvp::{MemoryLayer, MemoryRecordInput, MemoryStatus}, - planner::{classify_query_op, OpPlan}, + planner::{parse_structured_op_plan, OpPlan, STRUCTURED_QUERY_PLANNER_PROMPT}, }; pub use crate::evidence::EvidenceAtom; @@ -80,7 +80,7 @@ pub struct WriteTrace { pub source_refs: Vec, } -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub struct LaneRoute { pub lanes: Vec, pub reason: String, @@ -91,11 +91,13 @@ impl Default for LaneRoute { fn default() -> Self { Self { lanes: vec![ - MemoryLane::Semantic, + MemoryLane::Live, MemoryLane::Episodic, + MemoryLane::Semantic, MemoryLane::Profile, + MemoryLane::Procedural, ], - reason: "semantic_default".to_string(), + reason: "all_lanes_no_query_cues".to_string(), archival_allowed: true, } } @@ -207,6 +209,41 @@ impl MemoryPipeline { self } + async fn plan_query_op(&self, query: &str) -> OpPlan { + if self.llm.config.llm.url.trim().is_empty() || self.llm.config.llm.model.trim().is_empty() + { + return OpPlan::default(); + } + + let messages = [ + Message::system(STRUCTURED_QUERY_PLANNER_PROMPT), + Message::user(query), + ]; + let completion = tokio::time::timeout( + Duration::from_millis(self.config.rerank_timeout_ms), + self.llm.complete(&messages), + ) + .await; + + match completion { + Ok(Ok((reply, _))) => parse_structured_op_plan(&reply).unwrap_or_else(|| OpPlan { + reason: "planner_parse_failed_default_lookup".to_string(), + ..OpPlan::default() + }), + Ok(Err(error)) => { + tracing::warn!("structured query planner failed: {error:?}"); + OpPlan { + reason: "planner_request_failed_default_lookup".to_string(), + ..OpPlan::default() + } + } + Err(_) => OpPlan { + reason: "planner_timeout_default_lookup".to_string(), + ..OpPlan::default() + }, + } + } + pub async fn reset(&self, _run: BenchRun) -> Result<()> { self.memory.reset()?; self.memory.sync_reference_indexes()?; @@ -321,7 +358,7 @@ impl MemoryPipeline { .collect::>(); let route = route_memory_query(&query.text, !exact_refs.is_empty()); let routed_lanes = route.lanes.clone(); - let op_plan = classify_query_op(&query.text); + let op_plan = self.plan_query_op(&query.text).await; 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); @@ -355,17 +392,15 @@ impl MemoryPipeline { candidates, packet_limit.saturating_mul(4).max(candidate_limit), ); - let mut packet_plan = self.evidence_planner().build_packets_with_collect_mode( + let packet_plan = self.evidence_planner().build_packets_with_collect_mode( &query.text, &candidates, packet_limit, packet_token_budget, op_plan.collect_all, )?; - let constrained = constrain_packets_for_query(&query.text, packet_plan.packets); - packet_plan.omissions.extend(constrained.omissions); let (packets, packet_rerank) = self - .rerank_packets_if_enabled(&query.text, constrained.packets) + .rerank_packets_if_enabled(&query.text, packet_plan.packets) .await; Ok(PipelineRetrievalTrace { query_id: query.query_id, @@ -405,8 +440,6 @@ impl MemoryPipeline { retrieved.op_plan.collect_all, )? .packets; - let constrained = constrain_packets_for_query(&query.text, fallback_packets); - fallback_packets = constrained.packets; fallback_packets.as_slice() } else { retrieved.packets.as_slice() @@ -452,7 +485,7 @@ impl MemoryPipeline { ) -> Result { self.memory.sync_reference_indexes()?; let route = route_memory_query(&query.text, false); - let op_plan = classify_query_op(&query.text); + let op_plan = self.plan_query_op(&query.text).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); @@ -467,17 +500,15 @@ impl MemoryPipeline { .into_iter() .map(|entry| self.ref_search_entry_to_atom(entry)) .collect::>>()?; - let mut packet_plan = self.evidence_planner().build_packets_with_collect_mode( + let packet_plan = self.evidence_planner().build_packets_with_collect_mode( &query.text, &candidates, packet_limit, packet_token_budget, op_plan.collect_all, )?; - let constrained = constrain_packets_for_query(&query.text, packet_plan.packets); - packet_plan.omissions.extend(constrained.omissions); let (packets, packet_rerank) = self - .rerank_packets_if_enabled(&query.text, constrained.packets) + .rerank_packets_if_enabled(&query.text, packet_plan.packets) .await; let retrieved = PipelineRetrievalTrace { query_id: query.query_id.clone(), @@ -1015,288 +1046,6 @@ fn drain_source( } } -struct PacketConstraintOutcome { - packets: Vec, - omissions: Vec, -} - -#[derive(Debug, Clone)] -struct QueryEvidenceConstraints { - query_terms: HashSet, - negated_terms: Vec, -} - -impl QueryEvidenceConstraints { - fn from_query(query: &str) -> Self { - let negated_terms = negated_constraint_terms(query); - let negated = negated_terms.iter().cloned().collect::>(); - let query_terms = normalized_constraint_terms(query) - .into_iter() - .filter(|term| !negated.contains(term)) - .collect::>(); - Self { - query_terms, - negated_terms, - } - } - - fn is_empty(&self) -> bool { - self.negated_terms.is_empty() - && entity_boundary_classes_for_query(&self.query_terms).is_empty() - } - - fn rejection_reason(&self, packet: &EvidencePacket) -> Option { - if packet.packet_kind == EvidencePacketKind::ExactRef || self.is_empty() { - return None; - } - - let packet_terms = packet_terms(packet); - if let Some(term) = self - .negated_terms - .iter() - .find(|term| packet_terms.contains(term.as_str())) - { - return Some(format!("negated_query_term:{term}")); - } - - entity_boundary_mismatch(&self.query_terms, &packet_terms) - } -} - -fn constrain_packets_for_query( - query: &str, - packets: Vec, -) -> PacketConstraintOutcome { - let constraints = QueryEvidenceConstraints::from_query(query); - if constraints.is_empty() { - return PacketConstraintOutcome { - packets, - omissions: vec![], - }; - } - - let mut kept = Vec::new(); - let mut omissions = Vec::new(); - for packet in packets { - if let Some(reason) = constraints.rejection_reason(&packet) { - omissions.extend(packet.refs.iter().map(|ref_id| EvidenceOmittedRef { - packet_id: packet.packet_id.clone(), - ref_id: ref_id.clone(), - reason: reason.clone(), - })); - } else { - kept.push(packet); - } - } - - PacketConstraintOutcome { - packets: kept, - omissions, - } -} - -fn negated_constraint_terms(query: &str) -> Vec { - const NEGATION_MARKERS: &[&str] = &[ - "not touching", - "not about", - "not working on", - "not using", - "not related to", - "unrelated to", - "without", - "excluding", - "exclude", - ]; - - let lower = query - .chars() - .flat_map(char::to_lowercase) - .collect::(); - let mut out = Vec::new(); - for marker in NEGATION_MARKERS { - let mut cursor = 0usize; - while let Some(offset) = lower[cursor..].find(marker) { - let after_marker = cursor + offset + marker.len(); - for term in constraint_terms_until_boundary(&lower[after_marker..], 4) { - if !out.contains(&term) { - out.push(term); - } - } - cursor = after_marker; - } - } - out -} - -fn constraint_terms_until_boundary(value: &str, limit: usize) -> Vec { - let local = value - .split(|ch| [',', '.', ';', ':', '!', '?', '\n'].contains(&ch)) - .next() - .unwrap_or_default(); - normalized_constraint_terms(local) - .into_iter() - .take(limit) - .collect() -} - -fn packet_terms(packet: &EvidencePacket) -> HashSet { - packet - .bodies - .iter() - .flat_map(|body| normalized_constraint_terms(body)) - .collect() -} - -fn normalized_constraint_terms(value: &str) -> Vec { - value - .split(|ch: char| !ch.is_alphanumeric() && ch != '_') - .filter_map(|term| { - let term = term - .chars() - .flat_map(char::to_lowercase) - .collect::(); - let len = term.chars().count(); - (!constraint_stopword(&term) && len >= 3).then_some(term) - }) - .collect() -} - -fn constraint_stopword(term: &str) -> bool { - matches!( - term, - "the" - | "and" - | "for" - | "are" - | "but" - | "not" - | "you" - | "your" - | "our" - | "was" - | "were" - | "has" - | "had" - | "his" - | "her" - | "she" - | "him" - | "its" - | "what" - | "who" - | "why" - | "how" - | "did" - | "does" - | "can" - | "could" - | "would" - | "with" - | "from" - | "that" - | "this" - | "there" - | "today" - | "current" - | "currently" - | "working" - | "touching" - | "using" - | "about" - | "related" - | "unrelated" - | "without" - | "exclude" - | "excluding" - ) -} - -const ENTITY_BOUNDARY_CLASSES: &[&[&[&str]]] = &[ - &[ - &["film", "films", "movie", "movies"], - &[ - "camera", - "cameras", - "photo", - "photos", - "photograph", - "photographs", - ], - ], - &[ - &["uncle"], - &["aunt"], - &["niece"], - &["nephew"], - &["cousin"], - &["mother", "mom"], - &["father", "dad"], - &["sister"], - &["brother"], - ], - &[&["party", "parties"], &["cake", "cakes"]], -]; - -fn entity_boundary_classes_for_query(query_terms: &HashSet) -> Vec> { - ENTITY_BOUNDARY_CLASSES - .iter() - .map(|classes| { - classes - .iter() - .enumerate() - .filter_map(|(index, variants)| { - variants - .iter() - .any(|term| query_terms.contains(*term)) - .then_some(index) - }) - .collect::>() - }) - .filter(|indices| !indices.is_empty()) - .collect() -} - -fn entity_boundary_mismatch( - query_terms: &HashSet, - packet_terms: &HashSet, -) -> Option { - for classes in ENTITY_BOUNDARY_CLASSES { - let query_indices = classes - .iter() - .enumerate() - .filter_map(|(index, variants)| { - variants - .iter() - .any(|term| query_terms.contains(*term)) - .then_some(index) - }) - .collect::>(); - if query_indices.is_empty() { - continue; - } - - let packet_has_query_class = query_indices.iter().any(|index| { - classes[*index] - .iter() - .any(|term| packet_terms.contains(*term)) - }); - let conflicting_class = classes.iter().enumerate().find(|(index, variants)| { - !query_indices.contains(index) - && variants.iter().any(|term| packet_terms.contains(*term)) - }); - - if !packet_has_query_class { - if let Some((_, variants)) = conflicting_class { - return Some(format!( - "entity_boundary_mismatch:{}:{}", - classes[query_indices[0]][0], variants[0] - )); - } - } - } - None -} - fn xml_escape(value: &str) -> String { value .replace('&', "&") @@ -1504,97 +1253,22 @@ fn truncate_chars(value: &str, max_chars: usize) -> String { out } -fn route_memory_query(query: &str, has_explicit_refs: bool) -> LaneRoute { - if has_explicit_refs { - return LaneRoute { - lanes: vec![ - MemoryLane::Live, - MemoryLane::Episodic, - MemoryLane::Semantic, - MemoryLane::Profile, - MemoryLane::Procedural, - ], - reason: "explicit_ref".to_string(), - archival_allowed: true, - }; - } - let lower = query.to_lowercase(); - let words = lower.split_whitespace().count(); - if words <= 4 - && ["that", "this", "it", "yes", "do it", "same"] - .iter() - .any(|cue| lower.contains(cue)) - { - return LaneRoute { - lanes: vec![MemoryLane::Live], - reason: "discourse_local_live_context".to_string(), - archival_allowed: false, - }; - } - if [ - "when", - "before", - "after", - "last", - "yesterday", - "session", - "earlier", - "again", - "previous", - ] - .iter() - .any(|cue| lower.contains(cue)) - { - return LaneRoute { - lanes: vec![MemoryLane::Episodic, MemoryLane::Live, MemoryLane::Semantic], - reason: "temporal_episodic".to_string(), - archival_allowed: true, - }; - } - if [ - "prefer", - "preference", - "like", - "dislike", - "favorite", - "profile", - "who am i", - ] - .iter() - .any(|cue| lower.contains(cue)) - { - return LaneRoute { - lanes: vec![ - MemoryLane::Profile, - MemoryLane::Semantic, - MemoryLane::Episodic, - ], - reason: "profile_query".to_string(), - archival_allowed: true, - }; - } - if [ - "how do we", - "workflow", - "procedure", - "policy", - "always", - "habit", - ] - .iter() - .any(|cue| lower.contains(cue)) - { - return LaneRoute { - lanes: vec![ - MemoryLane::Procedural, - MemoryLane::Semantic, - MemoryLane::Profile, - ], - reason: "procedural_query".to_string(), - archival_allowed: true, - }; +fn route_memory_query(_query: &str, has_explicit_refs: bool) -> LaneRoute { + LaneRoute { + lanes: vec![ + MemoryLane::Live, + MemoryLane::Episodic, + MemoryLane::Semantic, + MemoryLane::Profile, + MemoryLane::Procedural, + ], + reason: if has_explicit_refs { + "explicit_ref".to_string() + } else { + "all_lanes_no_query_cues".to_string() + }, + archival_allowed: true, } - LaneRoute::default() } fn unix_timestamp() -> i64 { @@ -1878,50 +1552,7 @@ mod tests { } #[tokio::test] - async fn aggregate_query_trace_uses_collect_sized_plan() -> 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"); - let budget = ContextBudget { - max_tokens: 1_000, - top_k: 2, - graph_depth: 0, - }; - - for index in 0..6 { - pipeline - .observe_session(BenchSession { - session_id: format!("museum-session-{index}"), - timestamp: Some(index), - turns: vec![BenchTurn { - role: "user".to_string(), - content: format!( - "museum visit {index}: I went to the museum district stop" - ), - timestamp: Some(index), - }], - }) - .await?; - } - - let query = BenchQuery { - query_id: "q-museum-count".to_string(), - text: "how many museum visits did I mention?".to_string(), - reference_time: Some(10), - }; - let retrieval = pipeline.retrieve_evidence(query, budget).await?; - - assert_eq!(retrieval.op_plan.op, QueryOp::AggregateCount); - assert!(retrieval.op_plan.collect_all); - assert!(retrieval.candidates.len() > budget.top_k); - assert!(retrieval.packets.len() > budget.top_k); - Ok(()) - } - - #[tokio::test] - async fn lookup_query_trace_keeps_default_plan() -> Result<()> { + async fn query_trace_uses_lookup_plan_when_structured_planner_is_unavailable() -> Result<()> { let tmp = NamedTempFile::new()?; let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; let llm = LlmClient::new(crate::models::ModelsConfig::default()); @@ -1996,101 +1627,7 @@ mod tests { } #[tokio::test] - async fn entity_boundary_filter_blocks_sibling_entity_bleed() -> 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"); - let budget = ContextBudget { - max_tokens: 1_000, - top_k: 5, - graph_depth: 0, - }; - - pipeline - .observe_session(BenchSession { - session_id: "camera-session".to_string(), - timestamp: Some(1), - turns: vec![BenchTurn { - role: "user".to_string(), - content: "I have been collecting vintage cameras for three months.".to_string(), - timestamp: Some(1), - }], - }) - .await?; - - let query = BenchQuery { - query_id: "q-vintage-films".to_string(), - text: "how long have I been collecting vintage films?".to_string(), - reference_time: Some(3), - }; - let retrieval = pipeline.retrieve_evidence(query, budget).await?; - - assert!(!retrieval.candidates.is_empty()); - assert!(retrieval.packets.is_empty()); - assert!(retrieval.packet_omissions.iter().any(|omission| omission - .reason - .starts_with("entity_boundary_mismatch:film:camera"))); - Ok(()) - } - - #[tokio::test] - async fn negated_query_term_filter_blocks_excluded_project_bleed() -> 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"); - let budget = ContextBudget { - max_tokens: 1_000, - top_k: 5, - graph_depth: 0, - }; - - pipeline - .observe_session(BenchSession { - session_id: "dotfiles-klbr-session".to_string(), - timestamp: Some(1), - turns: vec![BenchTurn { - role: "user".to_string(), - content: "klbr note: the dotfiles syncer shares the rust cli harness." - .to_string(), - timestamp: Some(1), - }], - }) - .await?; - - let query = BenchQuery { - query_id: "q-not-klbr".to_string(), - text: "working on the dotfiles syncer today, not touching klbr".to_string(), - reference_time: Some(3), - }; - let retrieval = pipeline.retrieve_evidence(query, budget).await?; - - assert!(!retrieval.candidates.is_empty()); - assert!(retrieval.packets.is_empty()); - assert!(retrieval - .packet_omissions - .iter() - .any(|omission| omission.reason == "negated_query_term:klbr")); - let context = pipeline - .assemble_context( - BenchQuery { - query_id: "q-not-klbr".to_string(), - text: "working on the dotfiles syncer today, not touching klbr".to_string(), - reference_time: Some(3), - }, - &retrieval, - budget, - ) - .await?; - assert!(!context.content.contains("rust cli harness")); - Ok(()) - } - - #[tokio::test] - async fn exact_ref_packets_bypass_query_constraint_filters() -> Result<()> { + async fn exact_ref_packet_includes_exact_payload() -> Result<()> { let tmp = NamedTempFile::new()?; let store = MemoryStore::open(tmp.path().to_str().unwrap(), 4)?; let memory_id = store.store( @@ -2109,8 +1646,8 @@ mod tests { }; let query = BenchQuery { - query_id: "q-exact-negated".to_string(), - text: format!("what does [{alias}] say? not touching klbr"), + query_id: "q-exact".to_string(), + text: format!("what does [{alias}] say?"), reference_time: Some(3), }; let retrieval = pipeline.retrieve_evidence(query.clone(), budget).await?; @@ -2132,11 +1669,10 @@ mod tests { } #[test] - fn lane_routing_prefers_live_context_for_discourse_local_prompts() { + fn lane_routing_uses_all_lanes_without_query_cues() { let local = route_memory_query("yes do that", false); - assert_eq!(local.lanes, vec![MemoryLane::Live]); - assert_eq!(local.reason, "discourse_local_live_context"); - assert!(!local.archival_allowed); + assert_eq!(local, LaneRoute::default()); + assert!(local.archival_allowed); let exact = route_memory_query("what about [m1]?", true); assert_eq!(exact.lanes.first().copied(), Some(MemoryLane::Live)); diff --git a/klbr-core/src/planner.rs b/klbr-core/src/planner.rs index 1d63ed6..b395659 100644 --- a/klbr-core/src/planner.rs +++ b/klbr-core/src/planner.rs @@ -26,6 +26,8 @@ impl QueryOp { } } +const MODEL_REASON_MAX_CHARS: usize = 120; + #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct OpPlan { pub op: QueryOp, @@ -42,12 +44,31 @@ impl Default for OpPlan { collect_all: false, require_personal_support: false, max_passes: 1, - reason: "lookup_default".to_string(), + reason: "planner_unavailable_default_lookup".to_string(), } } } impl OpPlan { + pub fn for_op(op: QueryOp, reason: impl Into) -> Self { + let collect_all = matches!( + op, + QueryOp::AggregateCount + | QueryOp::AggregateSum + | QueryOp::AggregateAvg + | QueryOp::OrderOrRank + | QueryOp::UpdateResolution + ); + let require_personal_support = matches!(op, QueryOp::PreferenceRecommendation); + Self { + op, + collect_all, + require_personal_support, + max_passes: if collect_all { 2 } else { 1 }, + reason: reason.into(), + } + } + pub fn candidate_limit(&self, top_k: usize) -> usize { let base = top_k.max(1); if self.collect_all { @@ -78,202 +99,47 @@ impl OpPlan { } } -pub fn classify_query_op(query: &str) -> OpPlan { - let lower = normalize_query(query); - let tokens = lower.split_whitespace().collect::>(); - - if contains_any( - &lower, - &["don't know", "do not know", "unknown", "not sure"], - ) || contains_any(&lower, &["did i ever", "have i ever", "was there any"]) - { - return OpPlan { - op: QueryOp::AbstainOrFalsePremiseCheck, - collect_all: false, - require_personal_support: false, - max_passes: 1, - reason: "false_premise_or_abstain_cue".to_string(), - }; - } - - if contains_any( - &lower, - &[ - "previous", - "latest", - "current", - "most recent", - "last time", - "last week", - "last month", - "last year", - "last saturday", - "last sunday", - "last monday", - "last tuesday", - "last wednesday", - "last thursday", - "last friday", - "newer", - "older", - "before", - "after", - "ago", - "a month ago", - "a week ago", - ], - ) { - return OpPlan { - op: QueryOp::UpdateResolution, - collect_all: true, - require_personal_support: false, - max_passes: 2, - reason: "update_or_temporal_resolution_cue".to_string(), - }; - } +pub const STRUCTURED_QUERY_PLANNER_PROMPT: &str = r#"classify the user's memory question by the abstract operation needed to answer it. +return only a json object with: +{"op":"lookup|aggregate_count|aggregate_sum|aggregate_avg|order_or_rank|update_resolution|preference_recommendation|abstain_or_false_premise_check","reason":"short rationale"} +the classifier must work across languages and scripts. do not use keyword matching. do not answer the question."#; - if contains_any( - &lower, - &[ - "in order", - "what order", - "which order", - "ordered", - "chronological", - "rank", - "ranked", - "earliest", - "oldest", - "newest", - "first", - "second", - "third", - ], - ) { - return OpPlan { - op: QueryOp::OrderOrRank, - collect_all: true, - require_personal_support: false, - max_passes: 2, - reason: "order_or_rank_cue".to_string(), - }; - } - - if contains_any(&lower, &["average", "mean"]) { - return OpPlan { - op: QueryOp::AggregateAvg, - collect_all: true, - require_personal_support: false, - max_passes: 2, - reason: "average_cue".to_string(), - }; - } - - if contains_any(&lower, &["total", "sum", "combined"]) && has_numeric_or_unit_cue(&tokens) { - return OpPlan { - op: QueryOp::AggregateSum, - collect_all: true, - require_personal_support: false, - max_passes: 2, - reason: "sum_cue".to_string(), - }; - } - - if contains_any( - &lower, - &[ - "how many", - "number of", - "count", - "how much", - "how often", - "times did", - ], - ) { - return OpPlan { - op: QueryOp::AggregateCount, - collect_all: true, - require_personal_support: false, - max_passes: 2, - reason: "count_cue".to_string(), - }; - } - - if lower.contains("how old") && lower.contains("when") { - return OpPlan { - op: QueryOp::AggregateCount, - collect_all: true, - require_personal_support: false, - max_passes: 2, - reason: "age_at_event_cue".to_string(), - }; - } +#[derive(Debug, serde::Deserialize)] +struct StructuredPlannerReply { + op: QueryOp, + #[serde(default)] + reason: Option, +} - if contains_any( - &lower, - &[ - "recommend", - "recommendation", - "suggest", - "suggestion", - "tips", - "advice", - "ideas", - "best fit", - "would i like", - "should i try", - ], - ) || (contains_any( - &lower, - &["favorite", "prefer", "preference", "liked", "disliked"], - ) && contains_any(&lower, &["what", "which", "where"])) - { - return OpPlan { - op: QueryOp::PreferenceRecommendation, - collect_all: false, - require_personal_support: true, - max_passes: 1, - reason: "preference_recommendation_cue".to_string(), - }; - } +pub fn parse_structured_op_plan(reply: &str) -> Option { + let json = extract_json_object(reply)?; + let parsed = serde_json::from_str::(&json).ok()?; + let reason = parsed + .reason + .as_deref() + .map(sanitize_reason) + .filter(|reason| !reason.is_empty()) + .unwrap_or_else(|| "model_structured".to_string()); + Some(OpPlan::for_op( + parsed.op, + format!("model_structured:{reason}"), + )) +} - OpPlan::default() +fn extract_json_object(value: &str) -> Option { + let start = value.find('{')?; + let end = value.rfind('}')?; + (start <= end).then(|| value[start..=end].to_string()) } -fn normalize_query(query: &str) -> String { - query - .to_lowercase() - .replace(['?', '!', ',', '.', ';', ':', '"', '\''], " ") +fn sanitize_reason(value: &str) -> String { + value .split_whitespace() .collect::>() .join(" ") -} - -fn contains_any(haystack: &str, needles: &[&str]) -> bool { - needles.iter().any(|needle| haystack.contains(needle)) -} - -fn has_numeric_or_unit_cue(tokens: &[&str]) -> bool { - tokens.iter().any(|token| { - token.chars().any(|c| c.is_ascii_digit()) - || matches!( - *token, - "weeks" - | "days" - | "hours" - | "minutes" - | "miles" - | "kilometers" - | "km" - | "gpa" - | "score" - | "scores" - | "cost" - | "costs" - | "price" - | "prices" - ) - }) + .chars() + .take(MODEL_REASON_MAX_CHARS) + .collect() } #[cfg(test)] @@ -281,81 +147,46 @@ mod tests { use super::*; #[test] - fn planner_detects_aggregate_count() { - let plan = classify_query_op("how many competitive sports did I do?"); - - assert_eq!(plan.op, QueryOp::AggregateCount); - assert!(plan.collect_all); - assert_eq!(plan.candidate_limit(5), 15); - assert_eq!(plan.packet_limit(5), 10); - } - - #[test] - fn planner_detects_average() { - let plan = classify_query_op("what was the average of my two GPA values?"); + fn default_plan_is_conservative_lookup() { + let plan = OpPlan::default(); - assert_eq!(plan.op, QueryOp::AggregateAvg); - assert!(plan.collect_all); - } - - #[test] - fn planner_detects_temporal_update_resolution_before_ordering() { - let plan = classify_query_op("what was my previous personal best before the latest one?"); - - assert_eq!(plan.op, QueryOp::UpdateResolution); - assert!(plan.collect_all); + assert_eq!(plan.op, QueryOp::Lookup); + assert!(!plan.collect_all); + assert_eq!(plan.candidate_limit(5), 5); + assert_eq!(plan.packet_limit(5), 5); } #[test] - fn planner_detects_order_or_rank() { - let plan = classify_query_op("list the museums I visited in chronological order"); + fn collect_plan_expands_limits_without_query_text_rules() { + let plan = OpPlan::for_op(QueryOp::AggregateCount, "model_structured"); - assert_eq!(plan.op, QueryOp::OrderOrRank); assert!(plan.collect_all); + assert_eq!(plan.candidate_limit(5), 15); + assert_eq!(plan.packet_limit(5), 10); } #[test] - fn planner_detects_preference_recommendation() { - let plan = - classify_query_op("what music venue in Denver should I try based on what I liked?"); + fn preference_plan_requires_personal_support() { + let plan = OpPlan::for_op(QueryOp::PreferenceRecommendation, "model_structured"); - assert_eq!(plan.op, QueryOp::PreferenceRecommendation); assert!(plan.require_personal_support); assert!(!plan.collect_all); - assert_eq!(plan.candidate_limit(5), 10); - } - - #[test] - fn planner_detects_advice_as_preference_recommendation() { - let plan = classify_query_op( - "I was thinking about rearranging bedroom furniture this weekend. Any tips?", - ); - - assert_eq!(plan.op, QueryOp::PreferenceRecommendation); - assert!(plan.require_personal_support); } #[test] - fn planner_detects_relative_time_prompts() { - let plan = classify_query_op("who did I go with to the music event last Saturday?"); + fn parses_structured_model_reply() { + let plan = parse_structured_op_plan( + r#"{"op":"update_resolution","reason":"competing values over time"}"#, + ) + .expect("structured plan should parse"); assert_eq!(plan.op, QueryOp::UpdateResolution); assert!(plan.collect_all); + assert!(plan.reason.starts_with("model_structured:")); } #[test] - fn planner_detects_age_at_event_as_collect_mode() { - let plan = classify_query_op("how old was I when Alex was born?"); - - assert_eq!(plan.op, QueryOp::AggregateCount); - assert!(plan.collect_all); - } - - #[test] - fn planner_leaves_plain_lookup_alone() { - let plan = classify_query_op("where did I buy coffee creamer?"); - - assert_eq!(plan.op, QueryOp::Lookup); - assert!(!plan.collect_all); + fn ignores_non_json_model_reply() { + assert!(parse_structured_op_plan("lookup probably").is_none()); } } -- 2.51.2