diff --git a/src/openhuman/autocomplete/core.rs b/src/openhuman/autocomplete/core.rs index d1634eb32..d60a3db4d 100644 --- a/src/openhuman/autocomplete/core.rs +++ b/src/openhuman/autocomplete/core.rs @@ -418,6 +418,7 @@ impl AutocompleteEngine { show_overflow_badge("accepted", Some(&cleaned), None, None, None); // Persist acceptance for personalisation (fire-and-forget). + // Dual-write: KV (UI list) + local docs (semantic search). { let (ctx, app) = { let s = self.inner.lock().await; @@ -431,6 +432,12 @@ impl AutocompleteEngine { app.as_deref(), ) .await; + crate::openhuman::autocomplete::history::save_completion_to_local_docs( + &ctx, + &sug, + app.as_deref(), + ) + .await; }); } @@ -593,13 +600,32 @@ impl AutocompleteEngine { } let service = local_ai::global(&config); - // Build personalised style examples: dynamic (from accepted-history) + static (user config). - let dynamic_examples = - crate::openhuman::autocomplete::history::load_recent_examples(6).await; + // Build personalised style examples from three sources: + // 1. Semantically relevant past completions (local doc query) + // 2. Most recent past completions (KV recency signal / fallback) + // 3. Static user-configured examples + // Deduplicated and capped at 8 total. + let relevant_examples = + crate::openhuman::autocomplete::history::query_relevant_examples(&context, 4).await; + let recent_examples = + crate::openhuman::autocomplete::history::load_recent_examples(4).await; + let static_examples = config.autocomplete.style_examples.clone(); + let merged_examples: Vec = { - let mut v = dynamic_examples; - v.extend(config.autocomplete.style_examples.iter().cloned()); - v.truncate(8); + let mut seen = std::collections::HashSet::new(); + let mut v = Vec::new(); + for ex in relevant_examples + .into_iter() + .chain(recent_examples) + .chain(static_examples) + { + if seen.insert(ex.clone()) { + v.push(ex); + } + if v.len() >= 8 { + break; + } + } v }; @@ -695,6 +721,7 @@ impl AutocompleteEngine { show_overflow_badge("accepted", Some(&cleaned), None, None, None); // Persist acceptance for personalisation (fire-and-forget). + // Dual-write: KV (UI list) + local docs (semantic search). { let (ctx, app) = { let s = self.inner.lock().await; @@ -708,6 +735,12 @@ impl AutocompleteEngine { app.as_deref(), ) .await; + crate::openhuman::autocomplete::history::save_completion_to_local_docs( + &ctx, + &sug, + app.as_deref(), + ) + .await; }); } } diff --git a/src/openhuman/autocomplete/history.rs b/src/openhuman/autocomplete/history.rs index 16d61c220..fcf673b23 100644 --- a/src/openhuman/autocomplete/history.rs +++ b/src/openhuman/autocomplete/history.rs @@ -4,12 +4,15 @@ //! "autocomplete" namespace and fed back as dynamic style examples on the //! next inference cycle, giving the model in-context personalisation. -use crate::openhuman::memory::MemoryClient; +use crate::openhuman::memory::{MemoryClient, NamespaceDocumentInput}; use chrono::Utc; use serde::{Deserialize, Serialize}; +use serde_json::json; const AUTOCOMPLETE_KV_NAMESPACE: &str = "autocomplete"; +const AUTOCOMPLETE_DOC_NAMESPACE: &str = "autocomplete-memory"; const MAX_HISTORY_ENTRIES: usize = 50; +const MAX_DOC_ENTRIES: usize = 200; const CONTEXT_TAIL_CHARS: usize = 40; /// A single accepted completion record persisted in the KV store. @@ -73,6 +76,144 @@ pub async fn save_accepted_completion(context: &str, suggestion: &str, app_name: } } +/// Persist an accepted completion as a local memory document (fire-and-forget safe). +/// +/// Documents are stored in the `"autocomplete-memory"` namespace and are +/// searchable via `query_namespace`, enabling semantic matching of past +/// completions against the current typing context. +pub async fn save_completion_to_local_docs( + context: &str, + suggestion: &str, + app_name: Option<&str>, +) { + let client = match MemoryClient::new_local() { + Ok(c) => c, + Err(e) => { + log::warn!("[autocomplete:history] local doc — client init failed: {e}"); + return; + } + }; + + let ts_ms = Utc::now().timestamp_millis(); + let key = format!("completion:{ts_ms:018}"); + let app = app_name.unwrap_or("unknown"); + + // Build the same formatted string used by load_recent_examples so that + // query results are directly usable as style examples in inference. + let tail: String = context + .chars() + .rev() + .take(CONTEXT_TAIL_CHARS) + .collect::() + .chars() + .rev() + .collect(); + let formatted = format!("[{app}] ...{tail} → {suggestion}"); + + let mut tags = vec!["autocomplete".to_string(), "accepted".to_string()]; + if let Some(name) = app_name { + tags.push(name.to_string()); + } + + let input = NamespaceDocumentInput { + namespace: AUTOCOMPLETE_DOC_NAMESPACE.to_string(), + key, + title: format!("Accepted completion — {app}"), + content: formatted, + source_type: "autocomplete".to_string(), + priority: "low".to_string(), + tags, + metadata: json!({ + "context": context, + "suggestion": suggestion, + "app_name": app_name, + "timestamp_ms": ts_ms, + }), + category: "daily".to_string(), + session_id: None, + document_id: None, + }; + + if let Err(e) = client.put_doc(input).await { + log::warn!("[autocomplete:history] local doc put_doc failed: {e}"); + return; + } + + log::debug!("[autocomplete:history] saved local doc completion ts={ts_ms}"); + + // Trim to MAX_DOC_ENTRIES — delete oldest documents beyond the limit. + if let Ok(docs) = client + .list_documents(Some(AUTOCOMPLETE_DOC_NAMESPACE)) + .await + { + let items = docs + .get("documents") + .and_then(serde_json::Value::as_array) + .cloned() + .unwrap_or_default(); + if items.len() > MAX_DOC_ENTRIES { + for item in items.into_iter().skip(MAX_DOC_ENTRIES) { + if let Some(doc_id) = item.get("documentId").and_then(serde_json::Value::as_str) { + let _ = client + .delete_document(AUTOCOMPLETE_DOC_NAMESPACE, doc_id) + .await; + } + } + } + } +} + +/// Query the local document store for accepted completions semantically +/// relevant to the current typing `context`. +/// +/// Uses `query_namespace` (keyword + optional vector ranking) against the +/// `"autocomplete-memory"` namespace. Returns up to `n` formatted style +/// example strings ready for injection into the inference prompt. +pub async fn query_relevant_examples(context: &str, n: usize) -> Vec { + let client = match MemoryClient::new_local() { + Ok(c) => c, + Err(e) => { + log::warn!("[autocomplete:history] query_relevant — client init failed: {e}"); + return Vec::new(); + } + }; + + // Use the tail of the current context as the search query. + let tail: String = context + .chars() + .rev() + .take(80) + .collect::() + .chars() + .rev() + .collect(); + + let result = match client + .query_namespace(AUTOCOMPLETE_DOC_NAMESPACE, &tail, n as u32) + .await + { + Ok(r) if !r.is_empty() => r, + Ok(_) => return Vec::new(), + Err(e) => { + log::warn!("[autocomplete:history] query_namespace failed: {e}"); + return Vec::new(); + } + }; + + // query_namespace_context returns "key: content" entries joined by "\n\n". + // The content is already in "[app] ...tail → suggestion" format. + result + .split("\n\n") + .filter(|s| !s.is_empty()) + .filter_map(|entry| { + // Strip the "completion:XXXXXXXXXXXXXXXXXX: " key prefix. + let bracket_pos = entry.find('[')?; + Some(entry[bracket_pos..].to_string()) + }) + .take(n) + .collect() +} + /// Load the `n` most recent accepted completions as formatted style example strings. /// /// Each string has the form: `"[AppName] ...{tail} → suggestion"` @@ -129,16 +270,50 @@ pub async fn list_history(limit: usize) -> Result, Strin Ok(entries) } -/// Delete all accepted-completion entries. Returns the number of entries removed. +/// Delete all accepted-completion entries across all layers. +/// Returns the total number of entries removed (KV + local docs). pub async fn clear_history() -> Result { let client = MemoryClient::new_local()?; + + // 1. Clear KV entries (existing behaviour — powers the UI list). let rows = client.kv_list_namespace(AUTOCOMPLETE_KV_NAMESPACE).await?; - let count = rows.len(); + let kv_count = rows.len(); for row in &rows { if let Some(k) = row["key"].as_str() { let _ = client.kv_delete(Some(AUTOCOMPLETE_KV_NAMESPACE), k).await; } } - log::debug!("[autocomplete:history] cleared {count} entries"); - Ok(count) + + // 2. Clear local document entries (semantic search layer). + let doc_count = match client + .list_documents(Some(AUTOCOMPLETE_DOC_NAMESPACE)) + .await + { + Ok(docs) => { + let items = docs + .get("documents") + .and_then(serde_json::Value::as_array) + .cloned() + .unwrap_or_default(); + let count = items.len(); + for item in items { + if let Some(doc_id) = item.get("documentId").and_then(serde_json::Value::as_str) { + let _ = client + .delete_document(AUTOCOMPLETE_DOC_NAMESPACE, doc_id) + .await; + } + } + count + } + Err(e) => { + log::warn!("[autocomplete:history] clear docs — list_documents failed: {e}"); + 0 + } + }; + + let total = kv_count + doc_count; + log::debug!( + "[autocomplete:history] cleared {kv_count} KV + {doc_count} doc entries ({total} total)" + ); + Ok(total) } diff --git a/src/openhuman/autocomplete/mod.rs b/src/openhuman/autocomplete/mod.rs index 275c901bd..4866b9915 100644 --- a/src/openhuman/autocomplete/mod.rs +++ b/src/openhuman/autocomplete/mod.rs @@ -5,7 +5,8 @@ mod schemas; pub use core::*; pub use history::{ - clear_history, list_history, load_recent_examples, save_accepted_completion, AcceptedCompletion, + clear_history, list_history, load_recent_examples, query_relevant_examples, + save_accepted_completion, save_completion_to_local_docs, AcceptedCompletion, }; pub use ops as rpc; pub use ops::*; diff --git a/src/openhuman/memory/store/client.rs b/src/openhuman/memory/store/client.rs index bb2b1e7a9..cd1d39863 100644 --- a/src/openhuman/memory/store/client.rs +++ b/src/openhuman/memory/store/client.rs @@ -8,7 +8,6 @@ use crate::openhuman::memory::store::types::{ NamespaceDocumentInput, NamespaceMemoryHit, NamespaceRetrievalContext, }; use crate::openhuman::memory::store::unified::UnifiedMemory; - pub type MemoryClientRef = Arc; pub struct MemoryState(pub std::sync::Mutex>); diff --git a/tests/autocomplete_memory_e2e.rs b/tests/autocomplete_memory_e2e.rs new file mode 100644 index 000000000..0d05bfc88 --- /dev/null +++ b/tests/autocomplete_memory_e2e.rs @@ -0,0 +1,210 @@ +//! E2E tests for autocomplete memory storage (Issue #108). +//! +//! Validates the full accept → store → query → clear lifecycle against a real +//! local `MemoryClient` backed by SQLite in a temp workspace. +//! +//! Run with: `cargo test --test autocomplete_memory_e2e` + +use std::path::Path; +use std::sync::{Mutex, OnceLock}; + +use tempfile::tempdir; + +use openhuman_core::openhuman::autocomplete::history; + +// ── Env isolation ──────────────────────────────────────────────────── + +struct EnvVarGuard { + key: &'static str, + old: Option, +} + +impl EnvVarGuard { + fn set_to_path(key: &'static str, path: &Path) -> Self { + let old = std::env::var(key).ok(); + std::env::set_var(key, path.as_os_str()); + Self { key, old } + } +} + +impl Drop for EnvVarGuard { + fn drop(&mut self) { + match &self.old { + Some(v) => std::env::set_var(self.key, v), + None => std::env::remove_var(self.key), + } + } +} + +/// Serialises tests: `HOME` is process-global. +static ENV_LOCK: OnceLock> = OnceLock::new(); + +fn env_lock() -> std::sync::MutexGuard<'static, ()> { + ENV_LOCK + .get_or_init(|| Mutex::new(())) + .lock() + .expect("env lock poisoned") +} + +// ── Tests ──────────────────────────────────────────────────────────── + +/// Acceptance criteria 1 & 2: completions are written to memory and retrievable. +#[tokio::test] +async fn accepted_completions_stored_and_retrievable() { + let _lock = env_lock(); + let tmp = tempdir().expect("tempdir"); + let _home = EnvVarGuard::set_to_path("HOME", tmp.path()); + + // Write three completions with different contexts. + history::save_accepted_completion("fn main() { let x =", "42;", Some("VSCode")).await; + history::save_completion_to_local_docs("fn main() { let x =", "42;", Some("VSCode")).await; + + history::save_accepted_completion("def hello():", " print('hi')", Some("PyCharm")).await; + history::save_completion_to_local_docs("def hello():", " print('hi')", Some("PyCharm")) + .await; + + history::save_accepted_completion("const app = express", "()", Some("WebStorm")).await; + history::save_completion_to_local_docs("const app = express", "()", Some("WebStorm")).await; + + // KV history should contain all three (newest first). + let kv_entries = history::list_history(10).await.expect("list_history"); + assert_eq!( + kv_entries.len(), + 3, + "expected 3 KV entries, got {}", + kv_entries.len() + ); + + // Recent examples should be formatted correctly. + let recent = history::load_recent_examples(10).await; + assert_eq!(recent.len(), 3); + for ex in &recent { + assert!(ex.contains("→"), "example should contain arrow: {ex}"); + assert!(ex.starts_with('['), "example should start with [app]: {ex}"); + } + + // Semantic query: searching for "express" should return the JS completion. + let relevant = history::query_relevant_examples("const app = express", 5).await; + // With NoopEmbedding, keyword search should still match. + assert!( + !relevant.is_empty(), + "query_relevant_examples should return at least one result for matching context" + ); + let has_express = relevant.iter().any(|r| r.contains("()")); + assert!( + has_express, + "should find the express() completion via keyword match: {relevant:?}" + ); +} + +/// Acceptance criteria 3: completions are used for future improvement (merge pipeline). +#[tokio::test] +async fn completions_improve_future_suggestions_via_merge() { + let _lock = env_lock(); + let tmp = tempdir().expect("tempdir"); + let _home = EnvVarGuard::set_to_path("HOME", tmp.path()); + + // Populate with several completions. + for i in 0..5 { + let ctx = format!("context_{i} let value ="); + let sug = format!("suggestion_{i}"); + history::save_accepted_completion(&ctx, &sug, Some("TestApp")).await; + history::save_completion_to_local_docs(&ctx, &sug, Some("TestApp")).await; + } + + // Semantic query returns relevant results. + let relevant = history::query_relevant_examples("let value =", 4).await; + // Recent examples returns recent results. + let recent = history::load_recent_examples(4).await; + + // Simulate the merge pipeline from refresh(): relevant → recent → static, deduped, max 8. + let static_examples = vec!["[static] ...typing → completion".to_string()]; + let merged: Vec = { + let mut seen = std::collections::HashSet::new(); + let mut v = Vec::new(); + for ex in relevant.into_iter().chain(recent).chain(static_examples) { + if seen.insert(ex.clone()) { + v.push(ex); + } + if v.len() >= 8 { + break; + } + } + v + }; + + assert!(!merged.is_empty(), "merged examples should not be empty"); + assert!( + merged.len() <= 8, + "merged examples should be capped at 8, got {}", + merged.len() + ); + // Static example should be present (appended after dynamic ones). + let has_static = merged.iter().any(|e| e.contains("[static]")); + assert!( + has_static, + "static example should be in merged set: {merged:?}" + ); +} + +/// Acceptance criteria 4 (partial): clear_history removes all layers. +#[tokio::test] +async fn clear_history_removes_kv_and_docs() { + let _lock = env_lock(); + let tmp = tempdir().expect("tempdir"); + let _home = EnvVarGuard::set_to_path("HOME", tmp.path()); + + // Insert completions into both layers. + for i in 0..3 { + let ctx = format!("clear_test_{i}"); + history::save_accepted_completion(&ctx, "sug", None).await; + history::save_completion_to_local_docs(&ctx, "sug", None).await; + } + + // Verify they exist. + let before = history::list_history(10).await.expect("list before clear"); + assert_eq!(before.len(), 3); + + // Clear. + let cleared = history::clear_history().await.expect("clear_history"); + assert!( + cleared >= 3, + "should have cleared at least 3 entries, got {cleared}" + ); + + // Verify empty. + let after = history::list_history(10).await.expect("list after clear"); + assert!( + after.is_empty(), + "history should be empty after clear, got {}", + after.len() + ); + + // Semantic query should also return nothing. + let relevant = history::query_relevant_examples("clear_test", 5).await; + assert!( + relevant.is_empty(), + "query should return empty after clear: {relevant:?}" + ); +} + +/// Edge case: trimming keeps only MAX_HISTORY_ENTRIES (50) in KV. +#[tokio::test] +async fn kv_history_trims_beyond_max() { + let _lock = env_lock(); + let tmp = tempdir().expect("tempdir"); + let _home = EnvVarGuard::set_to_path("HOME", tmp.path()); + + // Insert 55 completions (MAX_HISTORY_ENTRIES = 50). + for i in 0..55 { + let ctx = format!("trim_test_{i:03}"); + history::save_accepted_completion(&ctx, "s", None).await; + } + + let entries = history::list_history(100).await.expect("list_history"); + assert!( + entries.len() <= 50, + "KV history should be trimmed to 50, got {}", + entries.len() + ); +}