diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 69166a4..dfce9ab 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -18,8 +18,9 @@ jobs: - uses: actions/checkout@v4 - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: + toolchain: "1.97" components: rustfmt - name: Check formatting @@ -37,8 +38,9 @@ jobs: sudo apt-get install -y protobuf-compiler libclang-dev - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: + toolchain: "1.97" components: clippy - name: Cache cargo registry @@ -72,7 +74,9 @@ jobs: echo "LIBCLANG_PATH=$(brew --prefix llvm)/lib" >> $GITHUB_ENV - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.97" - name: Cache cargo registry uses: Swatinem/rust-cache@v2 @@ -105,7 +109,9 @@ jobs: echo "LIBCLANG_PATH=$(brew --prefix llvm)/lib" >> $GITHUB_ENV - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.97" - name: Cache cargo registry uses: Swatinem/rust-cache@v2 @@ -127,7 +133,9 @@ jobs: sudo apt-get install -y protobuf-compiler libclang-dev - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.97" - name: Cache cargo registry uses: Swatinem/rust-cache@v2 @@ -151,7 +159,9 @@ jobs: sudo apt-get install -y protobuf-compiler libclang-dev - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.97" - name: Cache cargo registry uses: Swatinem/rust-cache@v2 @@ -194,7 +204,9 @@ jobs: sudo apt-get install -y protobuf-compiler libclang-dev - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.97" - name: Cache cargo registry uses: Swatinem/rust-cache@v2 diff --git a/.github/workflows/e2e-cli.yml b/.github/workflows/e2e-cli.yml index 7d789c7..2b98be4 100644 --- a/.github/workflows/e2e-cli.yml +++ b/.github/workflows/e2e-cli.yml @@ -38,7 +38,9 @@ jobs: echo "LIBCLANG_PATH=$(brew --prefix llvm)/lib" >> $GITHUB_ENV - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master + with: + toolchain: "1.97" - name: Cache cargo registry uses: Swatinem/rust-cache@v2 diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 9fabdcd..f14356a 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -93,8 +93,9 @@ jobs: cargo generate-lockfile - name: Install Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: + toolchain: "1.97" targets: ${{ matrix.target }} - name: Install cross diff --git a/.planning/ROADMAP.md b/.planning/ROADMAP.md index 039f2af..9a164ef 100644 --- a/.planning/ROADMAP.md +++ b/.planning/ROADMAP.md @@ -12,7 +12,7 @@ - ✅ **v2.6 Cognitive Retrieval** — Phases 39-44 (shipped 2026-03-16) - ✅ **v2.7 Multi-Runtime Portability** — Phases 45-50 (shipped 2026-03-22) - **v3.0 Competitive Parity & Benchmarks** — Phases 51-53 + Phase 51.5 (in progress; Phase 51.5 merged 2026-04-28) -- **v3.1 Make It True** — Phases 54-58 (in progress; Phase 54 merged 2026-08-30, Phase 55 executing) +- **v3.1 Make It True** — Phases 54-58 (in progress; Phase 54 merged 2026-08-30, Phase 55 merged 2026-08-30, Phase 54.5 executing) ## Phases @@ -258,22 +258,29 @@ See: `docs/plans/v3.1-make-it-true-plan.md` Close the claim/reality gap, then open the shop window. No new capabilities. -### Phase 54: Integration Truth (6/6 plans) — IN EXECUTION 2026-08-30 +### Phase 54: Integration Truth (6/6 plans) — COMPLETE 2026-08-30 (PR #32) -- [ ] 54-01: Wire orchestrator into RouteQuery + real LLM reranker -- [ ] 54-02: Fix BM25 outbox no-op (index events) -- [ ] 54-03: Make Hybrid layer hybrid -- [ ] 54-04: One rank-fusion implementation -- [ ] 54-05: Honest daemon flags + attach indexes -- [ ] 54-06: Lock-poisoning recover_lock policy +- [x] 54-01: Wire orchestrator into RouteQuery + real LLM reranker +- [x] 54-02: Fix BM25 outbox no-op (index events) +- [x] 54-03: Make Hybrid layer hybrid +- [x] 54-04: One rank-fusion implementation +- [x] 54-05: Honest daemon flags + attach indexes +- [x] 54-06: Lock-poisoning recover_lock policy -### Phase 55: Performance Truth (0/2) +### Phase 54.5: Truth leaks + CI pin (in execution 2026-08-30) -### Phase 56: Honest Benchmarks (0/3) +- [ ] CI: pin rust-toolchain 1.97; allow `result_large_err` on generated proto +- [ ] Explainability reports what actually ran +- [ ] LLM rerank order survives salience; BM25 events carry text +- [ ] Shared HNSW handle; concurrent fan-out; no per-event grip scan + +### Phase 55: Performance Truth — COMPLETE 2026-08-30 (PR #33) + +### Phase 56: Honest Benchmarks (PR #34 open) ### Phase 57: Shop Window & Positioning (0/3) ### Phase 58: Launch (side quest) -*Updated: 2026-08-30 — Phase 54 Integration Truth in execution* +*Updated: 2026-08-30 — Phase 54.5 truth-leaks + CI pin in execution* diff --git a/.planning/STATE.md b/.planning/STATE.md index 991d1f3..b69e950 100644 --- a/.planning/STATE.md +++ b/.planning/STATE.md @@ -3,11 +3,11 @@ gsd_state_version: 1.0 milestone_name: Make It True status: in_progress stopped_at: null -last_updated: "2026-08-30T17:30:00.000Z" -last_activity: 2026-08-30 — Phase 55 Performance Truth implemented (medium/warm/30 artifact) +last_updated: "2026-08-30T19:10:00.000Z" +last_activity: 2026-08-30 — Phase 54.5 truth-leaks + CI toolchain pin progress: - total_phases: 5 - completed_phases: 1 + total_phases: 6 + completed_phases: 2 total_plans: 14 completed_plans: 8 percent: 57 @@ -20,27 +20,30 @@ progress: See: .planning/PROJECT.md (updated 2026-03-22) **Core value:** Agent can answer "what were we talking about last week?" without scanning everything -**Current focus:** v3.1 Phase 55 — Performance Truth (setup vs query split; honest percentiles) +**Current focus:** v3.1 Phase 54.5 — close residual Phase 54 honesty leaks and pin CI toolchain ## Current Position -Phase: 55 of 58 (Performance Truth) -Plan: 01-02 implemented on `feature/phase-55-performance-truth` (PR pending) -Status: Phase 54 merged; Phase 55 code + medium/warm/30 artifact ready -Last activity: 2026-08-30 — `single.toc` query p50 = 0.13ms; 64.6s was `toc_build` +Phase: 54.5 of 58 (Truth leaks + CI pin) +Plan: implementing on `feature/phase-54.5-truth-leaks` +Status: Phase 54 merged (#32); Phase 55 merged (#33); Phase 56 PR #34 open (clippy red from toolchain drift) +Last activity: 2026-08-30 — explainability truth, shared HNSW, rust-toolchain.toml 1.97 -Progress: [██████░░░░] ~57% (8/14 plans; Phase 55 of 54-58) +Progress: [██████░░░░] ~57% (8/14 plans; Phase 54.5 cleanup) ## Out-of-band Work ### Open PRs -None. +| PR | What | Status | +|---|---|---| +| #34 | Phase 56 Honest Benchmarks | Open; Clippy red (1.98 `result_large_err` on generated tonic stubs) | ### Recently Merged | PR | What | Merged | |---|---|---| +| #33 | Phase 55 Performance Truth | 2026-08-30 | | #32 | Phase 54 Integration Truth | 2026-08-30 | | #31 | v3.1 Make It True design spec | 2026-08-30 | | #30 | Phase 53 Benchmark Suite | 2026-08-30 | @@ -51,5 +54,5 @@ None. ## Decisions - v3.1 scope: Make It True — no new capabilities; close claim/reality gap (Phases 54-58) -- Phase 55: split setup vs query in `perf_bench`; p90/p99 withheld below 10/30 samples -- Warm = one setup + N query samples; cold = new store per iteration +- Phase 54.5 before more measurement: explainability must report what ran; shared HNSW handle +- CI pins `rust-toolchain.toml` to 1.97 so floating stable cannot redden main diff --git a/.planning/phases/54.5-truth-leaks/54.5-CONTEXT.md b/.planning/phases/54.5-truth-leaks/54.5-CONTEXT.md new file mode 100644 index 0000000..95f97b3 --- /dev/null +++ b/.planning/phases/54.5-truth-leaks/54.5-CONTEXT.md @@ -0,0 +1,8 @@ +# Phase 54.5: Truth leaks + CI pin + +**Gathered:** 2026-08-30 +**Status:** In execution +**Source:** docs/plans/phase-54.5-truth-leaks-plan.md + +Close residual claim/reality leaks in Phase 54 wiring and unpin CI from +floating stable Clippy. No new capabilities. diff --git a/.planning/phases/54.5-truth-leaks/54.5-VERIFICATION.md b/.planning/phases/54.5-truth-leaks/54.5-VERIFICATION.md new file mode 100644 index 0000000..e55bc86 --- /dev/null +++ b/.planning/phases/54.5-truth-leaks/54.5-VERIFICATION.md @@ -0,0 +1,54 @@ +--- +phase: 54.5-truth-leaks +verified: 2026-08-30 +status: passed +--- + +# Phase 54.5: Truth leaks + CI toolchain pin + +**Phase Goal:** explainability reports what actually ran; LLM rerank survives to the +response; query/prune/dedup share one HNSW handle; Clippy on main is not a +floating-stable hostage. + +## Execution evidence + +| # | Truth | Status | Evidence | +|---|-------|--------|----------| +| 1 | Generated gRPC stubs do not trip `result_large_err` | VERIFIED | `crates/memory-service/src/lib.rs` `#![allow(clippy::result_large_err)]` on `include_proto!`; clippy `-D warnings` green on 1.97.1 | +| 2 | Toolchain is pinned | VERIFIED | `rust-toolchain.toml` channel `1.97`; CI/e2e/release use `dtolnay/rust-toolchain@master` so the pin is honored | +| 3 | LLM fail-open reports `rerank=heuristic` | VERIFIED | `test_llm_reranker_fail_open_on_completer_error`; `test_llm_reranker_fail_open_on_unparseable_order`; `test_llm_fail_open_reports_heuristic` | +| 4 | `layers_attempted` omits unsupported layers | VERIFIED | `test_layers_attempted_omits_unconfigured_layers`; `test_layers_attempted_omits_unavailable_indexes` | +| 5 | `stop_conditions` / `mode_override` are forwarded | VERIFIED | `query_ranked_with`; `test_sequential_timeout_stops_remaining_layers`; `test_mode_override_sequential_is_honored` | +| 6 | Unknown `--rerank` / `rerank_mode` is rejected | VERIFIED | clap `value_parser`; `test_parse_search_rejects_unknown_rerank`; `test_unknown_rerank_mode_is_rejected` | +| 7 | BM25 event hits carry `text` for LLM rerank | VERIFIED | schema `text` is `TEXT \| STORED`; `test_event_search_returns_text_preview`; `test_llm_rerank_reorders_bm25_hits` asserts non-empty `text_preview` | +| 8 | Salience cannot undo a successful LLM rerank | VERIFIED | `RankingConfig.preserve_order`; `test_preserve_order_skips_resort` | +| 9 | No per-event grip scan on BM25 outbox drain | VERIFIED | `find_grip_for_event` removed from `bm25_updater.rs` | +| 10 | Fan-out is concurrent unless Sequential | VERIFIED | `futures::join_all`; `test_parallel_fan_out_runs_layers_concurrently` (max in-flight ≥ 2) | +| 11 | Query / prune / dedup share one HNSW handle | VERIFIED | `open_query_indexes` runs before prune/dedup; `VectorTeleportHandler::index_handle` passed to both | +| 12 | Dimension comes from the embedder | VERIFIED | `embedder.info().dimension` / `EMBEDDING_DIM`; `HnswConfig::default` uses the constant | + +Existing BM25 indexes created before this change do not store `text` and still +need a rebuild for event previews to appear. New indexes (and tests) get the +stored field automatically. + +## Crate tests (local, 2026-08-30) + +| Crate | Result | +|---|---| +| memory-orchestrator | 34 passed | +| memory-retrieval | 78 passed | +| memory-search | 54 passed | +| memory-service | 125 passed | +| memory-indexing | 48 passed | +| memory-daemon | 93 passed | +| memory-cli | 63 passed | +| memory-vector | 28 passed | +| clippy `-D warnings` (workspace, exclude e2e) | green | +| rustfmt `--check` | green | + +## Human verification (blockers) + +- [x] `cargo clippy --workspace --all-targets --all-features --exclude e2e-tests -- -D warnings` green +- [x] Orchestrator + retrieval + search + service unit tests green +- [x] `rust-toolchain.toml` present and CI workflows use `@master` +- [ ] Full `cargo test --workspace --all-features` (CI) diff --git a/crates/e2e-tests/src/bin/perf_bench.rs b/crates/e2e-tests/src/bin/perf_bench.rs index 19b0f5d..021621c 100644 --- a/crates/e2e-tests/src/bin/perf_bench.rs +++ b/crates/e2e-tests/src/bin/perf_bench.rs @@ -693,8 +693,8 @@ async fn build_vector_index( )); } - let mut vector_id = 1_u64; for (idx, (text, agent, timestamp_ms)) in texts.iter().enumerate() { + let vector_id = idx as u64 + 1; let embedder_clone = embedder.clone(); let text_owned = text.clone(); let embedding = tokio::task::spawn_blocking(move || embedder_clone.embed(&text_owned)) @@ -714,7 +714,6 @@ async fn build_vector_index( let entry = VectorEntry::new(vector_id, DocType::TocNode, doc_id, *timestamp_ms, text) .with_agent(agent.clone()); metadata.put(&entry).map_err(|e| e.to_string())?; - vector_id += 1; } let index_lock = Arc::new(std::sync::RwLock::new(hnsw_index)); diff --git a/crates/e2e-tests/tests/ranking_test.rs b/crates/e2e-tests/tests/ranking_test.rs index 0bbc19c..0eea6bb 100644 --- a/crates/e2e-tests/tests/ranking_test.rs +++ b/crates/e2e-tests/tests/ranking_test.rs @@ -126,6 +126,7 @@ fn test_score_floor_prevents_collapse() { usage_decay_enabled: true, decay_factor: 0.15, score_floor: 0.50, + ..Default::default() }; // Worst case: low salience + extremely high access count @@ -157,6 +158,7 @@ fn test_combined_ranking_composition() { usage_decay_enabled: true, decay_factor: 0.15, score_floor: 0.50, + ..Default::default() }; // High-salience but heavily used vs low-salience but fresh diff --git a/crates/memory-cli/src/cli.rs b/crates/memory-cli/src/cli.rs index df5b724..7612eef 100644 --- a/crates/memory-cli/src/cli.rs +++ b/crates/memory-cli/src/cli.rs @@ -57,8 +57,8 @@ pub struct SearchArgs { #[arg(long, default_value_t = 10)] pub top: usize, - /// Rerank mode (e.g., "heuristic", "llm"). - #[arg(long)] + /// Rerank mode: "heuristic" or "llm". + #[arg(long, value_parser = ["heuristic", "llm"])] pub rerank: Option, /// Output format override. @@ -165,6 +165,17 @@ mod tests { } } + #[test] + fn test_parse_search_rejects_unknown_rerank() { + let err = Cli::try_parse_from(["memory", "search", "hello", "--rerank", "cross-encoder"]) + .unwrap_err(); + let msg = err.to_string(); + assert!( + msg.contains("invalid value") || msg.contains("possible values"), + "unexpected clap error: {msg}" + ); + } + #[test] fn test_parse_context() { let cli = Cli::try_parse_from(["memory", "context", "what happened"]).unwrap(); diff --git a/crates/memory-daemon/src/commands.rs b/crates/memory-daemon/src/commands.rs index e5c35c1..6a728d5 100644 --- a/crates/memory-daemon/src/commands.rs +++ b/crates/memory-daemon/src/commands.rs @@ -168,7 +168,11 @@ async fn register_indexing_job( /// Both jobs use per-level retention configured in lifecycle settings. /// BM25 pruning is DISABLED by default (per PRD append-only philosophy). /// Vector pruning is ENABLED by default. -async fn register_prune_jobs(scheduler: &SchedulerService, db_path: &Path) -> Result<()> { +async fn register_prune_jobs( + scheduler: &SchedulerService, + db_path: &Path, + shared_vector: Option>, +) -> Result<()> { use memory_embeddings::EmbeddingModel; use memory_scheduler::{ register_bm25_prune_job, register_bm25_rebuild_job, register_vector_prune_job, @@ -244,9 +248,30 @@ async fn register_prune_jobs(scheduler: &SchedulerService, db_path: &Path) -> Re info!("Search index not found, skipping BM25 prune job registration"); } - // Register vector prune job if vector index exists - if vector_dir.exists() { - // Try to create embedder + // Register vector prune job, sharing the query-path HNSW handle when + // available so prune mutations are visible without a daemon restart. + if let Some(vh) = shared_vector { + let pipeline = Arc::new(VectorIndexPipeline::new( + vh.embedder(), + vh.index_handle(), + vh.metadata_arc(), + VectorPipelineConfig::default(), + )); + let vector_job = VectorPruneJob::with_prune_fn( + VectorPruneJobConfig::default(), + move |age_days, level| { + let p = Arc::clone(&pipeline); + async move { + p.prune_level(age_days, level.as_deref()) + .map_err(|e| e.to_string()) + } + }, + ); + register_vector_prune_job(scheduler, vector_job) + .await + .context("Failed to register vector prune job")?; + info!("Vector prune job registered (shared HNSW handle)"); + } else if vector_dir.exists() { match memory_embeddings::CandleEmbedder::load_default() { Ok(embedder) => { let embedder = Arc::new(embedder); @@ -255,8 +280,6 @@ async fn register_prune_jobs(scheduler: &SchedulerService, db_path: &Path) -> Re match HnswIndex::open_or_create(hnsw_config) { Ok(hnsw_index) => { let hnsw_index = Arc::new(RwLock::new(hnsw_index)); - - // Open metadata store let metadata_path = vector_dir.join("metadata"); if metadata_path.exists() { match VectorMetadata::open(&metadata_path) { @@ -268,8 +291,6 @@ async fn register_prune_jobs(scheduler: &SchedulerService, db_path: &Path) -> Re metadata, VectorPipelineConfig::default(), )); - - // Create prune job with callback let vector_job = VectorPruneJob::with_prune_fn( VectorPruneJobConfig::default(), move |age_days, level| { @@ -280,11 +301,9 @@ async fn register_prune_jobs(scheduler: &SchedulerService, db_path: &Path) -> Re } }, ); - register_vector_prune_job(scheduler, vector_job) .await .context("Failed to register vector prune job")?; - info!("Vector prune job registered"); } Err(e) => { @@ -468,6 +487,7 @@ fn build_api_summarizer(settings: &SummarizerSettings) -> Option, db_path: &Path) -> QueryIndexBundle { + use memory_embeddings::EmbeddingModel; use memory_search::{SearchIndex, SearchIndexConfig, TeleportSearcher}; use memory_topics::TopicStorage; use memory_vector::{HnswConfig, HnswIndex, VectorMetadata}; @@ -497,7 +517,8 @@ fn open_query_indexes(storage: &Arc, db_path: &Path) -> QueryIndexBundl if vector_dir.exists() { match memory_embeddings::CandleEmbedder::load_default() { Ok(embedder) => { - let hnsw_config = HnswConfig::new(384, &vector_dir); + let dim = embedder.info().dimension; + let hnsw_config = HnswConfig::new(dim, &vector_dir); match HnswIndex::open_or_create(hnsw_config) { Ok(hnsw) => { let meta_path = vector_dir.join("metadata"); @@ -620,6 +641,10 @@ pub async fn start_daemon( .await .context("Failed to register compaction job")?; + // Attach live indexes before prune/dedup so they share one HNSW handle. + let mut indexes = open_query_indexes(&storage, &db_path); + indexes.api_summarizer = build_api_summarizer(&settings.summarizer); + // Register indexing job if search index exists // The indexing pipeline processes outbox entries into search indexes if let Err(e) = register_indexing_job(&scheduler, storage.clone(), &db_path).await { @@ -629,7 +654,7 @@ pub async fn start_daemon( // Register lifecycle prune jobs if indexes exist // These jobs prune old documents/vectors based on per-level retention policies - if let Err(e) = register_prune_jobs(&scheduler, &db_path).await { + if let Err(e) = register_prune_jobs(&scheduler, &db_path, indexes.vector.clone()).await { warn!("Prune jobs not fully registered: {}", e); } @@ -640,60 +665,56 @@ pub async fn start_daemon( // Create NoveltyChecker for dedup gate (DEDUP-02, DEDUP-03) let novelty_checker = if settings.dedup.enabled { - match memory_embeddings::CandleEmbedder::load_default() { - Ok(embedder) => { - let adapter = Arc::new(CandleEmbedderAdapter::new(embedder)) + let dim = indexes + .vector + .as_ref() + .map(|vh| vh.get_status().dimension.max(0) as usize) + .filter(|d| *d > 0) + .unwrap_or(memory_embeddings::EMBEDDING_DIM); + match indexes.vector.as_ref() { + Some(vh) => { + let adapter = Arc::new(CandleEmbedderAdapter::from_arc(vh.embedder())) as Arc; let buffer = Arc::new(RwLock::new(InFlightBuffer::new( settings.dedup.buffer_capacity, - 384, + dim, ))); - - // Try to open HNSW index for cross-session dedup (DEDUP-02) - let vector_dir = PathBuf::from(&settings.db_path).join("vector"); - let hnsw_opt = if vector_dir.exists() { - let hnsw_config = memory_vector::HnswConfig::new(384, &vector_dir); - match memory_vector::HnswIndex::open_or_create(hnsw_config) { - Ok(hnsw) => { - info!("HNSW index loaded for cross-session dedup"); - Some(Arc::new(std::sync::RwLock::new(hnsw))) - } - Err(e) => { - warn!("Failed to open HNSW for dedup, using buffer-only: {e}"); - None - } - } - } else { - info!("No vector index found, using buffer-only dedup"); - None - }; - - let has_hnsw = hnsw_opt.is_some(); - let checker = if let Some(hnsw_index) = hnsw_opt { - NoveltyChecker::with_composite_index( - Some(adapter), - buffer, - hnsw_index, - settings.dedup.clone(), - ) - } else { - NoveltyChecker::with_in_flight_buffer( - Some(adapter), - buffer, - settings.dedup.clone(), - ) - }; - + let checker = NoveltyChecker::with_composite_index( + Some(adapter), + buffer, + vh.index_handle(), + settings.dedup.clone(), + ); info!( - "Dedup gate enabled (threshold: {}, buffer: {}, hnsw: {})", - settings.dedup.threshold, settings.dedup.buffer_capacity, has_hnsw + "Dedup gate enabled (threshold: {}, buffer: {}, hnsw: shared)", + settings.dedup.threshold, settings.dedup.buffer_capacity ); Some(Arc::new(checker)) } - Err(e) => { - warn!("Failed to load CandleEmbedder for dedup, disabling: {e}"); - None - } + None => match memory_embeddings::CandleEmbedder::load_default() { + Ok(embedder) => { + let adapter = Arc::new(CandleEmbedderAdapter::new(embedder)) + as Arc; + let buffer = Arc::new(RwLock::new(InFlightBuffer::new( + settings.dedup.buffer_capacity, + dim, + ))); + let checker = NoveltyChecker::with_in_flight_buffer( + Some(adapter), + buffer, + settings.dedup.clone(), + ); + info!( + "Dedup gate enabled (threshold: {}, buffer: {}, hnsw: false)", + settings.dedup.threshold, settings.dedup.buffer_capacity + ); + Some(Arc::new(checker)) + } + Err(e) => { + warn!("Failed to load CandleEmbedder for dedup, disabling: {e}"); + None + } + }, } } else { tracing::debug!("Dedup gate disabled by config"); @@ -745,9 +766,7 @@ pub async fn start_daemon( settings.staleness.max_penalty ); - // Attach live indexes so RouteQuery is not agentic-only. - let mut indexes = open_query_indexes(&storage, &db_path); - indexes.api_summarizer = build_api_summarizer(&settings.summarizer); + // Indexes already opened above so prune/dedup share the HNSW handle. // Start server with scheduler let result = run_server_with_scheduler( @@ -1766,8 +1785,8 @@ fn handle_index_stats( if vector_path.exists() { // Try to get dimension from an existing index - // Default to 384 for all-MiniLM-L6-v2 - let dimension = 384; + // Use the embedder's native dimension (all-MiniLM-L6-v2 = EMBEDDING_DIM). + let dimension = memory_embeddings::EMBEDDING_DIM; let hnsw_config = HnswConfig::new(dimension, vector_path); match HnswIndex::open_or_create(hnsw_config) { @@ -2127,6 +2146,7 @@ async fn teleport_search(query: &str, doc_type: &str, limit: usize, addr: &str) let type_str = match result.doc_type { 1 => "TOC", 2 => "Grip", + 3 => "Event", _ => "?", }; diff --git a/crates/memory-embeddings/src/lib.rs b/crates/memory-embeddings/src/lib.rs index fb6068f..c8ad5f9 100644 --- a/crates/memory-embeddings/src/lib.rs +++ b/crates/memory-embeddings/src/lib.rs @@ -21,7 +21,7 @@ pub mod candle; pub mod error; pub mod model; -pub use crate::candle::CandleEmbedder; +pub use crate::candle::{CandleEmbedder, EMBEDDING_DIM}; pub use cache::{get_or_download_model, ModelCache, ModelPaths, DEFAULT_MODEL_REPO, MODEL_FILES}; pub use error::EmbeddingError; pub use model::{Embedding, EmbeddingModel, ModelInfo}; diff --git a/crates/memory-indexing/src/bm25_updater.rs b/crates/memory-indexing/src/bm25_updater.rs index a244092..4a14a66 100644 --- a/crates/memory-indexing/src/bm25_updater.rs +++ b/crates/memory-indexing/src/bm25_updater.rs @@ -96,10 +96,9 @@ impl Bm25IndexUpdater { } } - if let Some(grip) = self.find_grip_for_event(&entry.event_id)? { - self.index_grip(&grip)?; - indexed = true; - } + // Grips are indexed by rebuild / index_grip_direct when the TOC + // rollup creates them. Do not scan every grip per outbox entry + // (that is O(grips) per event and dominates drain time). if indexed { Ok(true) @@ -116,14 +115,6 @@ impl Bm25IndexUpdater { } } - /// Find a grip that references this event (start or end id). - fn find_grip_for_event(&self, event_id: &str) -> Result, IndexingError> { - let grips = crate::rebuild::iter_all_grips(&self.storage)?; - Ok(grips - .into_iter() - .find(|g| g.event_id_start == event_id || g.event_id_end == event_id)) - } - /// Process a batch of outbox entries. pub fn process_batch( &self, diff --git a/crates/memory-orchestrator/src/orchestrator.rs b/crates/memory-orchestrator/src/orchestrator.rs index bf90460..37f9a6d 100644 --- a/crates/memory-orchestrator/src/orchestrator.rs +++ b/crates/memory-orchestrator/src/orchestrator.rs @@ -9,8 +9,7 @@ use std::time::Instant; use anyhow::Result; use memory_retrieval::{ - CapabilityTier, ExecutionMode, FallbackChain, LayerExecutor, RetrievalExecutor, RetrievalLayer, - SearchResult, StopConditions, + ExecutionMode, LayerExecutor, RetrievalLayer, SearchResult, StopConditions, }; use crate::context_builder::ContextBuilder; @@ -28,7 +27,7 @@ pub struct OrchestratorOutput { pub fusion_stage: &'static str, /// Reranker that ran (`"heuristic"` or `"llm"`). pub rerank_mode: String, - /// Layers that returned at least one hit. + /// Layers that were supported and invoked (empty hits still count). pub layers_attempted: Vec, /// Wall-clock milliseconds for the pipeline. pub retrieval_ms: u64, @@ -97,10 +96,24 @@ impl MemoryOrchestrator { /// Execute expand → fan-out → RRF → rerank and return ranked hits. /// - /// Used by `RouteQuery` so gRPC callers get the orchestrator pipeline - /// without going through `MemoryContext` assembly. + /// Fan-out runs supported layers concurrently (independent lists for + /// rank fusion). Use [`Self::query_ranked_with`] to honor client stop + /// conditions and execution mode. pub async fn query_ranked(&self, query: &str) -> Result { + self.query_ranked_with(query, &StopConditions::default(), ExecutionMode::Parallel) + .await + } + + /// Like [`Self::query_ranked`] with caller-supplied stop conditions and mode. + pub async fn query_ranked_with( + &self, + query: &str, + conditions: &StopConditions, + mode: ExecutionMode, + ) -> Result { let start = Instant::now(); + let timeout = conditions.timeout(); + let limit = self.config.top_k.min(conditions.max_nodes as usize).max(1); let queries = if self.config.expand_query { expand_query(query) @@ -108,46 +121,78 @@ impl MemoryOrchestrator { vec![query.to_string()] }; - let layers = [ + let all_layers = [ RetrievalLayer::Topics, RetrievalLayer::Vector, RetrievalLayer::BM25, RetrievalLayer::Agentic, ]; + let layers: Vec = all_layers + .into_iter() + .filter(|&layer| self.executor.supports(layer)) + .collect(); - let re = RetrievalExecutor::new(self.executor.clone()); let mut all_lists: Vec> = Vec::new(); let mut layers_attempted: Vec = Vec::new(); - for q in &queries { - for &layer in &layers { - let chain = FallbackChain { - layers: vec![layer], - merge_results: false, - max_layers: 1, - }; - let conds = StopConditions::default(); - let result = re - .execute( - q, - chain, - &conds, - ExecutionMode::Sequential, - CapabilityTier::Full, - ) - .await; - if !layers_attempted.contains(&layer) { - layers_attempted.push(layer); + let parallel = !matches!(mode, ExecutionMode::Sequential); + + if parallel { + let mut tasks = Vec::new(); + for q in &queries { + for &layer in &layers { + let exec = Arc::clone(&self.executor); + let q = q.clone(); + tasks.push(async move { + match exec.execute(&q, layer, limit).await { + Ok(results) => (layer, results), + Err(e) => { + tracing::debug!(layer = ?layer, error = %e, "layer failed; fail-open"); + (layer, Vec::new()) + } + } + }); } - if result.has_results() { - all_lists.push(result.results); + } + match tokio::time::timeout(timeout, futures::future::join_all(tasks)).await { + Ok(pairs) => { + for (layer, results) in pairs { + if !layers_attempted.contains(&layer) { + layers_attempted.push(layer); + } + if !results.is_empty() { + all_lists.push(results); + } + } + } + Err(_) => { + tracing::warn!("orchestrator fan-out timed out"); + } + } + } else { + 'fanout: for q in &queries { + for &layer in &layers { + if start.elapsed() >= timeout { + tracing::warn!("orchestrator sequential fan-out hit stop timeout"); + break 'fanout; + } + if !layers_attempted.contains(&layer) { + layers_attempted.push(layer); + } + match self.executor.execute(q, layer, limit).await { + Ok(results) if !results.is_empty() => all_lists.push(results), + Ok(_) => {} + Err(e) => { + tracing::debug!(layer = ?layer, error = %e, "layer failed; fail-open"); + } + } } - // fail-open: skip empty/failed layers silently } } let fused = fuse(all_lists, self.config.fusion_k); let reranked = self.reranker.rerank(query, fused).await?; + let rerank_mode = self.reranker.mode_name().to_string(); let results: Vec = reranked .into_iter() @@ -158,12 +203,6 @@ impl MemoryOrchestrator { }) .collect(); - let rerank_mode = match self.config.rerank_mode { - crate::types::RerankMode::Heuristic => "heuristic", - crate::types::RerankMode::Llm => "llm", - } - .to_string(); - Ok(OrchestratorOutput { results, fusion_stage: "rank_fusion", @@ -218,6 +257,10 @@ mod tests { out.reverse(); Ok(out) } + + fn mode_name(&self) -> &'static str { + "llm" + } } #[tokio::test] @@ -323,4 +366,154 @@ mod tests { assert!(!output.results.is_empty()); assert!(output.layers_attempted.contains(&RetrievalLayer::BM25)); } + + #[tokio::test] + async fn test_layers_attempted_omits_unconfigured_layers() { + let executor = MockLayerExecutor::default().with_results( + RetrievalLayer::BM25, + vec![mock_result("doc-x", 0.7, RetrievalLayer::BM25)], + ); + let orch = MemoryOrchestrator::new(Arc::new(executor), OrchestratorConfig::default()); + let output = orch.query_ranked("test").await.unwrap(); + assert_eq!(output.layers_attempted, vec![RetrievalLayer::BM25]); + } + + #[tokio::test] + async fn test_parallel_fan_out_runs_layers_concurrently() { + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::time::Duration; + + struct CountingExecutor { + in_flight: AtomicUsize, + max_in_flight: Arc, + } + + #[async_trait] + impl LayerExecutor for CountingExecutor { + async fn execute( + &self, + _query: &str, + layer: RetrievalLayer, + _limit: usize, + ) -> Result, String> { + let now = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1; + self.max_in_flight.fetch_max(now, Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(40)).await; + self.in_flight.fetch_sub(1, Ordering::SeqCst); + Ok(vec![mock_result("doc", 0.5, layer)]) + } + + fn supports(&self, layer: RetrievalLayer) -> bool { + matches!( + layer, + RetrievalLayer::BM25 | RetrievalLayer::Vector | RetrievalLayer::Topics + ) + } + } + + let max_in_flight = Arc::new(AtomicUsize::new(0)); + let executor = CountingExecutor { + in_flight: AtomicUsize::new(0), + max_in_flight: Arc::clone(&max_in_flight), + }; + let orch = MemoryOrchestrator::new(Arc::new(executor), OrchestratorConfig::default()); + let output = orch + .query_ranked_with("test", &StopConditions::default(), ExecutionMode::Parallel) + .await + .unwrap(); + assert!( + output.layers_attempted.len() >= 3, + "expected BM25+Vector+Topics, got {:?}", + output.layers_attempted + ); + assert!( + max_in_flight.load(Ordering::SeqCst) >= 2, + "parallel fan-out must overlap at least two layers (max in-flight {})", + max_in_flight.load(Ordering::SeqCst) + ); + } + + #[tokio::test] + async fn test_sequential_timeout_stops_remaining_layers() { + use std::time::Duration; + + let executor = MockLayerExecutor::default() + .with_results( + RetrievalLayer::BM25, + vec![mock_result("doc-a", 0.9, RetrievalLayer::BM25)], + ) + .with_delay(RetrievalLayer::BM25, Duration::from_millis(80)) + .with_results( + RetrievalLayer::Vector, + vec![mock_result("doc-b", 0.8, RetrievalLayer::Vector)], + ) + .with_delay(RetrievalLayer::Vector, Duration::from_millis(80)); + + let orch = MemoryOrchestrator::new(Arc::new(executor), OrchestratorConfig::default()); + let output = orch + .query_ranked_with( + "test", + &StopConditions::with_timeout(Duration::from_millis(30)), + ExecutionMode::Sequential, + ) + .await + .unwrap(); + // Layer order is Topics, Vector, BM25, Agentic. Vector is first supported + // layer; its delay trips the timeout so BM25 must not run. + assert_eq!(output.layers_attempted, vec![RetrievalLayer::Vector]); + assert_eq!(output.results.len(), 1); + assert_eq!(output.results[0].doc_id, "doc-b"); + } + + #[tokio::test] + async fn test_mode_override_sequential_is_honored() { + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::time::Duration; + + struct CountingExecutor { + in_flight: AtomicUsize, + max_in_flight: Arc, + } + + #[async_trait] + impl LayerExecutor for CountingExecutor { + async fn execute( + &self, + _query: &str, + layer: RetrievalLayer, + _limit: usize, + ) -> Result, String> { + let now = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1; + self.max_in_flight.fetch_max(now, Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(25)).await; + self.in_flight.fetch_sub(1, Ordering::SeqCst); + Ok(vec![mock_result("doc", 0.5, layer)]) + } + + fn supports(&self, layer: RetrievalLayer) -> bool { + matches!(layer, RetrievalLayer::BM25 | RetrievalLayer::Vector) + } + } + + let max_in_flight = Arc::new(AtomicUsize::new(0)); + let executor = CountingExecutor { + in_flight: AtomicUsize::new(0), + max_in_flight: Arc::clone(&max_in_flight), + }; + let orch = MemoryOrchestrator::new(Arc::new(executor), OrchestratorConfig::default()); + let output = orch + .query_ranked_with( + "test", + &StopConditions::default(), + ExecutionMode::Sequential, + ) + .await + .unwrap(); + assert_eq!(output.layers_attempted.len(), 2); + assert_eq!( + max_in_flight.load(Ordering::SeqCst), + 1, + "sequential mode must never overlap layer executions" + ); + } } diff --git a/crates/memory-orchestrator/src/rerank.rs b/crates/memory-orchestrator/src/rerank.rs index 7879fe3..d3e741b 100644 --- a/crates/memory-orchestrator/src/rerank.rs +++ b/crates/memory-orchestrator/src/rerank.rs @@ -8,6 +8,7 @@ //! — never warn-and-fallback. use std::collections::HashMap; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use anyhow::Result; @@ -57,6 +58,14 @@ impl RerankedResult { pub trait Reranker: Send + Sync { /// Rerank fused results, returning a sorted and potentially trimmed list. async fn rerank(&self, query: &str, results: Vec) -> Result>; + + /// Label of the strategy that actually produced the last `rerank` output. + /// + /// LLM fail-open reports `"heuristic"` so explainability cannot claim a + /// rerank that did not run. + fn mode_name(&self) -> &'static str { + "heuristic" + } } /// LLM text completion used by [`LlmReranker`]. @@ -118,6 +127,7 @@ impl Reranker for HeuristicReranker { pub struct LlmReranker { completer: Arc, max_results: usize, + fell_back: AtomicBool, } impl LlmReranker { @@ -126,6 +136,7 @@ impl LlmReranker { Self { completer, max_results: max_results.max(1), + fell_back: AtomicBool::new(false), } } @@ -151,24 +162,19 @@ Include every doc_id exactly once, most relevant first."# ) } - fn parse_order(text: &str, fallback: &[FusedResult]) -> Vec { + fn parse_order(text: &str) -> Option> { let json_str = extract_json_object(text); - let parsed: serde_json::Value = match serde_json::from_str(&json_str) { - Ok(v) => v, - Err(_) => { - return fallback.iter().map(|r| r.inner.doc_id.clone()).collect(); - } - }; - if let Some(arr) = parsed.get("order").and_then(|v| v.as_array()) { - let ids: Vec = arr - .iter() - .filter_map(|v| v.as_str().map(ToOwned::to_owned)) - .collect(); - if !ids.is_empty() { - return ids; - } + let parsed: serde_json::Value = serde_json::from_str(&json_str).ok()?; + let arr = parsed.get("order")?.as_array()?; + let ids: Vec = arr + .iter() + .filter_map(|v| v.as_str().map(ToOwned::to_owned)) + .collect(); + if ids.is_empty() { + None + } else { + Some(ids) } - fallback.iter().map(|r| r.inner.doc_id.clone()).collect() } } @@ -185,7 +191,16 @@ fn extract_json_object(text: &str) -> String { #[async_trait] impl Reranker for LlmReranker { + fn mode_name(&self) -> &'static str { + if self.fell_back.load(Ordering::Relaxed) { + "heuristic" + } else { + "llm" + } + } + async fn rerank(&self, query: &str, results: Vec) -> Result> { + self.fell_back.store(false, Ordering::Relaxed); if results.is_empty() { return Ok(Vec::new()); } @@ -194,12 +209,19 @@ impl Reranker for LlmReranker { Ok(text) => text, Err(e) => { tracing::warn!(error = %e, "LLM rerank failed; keeping RRF order"); + self.fell_back.store(true, Ordering::Relaxed); return HeuristicReranker::new(self.max_results) .rerank(query, results) .await; } }; - let order = Self::parse_order(&response, &results); + let Some(order) = Self::parse_order(&response) else { + tracing::warn!("LLM rerank returned unparseable order; keeping RRF order"); + self.fell_back.store(true, Ordering::Relaxed); + return HeuristicReranker::new(self.max_results) + .rerank(query, results) + .await; + }; let mut by_id: HashMap = results .into_iter() .map(|r| (r.inner.doc_id.clone(), r)) @@ -237,6 +259,10 @@ pub struct CrossEncoderReranker; #[async_trait] impl Reranker for CrossEncoderReranker { + fn mode_name(&self) -> &'static str { + "cross-encoder" + } + async fn rerank( &self, _query: &str, @@ -311,6 +337,7 @@ mod tests { let reranked = reranker.rerank("q", results).await.unwrap(); assert_eq!(reranked[0].doc_id, "b"); assert_eq!(reranked[1].doc_id, "a"); + assert_eq!(reranker.mode_name(), "llm"); } struct BrokenCompleter; @@ -329,5 +356,28 @@ mod tests { let reranked = reranker.rerank("q", results).await.unwrap(); assert_eq!(reranked[0].doc_id, "a"); assert_eq!(reranked[1].doc_id, "b"); + assert_eq!( + reranker.mode_name(), + "heuristic", + "fail-open must not claim llm ran" + ); + } + + struct GarbageCompleter; + + #[async_trait] + impl Completer for GarbageCompleter { + async fn complete(&self, _prompt: &str) -> Result { + Ok("not json".to_string()) + } + } + + #[tokio::test] + async fn test_llm_reranker_fail_open_on_unparseable_order() { + let results = vec![make_fused("a", 0.9), make_fused("b", 0.5)]; + let reranker = LlmReranker::new(Arc::new(GarbageCompleter), 10); + let reranked = reranker.rerank("q", results).await.unwrap(); + assert_eq!(reranked[0].doc_id, "a"); + assert_eq!(reranker.mode_name(), "heuristic"); } } diff --git a/crates/memory-retrieval/src/executor.rs b/crates/memory-retrieval/src/executor.rs index 549ca5e..355ce42 100644 --- a/crates/memory-retrieval/src/executor.rs +++ b/crates/memory-retrieval/src/executor.rs @@ -678,8 +678,11 @@ impl LayerExecutor for MockLayerExecutor { Ok(results.into_iter().take(limit).collect()) } - fn supports(&self, _layer: RetrievalLayer) -> bool { - true // Mock supports all layers + fn supports(&self, layer: RetrievalLayer) -> bool { + // Only layers that were explicitly configured are "available". + self.results.contains_key(&layer) + || self.fail_layers.contains(&layer) + || self.delays.contains_key(&layer) } } diff --git a/crates/memory-retrieval/src/ranking.rs b/crates/memory-retrieval/src/ranking.rs index 75010fc..9e102e9 100644 --- a/crates/memory-retrieval/src/ranking.rs +++ b/crates/memory-retrieval/src/ranking.rs @@ -24,6 +24,9 @@ pub struct RankingConfig { pub decay_factor: f32, /// Minimum score floor as fraction of original similarity (0.0-1.0). pub score_floor: f32, + /// When true, apply score multipliers but do not re-sort. + /// Used after LLM rerank so salience cannot undo the model order. + pub preserve_order: bool, } impl Default for RankingConfig { @@ -33,6 +36,7 @@ impl Default for RankingConfig { usage_decay_enabled: false, // Off by default until validated decay_factor: 0.15, score_floor: 0.50, + preserve_order: false, } } } @@ -82,12 +86,13 @@ pub fn apply_combined_ranking( result.score = combined.max(floor); } - // Re-sort by adjusted score - results.sort_by(|a, b| { - b.score - .partial_cmp(&a.score) - .unwrap_or(std::cmp::Ordering::Equal) - }); + if !config.preserve_order { + results.sort_by(|a, b| { + b.score + .partial_cmp(&a.score) + .unwrap_or(std::cmp::Ordering::Equal) + }); + } results } @@ -170,6 +175,7 @@ mod tests { usage_decay_enabled: true, decay_factor: 0.15, score_floor: 0.50, + preserve_order: false, }; // Very low salience + high usage: combined would be very low @@ -194,6 +200,7 @@ mod tests { usage_decay_enabled: true, decay_factor: 0.15, score_floor: 0.50, + preserve_order: false, }; let results = vec![make_result("test", 0.8, 0.7, 3)]; @@ -228,4 +235,21 @@ mod tests { "Score should be unchanged when both disabled" ); } + + #[test] + fn test_preserve_order_skips_resort() { + let config = RankingConfig { + salience_enabled: true, + preserve_order: true, + ..Default::default() + }; + let results = vec![ + make_result("low_sal", 0.8, 0.0, 0), + make_result("high_sal", 0.8, 1.0, 0), + ]; + let ranked = apply_combined_ranking(results, &config); + assert_eq!(ranked[0].doc_id, "low_sal"); + assert_eq!(ranked[1].doc_id, "high_sal"); + assert!(ranked[1].score > ranked[0].score); + } } diff --git a/crates/memory-search/src/schema.rs b/crates/memory-search/src/schema.rs index 66253a8..534da0d 100644 --- a/crates/memory-search/src/schema.rs +++ b/crates/memory-search/src/schema.rs @@ -55,7 +55,7 @@ pub struct SearchSchema { pub doc_id: Field, /// TOC level for toc_node: "year", "month", etc. (STRING) pub level: Field, - /// Searchable text: title+bullets for TOC, excerpt for grip (TEXT) + /// Searchable text: title+bullets for TOC, excerpt for grip, body for event (TEXT | STORED) pub text: Field, /// Keywords/tags (TEXT | STORED) pub keywords: Field, @@ -114,7 +114,7 @@ impl SearchSchema { /// - doc_type: STRING | STORED - "toc_node" or "grip" /// - doc_id: STRING | STORED - node_id or grip_id /// - level: STRING - TOC level (for filtering) -/// - text: TEXT - searchable content +/// - text: TEXT | STORED - searchable content (also the LLM-rerank preview) /// - keywords: TEXT | STORED - keywords/tags /// - timestamp_ms: STRING | STORED - for recency info pub fn build_teleport_schema() -> SearchSchema { @@ -129,8 +129,10 @@ pub fn build_teleport_schema() -> SearchSchema { // TOC level (for toc_node only): "year", "month", "week", "day", "segment" let level = schema_builder.add_text_field("level", STRING | STORED); - // Searchable text content (title + bullets for TOC, excerpt for grip) - let text = schema_builder.add_text_field("text", TEXT); + // Searchable text content (title + bullets for TOC, excerpt for grip, body for event). + // STORED so BM25 hits expose a preview — events have empty keywords, so the + // LLM reranker must read this field rather than keywords. + let text = schema_builder.add_text_field("text", TEXT | STORED); // Keywords (indexed and stored for retrieval) let keywords = schema_builder.add_text_field("keywords", TEXT | STORED); @@ -165,6 +167,11 @@ mod tests { assert!(schema.schema.get_field("doc_type").is_ok()); assert!(schema.schema.get_field("doc_id").is_ok()); assert!(schema.schema.get_field("text").is_ok()); + let text_entry = schema.schema.get_field_entry(schema.text); + assert!( + text_entry.is_stored(), + "text must be STORED so event hits have a preview" + ); } #[test] diff --git a/crates/memory-search/src/searcher.rs b/crates/memory-search/src/searcher.rs index c41893e..a14f0f6 100644 --- a/crates/memory-search/src/searcher.rs +++ b/crates/memory-search/src/searcher.rs @@ -21,6 +21,9 @@ pub struct TeleportResult { pub doc_type: DocType, /// BM25 relevance score pub score: f32, + /// Indexed body text (truncated). Events have no keywords, so callers + /// must use this field — not `keywords` — as the preview. + pub text: String, /// Keywords from the document (if stored) pub keywords: Option, /// Timestamp in milliseconds @@ -152,6 +155,16 @@ impl TeleportSearcher { .map(|s| s.to_string()) .filter(|s| !s.is_empty()); + let text_raw = doc + .get_first(self.schema.text) + .and_then(|v| v.as_str()) + .unwrap_or(""); + let text = if text_raw.is_empty() { + keywords.clone().unwrap_or_default() + } else { + truncate_preview(text_raw, 240) + }; + let timestamp_ms = doc .get_first(self.schema.timestamp_ms) .and_then(|v| v.as_str()) @@ -169,6 +182,7 @@ impl TeleportSearcher { doc_id, doc_type, score, + text, keywords, timestamp_ms, agent, @@ -290,6 +304,16 @@ impl TeleportSearcher { unsafe impl Send for TeleportSearcher {} unsafe impl Sync for TeleportSearcher {} +fn truncate_preview(s: &str, max_chars: usize) -> String { + let mut iter = s.chars(); + let preview: String = iter.by_ref().take(max_chars).collect(); + if iter.next().is_some() { + format!("{preview}...") + } else { + preview + } +} + #[cfg(test)] mod tests { use super::*; @@ -503,6 +527,34 @@ mod tests { assert!(results[0].keywords.is_some()); } + #[test] + fn test_event_search_returns_text_preview() { + use memory_types::{Event, EventRole, EventType}; + + let (_temp_dir, index) = setup_index(); + let indexer = SearchIndexer::new(&index).unwrap(); + indexer + .index_event(&Event::new( + "evt-1".into(), + "s".into(), + Utc::now(), + EventType::UserMessage, + EventRole::User, + "zebra unique token event body".into(), + )) + .unwrap(); + indexer.commit().unwrap(); + + let searcher = TeleportSearcher::new(&index).unwrap(); + let results = searcher + .search("zebra unique", SearchOptions::new().with_limit(10)) + .unwrap(); + assert_eq!(results.len(), 1); + assert_eq!(results[0].doc_id, "evt-1"); + assert!(results[0].text.contains("zebra unique token")); + assert!(results[0].keywords.is_none()); + } + #[test] fn test_reload() { let (_temp_dir, index) = setup_index(); diff --git a/crates/memory-service/src/hybrid.rs b/crates/memory-service/src/hybrid.rs index e483149..a330d1b 100644 --- a/crates/memory-service/src/hybrid.rs +++ b/crates/memory-service/src/hybrid.rs @@ -149,7 +149,7 @@ impl HybridSearchHandler { doc_id: r.doc_id, doc_type: r.doc_type.as_str().to_string(), score: r.score, - text_preview: r.keywords.unwrap_or_default(), + text_preview: r.text, timestamp_ms: r.timestamp_ms.unwrap_or(0), agent: r.agent, }) diff --git a/crates/memory-service/src/lib.rs b/crates/memory-service/src/lib.rs index 06b6667..e50c6b1 100644 --- a/crates/memory-service/src/lib.rs +++ b/crates/memory-service/src/lib.rs @@ -26,6 +26,9 @@ pub mod topics; pub mod vector; pub mod pb { + // tonic::Status is ≥176 bytes; generated gRPC stubs trip clippy::result_large_err + // (Rust 1.98+). We do not control the generated signatures. + #![allow(clippy::result_large_err)] tonic::include_proto!("memory"); pub const FILE_DESCRIPTOR_SET: &[u8] = tonic::include_file_descriptor_set!("memory_descriptor"); diff --git a/crates/memory-service/src/novelty.rs b/crates/memory-service/src/novelty.rs index 2d7a738..b84e9f5 100644 --- a/crates/memory-service/src/novelty.rs +++ b/crates/memory-service/src/novelty.rs @@ -119,6 +119,11 @@ impl CandleEmbedderAdapter { embedder: Arc::new(embedder), } } + + /// Wrap an already-shared embedder (query path / prune job). + pub fn from_arc(embedder: Arc) -> Self { + Self { embedder } + } } #[async_trait::async_trait] diff --git a/crates/memory-service/src/retrieval.rs b/crates/memory-service/src/retrieval.rs index b87c7be..196e182 100644 --- a/crates/memory-service/src/retrieval.rs +++ b/crates/memory-service/src/retrieval.rs @@ -264,10 +264,9 @@ impl RetrievalHandler { return Err(Status::invalid_argument("Query is required")); } - // Get stop conditions (honored by the fallback executor inside the orchestrator - // via StopConditions::default today; parsed here so the proto field is not dropped - // silently if we thread it through later). - let _stop_conditions = req + // Get stop conditions and execution mode — both are forwarded into + // the orchestrator so the proto fields are not echo-only. + let stop_conditions = req .stop_conditions .map(|sc| proto_to_stop_conditions(&sc)) .unwrap_or_default(); @@ -314,8 +313,14 @@ impl RetrievalHandler { )); let requested_rerank = match req.rerank_mode.as_deref() { + None | Some("") => RerankMode::Heuristic, + Some(mode) if mode.eq_ignore_ascii_case("heuristic") => RerankMode::Heuristic, Some(mode) if mode.eq_ignore_ascii_case("llm") => RerankMode::Llm, - _ => RerankMode::Heuristic, + Some(other) => { + return Err(Status::invalid_argument(format!( + "unknown rerank_mode '{other}'; expected 'heuristic' or 'llm'" + ))); + } }; let (rerank_mode, reranker): (RerankMode, Box) = match requested_rerank { RerankMode::Llm => { @@ -348,7 +353,7 @@ impl RetrievalHandler { }; let orchestrator = MemoryOrchestrator::with_reranker(executor, orch_config, reranker); let output = orchestrator - .query_ranked(&req.query) + .query_ranked_with(&req.query, &stop_conditions, mode) .await .map_err(|e| Status::internal(format!("orchestrator error: {e}")))?; @@ -373,8 +378,12 @@ impl RetrievalHandler { enriched_results }; - // Apply combined ranking (salience + usage decay) after stale filter - let ranking_config = RankingConfig::default(); + // Apply combined ranking (salience + usage decay) after stale filter. + // LLM order is preserved: salience may adjust scores but must not re-sort. + let ranking_config = RankingConfig { + preserve_order: orch_rerank_mode == "llm", + ..Default::default() + }; let ranked_results = apply_combined_ranking(filtered_results, &ranking_config); let total_time_ms = start.elapsed().as_millis() as u64; @@ -608,7 +617,7 @@ impl LayerExecutor for SimpleLayerExecutor { doc_id: r.doc_id, doc_type: format!("{:?}", r.doc_type).to_lowercase(), score: r.score, - text_preview: r.keywords.unwrap_or_default(), + text_preview: r.text, source_layer: CrateLayer::BM25, metadata: build_metadata( r.timestamp_ms, @@ -676,7 +685,7 @@ impl LayerExecutor for SimpleLayerExecutor { doc_id: r.doc_id, doc_type: format!("{:?}", r.doc_type).to_lowercase(), score: r.score, - text_preview: r.keywords.unwrap_or_default(), + text_preview: r.text, source_layer: CrateLayer::BM25, metadata: build_metadata( r.timestamp_ms, @@ -1359,5 +1368,118 @@ mod tests { ); // First of LLM should be last of heuristic (full reverse of two hits). assert_eq!(llm_ids.first(), heuristic_ids.last()); + assert!( + llm_resp.results.iter().all(|r| !r.text_preview.is_empty()), + "BM25 event hits must carry a text preview for LLM rerank" + ); + } + + #[tokio::test] + async fn test_unknown_rerank_mode_is_rejected() { + let (handler, _temp) = create_test_handler(); + let err = handler + .route_query(Request::new(RouteQueryRequest { + query: "what is rust?".to_string(), + intent_override: None, + stop_conditions: None, + mode_override: None, + limit: 10, + agent_filter: None, + all_projects: false, + rerank_mode: Some("cross-encoder".into()), + expand_query: false, + })) + .await + .unwrap_err(); + assert_eq!(err.code(), tonic::Code::InvalidArgument); + assert!(err.message().contains("rerank_mode")); + } + + #[tokio::test] + async fn test_layers_attempted_omits_unavailable_indexes() { + let (handler, _temp) = create_test_handler(); + let resp = handler + .route_query(Request::new(RouteQueryRequest { + query: "what is rust?".to_string(), + intent_override: None, + stop_conditions: None, + mode_override: None, + limit: 10, + agent_filter: None, + all_projects: false, + rerank_mode: None, + expand_query: false, + })) + .await + .unwrap() + .into_inner(); + // No BM25/vector/topics attached: only Agentic is supported. + assert_eq!(resp.layers_attempted, vec![ProtoLayer::Agentic as i32]); + } + + struct BrokenCompleter; + + #[async_trait] + impl Completer for BrokenCompleter { + async fn complete(&self, _prompt: &str) -> anyhow::Result { + anyhow::bail!("network down"); + } + } + + #[tokio::test] + async fn test_llm_fail_open_reports_heuristic() { + use chrono::Utc; + use memory_search::{SearchIndex, SearchIndexConfig, SearchIndexer}; + use memory_types::{Event, EventRole, EventType}; + + let temp_dir = TempDir::new().unwrap(); + let index_path = temp_dir.path().join("search"); + std::fs::create_dir_all(&index_path).unwrap(); + let index = SearchIndex::open_or_create(SearchIndexConfig::new(&index_path)).unwrap(); + let indexer = SearchIndexer::new(&index).unwrap(); + indexer + .index_event(&Event::new( + "alpha-event".into(), + "s".into(), + Utc::now(), + EventType::UserMessage, + EventRole::User, + "zebra unique token alpha document".into(), + )) + .unwrap(); + indexer.commit().unwrap(); + let searcher = Arc::new(TeleportSearcher::new(&index).unwrap()); + let storage = Arc::new(Storage::open(temp_dir.path()).unwrap()); + let handler = RetrievalHandler::with_services( + storage, + Some(searcher), + None, + None, + Default::default(), + ) + .with_completer(Arc::new(BrokenCompleter)); + + let resp = handler + .route_query(Request::new(RouteQueryRequest { + query: "zebra unique token".to_string(), + intent_override: None, + stop_conditions: None, + mode_override: None, + limit: 10, + agent_filter: None, + all_projects: false, + rerank_mode: Some("llm".into()), + expand_query: false, + })) + .await + .unwrap() + .into_inner(); + assert_eq!( + resp.explanation + .as_ref() + .and_then(|e| e.rerank_mode.as_deref()), + Some("heuristic"), + "completer failure must not report rerank=llm" + ); } } diff --git a/crates/memory-service/src/teleport_service.rs b/crates/memory-service/src/teleport_service.rs index b7cb4e9..1d991d3 100644 --- a/crates/memory-service/src/teleport_service.rs +++ b/crates/memory-service/src/teleport_service.rs @@ -63,7 +63,10 @@ pub async fn handle_teleport_search( DocType::Event => TeleportDocType::Event as i32, }, score: r.score, - keywords: r.keywords, + keywords: r + .keywords + .filter(|k| !k.is_empty()) + .or_else(|| (!r.text.is_empty()).then_some(r.text)), timestamp_ms: r.timestamp_ms, agent: r.agent, }) diff --git a/crates/memory-service/src/vector.rs b/crates/memory-service/src/vector.rs index d26510a..a4d69f9 100644 --- a/crates/memory-service/src/vector.rs +++ b/crates/memory-service/src/vector.rs @@ -40,6 +40,21 @@ impl VectorTeleportHandler { } } + /// Shared HNSW handle (same lock the prune job must mutate). + pub fn index_handle(&self) -> Arc> { + Arc::clone(&self.index) + } + + /// Shared metadata store handle. + pub fn metadata_arc(&self) -> Arc { + Arc::clone(&self.metadata) + } + + /// Shared embedder. + pub fn embedder(&self) -> Arc { + Arc::clone(&self.embedder) + } + /// Get a reference to the vector metadata store. pub fn metadata(&self) -> &Arc { &self.metadata diff --git a/crates/memory-vector/src/hnsw.rs b/crates/memory-vector/src/hnsw.rs index c0125ba..6210378 100644 --- a/crates/memory-vector/src/hnsw.rs +++ b/crates/memory-vector/src/hnsw.rs @@ -36,7 +36,7 @@ pub struct HnswConfig { impl Default for HnswConfig { fn default() -> Self { Self { - dimension: 384, // all-MiniLM-L6-v2 + dimension: memory_embeddings::EMBEDDING_DIM, connectivity: 16, expansion_add: 200, expansion_search: 100, diff --git a/docs/plans/phase-54.5-truth-leaks-plan.md b/docs/plans/phase-54.5-truth-leaks-plan.md new file mode 100644 index 0000000..3631de6 --- /dev/null +++ b/docs/plans/phase-54.5-truth-leaks-plan.md @@ -0,0 +1,37 @@ +# Phase 54.5 — Truth leaks + CI toolchain pin + +**Date:** 2026-08-30 +**Status:** Implemented (PR pending) +**Depends on:** Phase 54 (merged #32), Phase 55 (merged #33) +**Branch:** `feature/phase-54.5-truth-leaks` + +## Why + +Phase 54's public claims were true. A review of the wiring found second-order +honesty bugs in the explainability payload and two landmines that would +contaminate Phase 55 measurements. Main is also red on Clippy because CI +tracks floating `stable` (Rust 1.98 `result_large_err` on generated tonic stubs). + +## Scope + +1. **CI:** `#[allow(clippy::result_large_err)]` on `include_proto!`; pin + `rust-toolchain.toml` to 1.97; CI/release/e2e workflows honor the pin. +2. **Explainability tells the truth** + - LLM fail-open (completer error or unparseable order) reports `rerank=heuristic` + - `layers_attempted` only lists layers that were supported and invoked + - `stop_conditions` and `mode_override` are forwarded into the orchestrator + - unknown `rerank_mode` is `InvalidArgument` (CLI clap-rejects too) +3. **LLM rerank actually reranks** + - BM25 hits expose indexed `text` (events have empty keywords) + - salience ranking `preserve_order` after a successful LLM rerank +4. **Don't plant Phase 55 landmines** + - drop per-event full grip scan on BM25 outbox drain + - fusion fan-out is concurrent (`join_all`) unless mode is Sequential + - query path, prune job, and dedup share one HNSW `Arc>` + - dimension comes from the embedder / `EMBEDDING_DIM`, not a magic 384 + +## Non-goals + +- Phase 55/56 work +- Implementing cross-encoder reranking +- Background daemonization diff --git a/rust-toolchain.toml b/rust-toolchain.toml new file mode 100644 index 0000000..33161a1 --- /dev/null +++ b/rust-toolchain.toml @@ -0,0 +1,8 @@ +# Pin the toolchain so a new stable Clippy lint cannot redden main overnight. +# Bump deliberately (and run clippy) rather than tracking floating `stable`. +# CI must pass the same version as `toolchain:` to dtolnay/rust-toolchain@master +# — that action does not read this file. +[toolchain] +channel = "1.97" +components = ["clippy", "rustfmt"] +profile = "minimal"