Skip to main content

khive_runtime/
fusion.rs

1//! Fusion strategies for combining ranked result lists.
2
3use std::collections::{hash_map::Entry, HashMap, HashSet};
4
5use uuid::Uuid;
6
7use khive_score::DeterministicScore;
8use khive_storage::types::{
9    PageRequest, TextFilter, TextQueryMode, TextSearchHit, TextSearchRequest, VectorSearchHit,
10};
11use khive_storage::EntityFilter;
12use khive_types::SubstrateKind;
13
14use crate::error::{RuntimeError, RuntimeResult};
15use crate::retrieval::{RankScoreKind, SearchHit, SearchSignals, SearchSource};
16use crate::runtime::{KhiveRuntime, NamespaceToken};
17
18pub use khive_fusion::FusionStrategy;
19
20/// A single ranked candidate stream fed into a [`FusionExecutor`] โ€” the
21/// entity/note ID keyed shape used throughout hybrid search (ADR-012
22/// `FusionStrategy::Custom` ยง"strategy executor").
23pub type CandidateStream = Vec<(Uuid, DeterministicScore)>;
24
25/// One fused, ranked `(id, score)` pair returned by a [`FusionExecutor`].
26pub type RankedHit = (Uuid, DeterministicScore);
27
28/// Runtime-registered custom fusion strategy (ADR-012).
29///
30/// Packs implement this to plug a strategy into `FusionStrategy::Custom { name,
31/// .. }` via [`KhiveRuntime::register_fusion_strategy`] โ€” the seam a
32/// learned-sparse (SPLADE) retrieval leg plugs into. Async so an executor can
33/// perform I/O (e.g. a decay/posterior lookup) while fusing, and fallible so
34/// it can reject malformed `params` instead of degrading silently.
35#[async_trait::async_trait]
36pub trait FusionExecutor: Send + Sync + 'static {
37    /// Declare the strategy represented by the executor's ordering score.
38    fn rank_score_kind(&self) -> RankScoreKind;
39
40    /// Combine `streams` into a single ranked list, honoring `limit` as a
41    /// hint (the dispatch boundary re-sorts and truncates the result with the
42    /// crate's canonical comparator regardless, so an executor need not sort
43    /// or truncate defensively itself).
44    async fn fuse(
45        &self,
46        streams: Vec<CandidateStream>,
47        params: &serde_json::Value,
48        limit: usize,
49    ) -> RuntimeResult<Vec<RankedHit>>;
50}
51
52const CANDIDATE_MULTIPLIER: u32 = 4;
53
54/// RRF convenience wrapper used by operations.rs (k=60 note search path).
55pub(crate) async fn rrf_fuse_k(
56    rt: &KhiveRuntime,
57    text_hits: Vec<TextSearchHit>,
58    vector_hits: Vec<VectorSearchHit>,
59    k: usize,
60    limit: usize,
61) -> RuntimeResult<Vec<SearchHit>> {
62    rt.fuse_with_strategy(text_hits, vector_hits, &FusionStrategy::Rrf { k }, limit)
63        .await
64}
65
66impl KhiveRuntime {
67    /// Fuse text and vector hits using the given strategy, returning at most
68    /// `limit` results. Positional weighted strategies use `[vector, keyword]`
69    /// order.
70    ///
71    /// `FusionStrategy::Custom { name, .. }` is resolved against this
72    /// runtime's registered executors (see
73    /// [`register_fusion_strategy`](KhiveRuntime::register_fusion_strategy)).
74    /// An unregistered name fails closed with
75    /// `RuntimeError::UnknownFusionStrategy` rather than silently falling
76    /// back to RRF.
77    pub(crate) async fn fuse_with_strategy(
78        &self,
79        text_hits: Vec<TextSearchHit>,
80        vector_hits: Vec<VectorSearchHit>,
81        strategy: &FusionStrategy,
82        limit: usize,
83    ) -> RuntimeResult<Vec<SearchHit>> {
84        match strategy {
85            FusionStrategy::VectorOnly => {
86                self.fuse_sources(Vec::new(), vector_hits, strategy, limit)
87                    .await
88            }
89            FusionStrategy::KeywordOnly => {
90                self.fuse_sources(text_hits, Vec::new(), strategy, limit)
91                    .await
92            }
93            FusionStrategy::Rrf { .. }
94            | FusionStrategy::Weighted { .. }
95            | FusionStrategy::Union
96            | FusionStrategy::Custom { .. } => {
97                self.fuse_sources(text_hits, vector_hits, strategy, limit)
98                    .await
99            }
100        }
101    }
102
103    async fn fuse_sources(
104        &self,
105        text_hits: Vec<TextSearchHit>,
106        vector_hits: Vec<VectorSearchHit>,
107        strategy: &FusionStrategy,
108        limit: usize,
109    ) -> RuntimeResult<Vec<SearchHit>> {
110        let mut metadata: HashMap<Uuid, SearchHit> =
111            HashMap::with_capacity(text_hits.len() + vector_hits.len());
112        let prefer_maximum_signal = matches!(
113            strategy,
114            FusionStrategy::Weighted { .. } | FusionStrategy::Union
115        );
116
117        let text_source: Vec<(Uuid, DeterministicScore)> = text_hits
118            .into_iter()
119            .map(|h| {
120                let hit = SearchHit {
121                    entity_id: h.subject_id,
122                    score: h.score,
123                    rank_score_kind: RankScoreKind::Keyword,
124                    signals: SearchSignals {
125                        vector_similarity: None,
126                        keyword_score: Some(h.score),
127                    },
128                    source: SearchSource::Text,
129                    title: h.title,
130                    snippet: h.snippet,
131                };
132                let id = hit.entity_id;
133                let score = hit.score;
134                merge_metadata(&mut metadata, hit, prefer_maximum_signal);
135                (id, score)
136            })
137            .collect();
138
139        let vector_source: Vec<(Uuid, DeterministicScore)> = vector_hits
140            .into_iter()
141            .map(|h| {
142                let hit = SearchHit {
143                    entity_id: h.subject_id,
144                    score: h.score,
145                    rank_score_kind: RankScoreKind::Vector,
146                    signals: SearchSignals {
147                        vector_similarity: Some(h.score),
148                        keyword_score: None,
149                    },
150                    source: SearchSource::Vector,
151                    title: None,
152                    snippet: None,
153                };
154                let id = hit.entity_id;
155                let score = hit.score;
156                merge_metadata(&mut metadata, hit, prefer_maximum_signal);
157                (id, score)
158            })
159            .collect();
160
161        // Canonical positional order is [vector, keyword]. Empty arms remain in
162        // place: removing one would shift the surviving arm onto the wrong weight.
163        let sources: Vec<Vec<(Uuid, DeterministicScore)>> = vec![vector_source, text_source];
164
165        let (rank_score_kind, fused) = self.dispatch_fusion(sources, strategy, limit).await?;
166
167        Ok(fused
168            .into_iter()
169            .filter_map(|(id, score)| {
170                let mut hit = metadata.remove(&id)?;
171                hit.score = score;
172                hit.rank_score_kind = rank_score_kind;
173                Some(hit)
174            })
175            .collect())
176    }
177
178    /// Resolve `strategy` against either the built-in `khive-fusion`
179    /// dispatcher or a registered [`FusionExecutor`], applying the crate's
180    /// canonical score-desc/id-asc ordering at the boundary either way.
181    ///
182    /// `Custom` names are resolved *before* the empty-input/zero-limit short
183    /// circuit, so a misconfigured name errors on every call -- including
184    /// zero-result ones -- rather than being indistinguishable from a valid
185    /// empty result.
186    async fn dispatch_fusion(
187        &self,
188        sources: Vec<Vec<(Uuid, DeterministicScore)>>,
189        strategy: &FusionStrategy,
190        limit: usize,
191    ) -> RuntimeResult<(RankScoreKind, Vec<RankedHit>)> {
192        let rank_score_kind = match strategy {
193            FusionStrategy::Rrf { .. } => RankScoreKind::Rrf,
194            FusionStrategy::VectorOnly => RankScoreKind::Vector,
195            FusionStrategy::KeywordOnly => RankScoreKind::Keyword,
196            FusionStrategy::Weighted { .. } => RankScoreKind::Weighted,
197            FusionStrategy::Union => RankScoreKind::Union,
198            FusionStrategy::Custom { name, params } => {
199                let executor = self.fusion_executor(name)?;
200                let rank_score_kind = executor.rank_score_kind();
201                if limit == 0 || sources.iter().all(Vec::is_empty) {
202                    return Ok((rank_score_kind, Vec::new()));
203                }
204                let mut hits = executor.fuse(sources, params, limit).await?;
205                hits.sort_by(khive_fusion::cmp_desc_then_id);
206                hits.truncate(limit);
207                return Ok((rank_score_kind, hits));
208            }
209        };
210        Ok((
211            rank_score_kind,
212            khive_fusion::fuse(sources, strategy, limit)?,
213        ))
214    }
215}
216
217fn merge_metadata(
218    metadata: &mut HashMap<Uuid, SearchHit>,
219    hit: SearchHit,
220    prefer_maximum_signal: bool,
221) {
222    match metadata.entry(hit.entity_id) {
223        Entry::Occupied(mut entry) => {
224            let existing = entry.get_mut();
225            existing.source = merge_sources(existing.source, hit.source);
226            // RRF and pass-through retain the first occurrence; weighted and
227            // union use the maximum contribution from each retrieval leg.
228            existing.signals.vector_similarity = if prefer_maximum_signal {
229                existing
230                    .signals
231                    .vector_similarity
232                    .max(hit.signals.vector_similarity)
233            } else {
234                existing
235                    .signals
236                    .vector_similarity
237                    .or(hit.signals.vector_similarity)
238            };
239            existing.signals.keyword_score = if prefer_maximum_signal {
240                existing
241                    .signals
242                    .keyword_score
243                    .max(hit.signals.keyword_score)
244            } else {
245                existing.signals.keyword_score.or(hit.signals.keyword_score)
246            };
247            if existing.title.is_none() {
248                existing.title = hit.title;
249            }
250            if existing.snippet.is_none() {
251                existing.snippet = hit.snippet;
252            }
253        }
254        Entry::Vacant(entry) => {
255            entry.insert(hit);
256        }
257    }
258}
259
260fn merge_sources(left: SearchSource, right: SearchSource) -> SearchSource {
261    match (left, right) {
262        (SearchSource::Both, _) | (_, SearchSource::Both) => SearchSource::Both,
263        (SearchSource::Text, SearchSource::Vector) | (SearchSource::Vector, SearchSource::Text) => {
264            SearchSource::Both
265        }
266        (SearchSource::Text, SearchSource::Text) => SearchSource::Text,
267        (SearchSource::Vector, SearchSource::Vector) => SearchSource::Vector,
268    }
269}
270
271impl KhiveRuntime {
272    async fn retain_alive_search_hits(
273        &self,
274        token: &NamespaceToken,
275        mut fused: Vec<SearchHit>,
276        limit: usize,
277    ) -> RuntimeResult<Vec<SearchHit>> {
278        // Filter out soft-deleted entities. A single query fetches all alive IDs from the
279        // fused candidate pool; any ID absent from the result has been soft-deleted.
280        if !fused.is_empty() {
281            let candidate_ids: Vec<Uuid> = fused.iter().map(|h| h.entity_id).collect();
282            let alive_page = self
283                .entities(token)?
284                .query_entities(
285                    token.namespace().as_str(),
286                    EntityFilter {
287                        ids: candidate_ids,
288                        ..EntityFilter::default()
289                    },
290                    PageRequest {
291                        offset: 0,
292                        limit: u32::try_from(fused.len()).unwrap_or(u32::MAX),
293                    },
294                )
295                .await?;
296            let alive: HashSet<Uuid> = alive_page.items.into_iter().map(|e| e.id).collect();
297            fused.retain(|h| alive.contains(&h.entity_id));
298        }
299
300        fused.truncate(limit);
301        Ok(fused)
302    }
303
304    /// Hybrid search with a caller-supplied fusion strategy.
305    ///
306    /// `FusionStrategy::Custom { name, .. }` is resolved against this
307    /// runtime's registered executors (see
308    /// [`register_fusion_strategy`](KhiveRuntime::register_fusion_strategy));
309    /// an unregistered name fails closed with
310    /// `RuntimeError::UnknownFusionStrategy`.
311    pub async fn hybrid_search_with_strategy(
312        &self,
313        token: &NamespaceToken,
314        query_text: &str,
315        query_vector: Option<Vec<f32>>,
316        strategy: FusionStrategy,
317        limit: u32,
318    ) -> RuntimeResult<Vec<SearchHit>> {
319        let candidates = limit.saturating_mul(CANDIDATE_MULTIPLIER).max(limit);
320
321        let text_hits = if matches!(&strategy, FusionStrategy::VectorOnly) {
322            Vec::new()
323        } else {
324            let ns = token.namespace().as_str().to_owned();
325            // sanitize_fts5_query strips known-unsafe metacharacters, but residual
326            // punctuation can still trip the FTS5 parser at runtime; that error must
327            // fail loud rather than silently degrade to vector-only fusion. Errors
328            // from other legs (vector search) still propagate normally.
329            let text_search_result = self
330                .text(token)?
331                .search(TextSearchRequest {
332                    query: query_text.to_string(),
333                    mode: TextQueryMode::Plain,
334                    filter: Some(TextFilter {
335                        namespaces: vec![ns],
336                        ..TextFilter::default()
337                    }),
338                    top_k: candidates,
339                    snippet_chars: 200,
340                })
341                .await;
342            crate::error::fts_text_leg_or_err(
343                text_search_result.map_err(RuntimeError::from),
344                "hybrid_search_with_strategy",
345                query_text,
346            )?
347        };
348
349        let vector_hits = if !matches!(&strategy, FusionStrategy::KeywordOnly)
350            && (query_vector.is_some() || self.config().embedding_model.is_some())
351        {
352            self.vector_search(
353                token,
354                query_vector,
355                Some(query_text),
356                candidates,
357                Some(SubstrateKind::Entity),
358            )
359            .await?
360        } else {
361            Vec::new()
362        };
363
364        // Each arm fetched `candidates` independently, so their union can contain
365        // twice that many distinct IDs. Keep the complete fetched pool through
366        // ranking and the alive check; truncating it first lets stale hits hide
367        // live candidates from the other arm.
368        let fusion_limit = text_hits.len().saturating_add(vector_hits.len());
369        let fused = self
370            .fuse_with_strategy(text_hits, vector_hits, &strategy, fusion_limit)
371            .await?;
372        self.retain_alive_search_hits(token, fused, limit as usize)
373            .await
374    }
375}
376
377#[cfg(test)]
378mod tests {
379    use super::*;
380    use chrono::Utc;
381    use khive_storage::types::{TextDocument, TextSearchHit, VectorSearchHit, VectorSearchRequest};
382    use khive_storage::Entity;
383    use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService};
384    use std::sync::atomic::{AtomicUsize, Ordering};
385    use std::sync::Arc;
386
387    use crate::{EmbedderProvider, RuntimeConfig};
388
389    fn text_hit(id: Uuid, score: f64, title: &str) -> TextSearchHit {
390        TextSearchHit {
391            subject_id: id,
392            score: DeterministicScore::from_f64(score),
393            rank: 1,
394            title: Some(title.to_string()),
395            snippet: Some("...".to_string()),
396        }
397    }
398
399    fn vector_hit(id: Uuid, score: f64) -> VectorSearchHit {
400        VectorSearchHit {
401            subject_id: id,
402            score: DeterministicScore::from_f64(score),
403            rank: 1,
404        }
405    }
406
407    fn evidence_runtime() -> KhiveRuntime {
408        let backend = Arc::new(crate::StorageBackend::memory().expect("in-memory backend"));
409        backend.prepare_core_schema().expect("core schema");
410        KhiveRuntime::from_backend(
411            backend,
412            RuntimeConfig {
413                db_path: None,
414                events_split: None,
415                actor_id: Some("test:fusion-evidence".into()),
416                ..RuntimeConfig::no_embeddings()
417            },
418        )
419    }
420
421    struct CountingEmbeddingService {
422        calls: Arc<AtomicUsize>,
423        dimensions: usize,
424    }
425
426    #[async_trait::async_trait]
427    impl EmbeddingService for CountingEmbeddingService {
428        async fn embed(
429            &self,
430            texts: &[String],
431            _model: EmbeddingModel,
432        ) -> Result<Vec<Vec<f32>>, EmbedError> {
433            self.calls.fetch_add(1, Ordering::SeqCst);
434            Ok(texts.iter().map(|_| vec![1.0; self.dimensions]).collect())
435        }
436
437        fn supports_model(&self, _model: EmbeddingModel) -> bool {
438            true
439        }
440
441        fn name(&self) -> &'static str {
442            "fusion-counting-embedding"
443        }
444    }
445
446    struct CountingEmbedderProvider {
447        name: String,
448        calls: Arc<AtomicUsize>,
449        dimensions: usize,
450    }
451
452    #[async_trait::async_trait]
453    impl EmbedderProvider for CountingEmbedderProvider {
454        fn name(&self) -> &str {
455            &self.name
456        }
457
458        fn dimensions(&self) -> usize {
459            self.dimensions
460        }
461
462        async fn build(&self) -> RuntimeResult<Arc<dyn EmbeddingService>> {
463            Ok(Arc::new(CountingEmbeddingService {
464                calls: Arc::clone(&self.calls),
465                dimensions: self.dimensions,
466            }))
467        }
468    }
469
470    fn counting_embedding_runtime() -> (KhiveRuntime, Arc<AtomicUsize>) {
471        let model = EmbeddingModel::AllMiniLmL6V2;
472        let rt = KhiveRuntime::new(RuntimeConfig {
473            db_path: None,
474            embedding_model: Some(model),
475            packs: vec!["kg".to_string()],
476            ..RuntimeConfig::no_embeddings()
477        })
478        .expect("in-memory runtime");
479        let calls = Arc::new(AtomicUsize::new(0));
480        rt.register_embedder(CountingEmbedderProvider {
481            name: model.to_string(),
482            calls: Arc::clone(&calls),
483            dimensions: model.dimensions(),
484        });
485        (rt, calls)
486    }
487
488    #[tokio::test]
489    async fn keyword_only_search_skips_embedding_and_vector_arm() {
490        let (rt, embed_calls) = counting_embedding_runtime();
491        let tok = NamespaceToken::local();
492        let entity = Entity::new("local", "concept", "keyword candidate");
493        rt.entities(&tok)
494            .unwrap()
495            .upsert_entities(vec![entity.clone()])
496            .await
497            .unwrap();
498        rt.text(&tok)
499            .unwrap()
500            .upsert_document(TextDocument {
501                subject_id: entity.id,
502                kind: SubstrateKind::Entity,
503                record_kind: None,
504                namespace: "local".to_string(),
505                title: None,
506                body: "fusionkeyword".to_string(),
507                tags: vec![],
508                metadata: None,
509                updated_at: Utc::now(),
510            })
511            .await
512            .unwrap();
513
514        let hits = rt
515            .hybrid_search_with_strategy(
516                &tok,
517                "fusionkeyword",
518                None,
519                FusionStrategy::KeywordOnly,
520                10,
521            )
522            .await
523            .unwrap();
524        assert_eq!(hits.len(), 1);
525        assert_eq!(hits[0].entity_id, entity.id);
526        assert_eq!(hits[0].source, SearchSource::Text);
527        assert_eq!(embed_calls.load(Ordering::SeqCst), 0);
528
529        let mixed = rt
530            .hybrid_search_with_strategy(&tok, "fusionkeyword", None, FusionStrategy::rrf(), 10)
531            .await
532            .unwrap();
533        assert_eq!(mixed[0].entity_id, entity.id);
534        assert_eq!(embed_calls.load(Ordering::SeqCst), 1);
535    }
536
537    #[tokio::test]
538    async fn vector_only_search_skips_unavailable_text_arm() {
539        let (rt, embed_calls) = counting_embedding_runtime();
540        let tok = NamespaceToken::local();
541        let entity = Entity::new("local", "concept", "vector candidate");
542        rt.entities(&tok)
543            .unwrap()
544            .upsert_entities(vec![entity.clone()])
545            .await
546            .unwrap();
547        let query_vector = vec![1.0; EmbeddingModel::AllMiniLmL6V2.dimensions()];
548        rt.vectors(&tok)
549            .unwrap()
550            .insert(
551                entity.id,
552                SubstrateKind::Entity,
553                "local",
554                "entity.body",
555                vec![query_vector.clone()],
556            )
557            .await
558            .unwrap();
559        rt.text(&tok).unwrap();
560        let mut writer = rt.sql().writer().await.unwrap();
561        writer
562            .execute_script(
563                "DROP TABLE fts_entities; CREATE TABLE fts_entities (id INTEGER PRIMARY KEY);"
564                    .to_string(),
565            )
566            .await
567            .unwrap();
568        drop(writer);
569
570        let hits = rt
571            .hybrid_search_with_strategy(
572                &tok,
573                "fusionkeyword",
574                Some(query_vector),
575                FusionStrategy::VectorOnly,
576                10,
577            )
578            .await
579            .unwrap();
580        assert_eq!(hits.len(), 1);
581        assert_eq!(hits[0].entity_id, entity.id);
582        assert_eq!(hits[0].source, SearchSource::Vector);
583        assert_eq!(embed_calls.load(Ordering::SeqCst), 0);
584    }
585
586    #[tokio::test]
587    async fn fusion_evidence_labels_builtin_strategies_and_preserves_components() {
588        let rt = evidence_runtime();
589        let id = Uuid::from_u128(1);
590        let keyword = DeterministicScore::from_raw(1_i64 << 30);
591        let vector = DeterministicScore::from_raw(3_i64 << 30);
592        for (strategy, kind, raw_score, signals) in [
593            (
594                FusionStrategy::Rrf { k: 60 },
595                RankScoreKind::Rrf,
596                140_818_600,
597                SearchSignals {
598                    vector_similarity: Some(vector),
599                    keyword_score: Some(keyword),
600                },
601            ),
602            (
603                FusionStrategy::VectorOnly,
604                RankScoreKind::Vector,
605                vector.to_raw(),
606                SearchSignals {
607                    vector_similarity: Some(vector),
608                    keyword_score: None,
609                },
610            ),
611            (
612                FusionStrategy::KeywordOnly,
613                RankScoreKind::Keyword,
614                keyword.to_raw(),
615                SearchSignals {
616                    vector_similarity: None,
617                    keyword_score: Some(keyword),
618                },
619            ),
620            (
621                FusionStrategy::weighted(vec![0.5, 0.5]),
622                RankScoreKind::Weighted,
623                1_i64 << 32,
624                SearchSignals {
625                    vector_similarity: Some(vector),
626                    keyword_score: Some(keyword),
627                },
628            ),
629            (
630                FusionStrategy::Union,
631                RankScoreKind::Union,
632                vector.to_raw(),
633                SearchSignals {
634                    vector_similarity: Some(vector),
635                    keyword_score: Some(keyword),
636                },
637            ),
638        ] {
639            let hits = rt
640                .fuse_with_strategy(
641                    vec![text_hit(id, 0.25, "candidate")],
642                    vec![vector_hit(id, 0.75)],
643                    &strategy,
644                    10,
645                )
646                .await
647                .unwrap();
648            assert_eq!(hits.len(), 1);
649            assert_eq!(hits[0].entity_id, id);
650            assert_eq!(hits[0].score.to_raw(), raw_score);
651            assert_eq!(hits[0].rank_score_kind, kind);
652            assert_eq!(hits[0].signals, signals);
653        }
654        assert_eq!(RankScoreKind::Rrf.as_str(), "rrf");
655        assert_eq!(RankScoreKind::Vector.as_str(), "vector");
656        assert_eq!(RankScoreKind::Keyword.as_str(), "keyword");
657        assert_eq!(RankScoreKind::Weighted.as_str(), "weighted");
658        assert_eq!(RankScoreKind::Union.as_str(), "union");
659    }
660
661    #[tokio::test]
662    async fn fusion_evidence_distinguishes_absence_from_zero() {
663        let rt = evidence_runtime();
664        let id = Uuid::from_u128(1);
665        for (text, vector, signals) in [
666            (
667                vec![text_hit(id, 0.0, "zero keyword")],
668                vec![],
669                SearchSignals {
670                    vector_similarity: None,
671                    keyword_score: Some(DeterministicScore::ZERO),
672                },
673            ),
674            (
675                vec![],
676                vec![vector_hit(id, 0.0)],
677                SearchSignals {
678                    vector_similarity: Some(DeterministicScore::ZERO),
679                    keyword_score: None,
680                },
681            ),
682        ] {
683            let hits = rt
684                .fuse_with_strategy(text, vector, &FusionStrategy::rrf(), 10)
685                .await
686                .unwrap();
687            assert_eq!(hits.len(), 1);
688            assert_eq!(hits[0].signals, signals);
689        }
690        assert_eq!(
691            SearchSignals::default(),
692            SearchSignals {
693                vector_similarity: None,
694                keyword_score: None,
695            }
696        );
697    }
698
699    #[tokio::test]
700    async fn fusion_evidence_golden_preserves_true_ties_across_permutations() {
701        let rt = evidence_runtime();
702        let a = Uuid::from_u128(1);
703        let b = Uuid::from_u128(2);
704        let expected = vec![
705            (
706                a,
707                139_682_966,
708                RankScoreKind::Rrf,
709                SearchSignals {
710                    vector_similarity: Some(DeterministicScore::from_raw(1_i64 << 31)),
711                    keyword_score: Some(DeterministicScore::from_raw(1_i64 << 30)),
712                },
713            ),
714            (
715                b,
716                139_682_966,
717                RankScoreKind::Rrf,
718                SearchSignals {
719                    vector_similarity: Some(DeterministicScore::from_raw(1_i64 << 32)),
720                    keyword_score: Some(DeterministicScore::from_raw(3_i64 << 30)),
721                },
722            ),
723        ];
724        for (text_ids, vector_ids) in [([a, b], [b, a]), ([b, a], [a, b])] {
725            for _ in 0..4 {
726                let text = text_ids
727                    .into_iter()
728                    .map(|id| text_hit(id, if id == a { 0.25 } else { 0.75 }, "candidate"))
729                    .collect();
730                let vector = vector_ids
731                    .into_iter()
732                    .map(|id| vector_hit(id, if id == a { 0.5 } else { 1.0 }))
733                    .collect();
734                let hits = rt
735                    .fuse_with_strategy(text, vector, &FusionStrategy::Rrf { k: 60 }, 10)
736                    .await
737                    .unwrap();
738                assert_eq!(hits.len(), 2);
739                assert_eq!(hits[0].score, hits[1].score);
740                assert!(hits.iter().all(|hit| hit.source == SearchSource::Both));
741                let snapshot: Vec<_> = hits
742                    .iter()
743                    .map(|hit| {
744                        (
745                            hit.entity_id,
746                            hit.score.to_raw(),
747                            hit.rank_score_kind,
748                            hit.signals,
749                        )
750                    })
751                    .collect();
752                assert_eq!(snapshot, expected);
753            }
754        }
755    }
756
757    #[tokio::test]
758    async fn fusion_evidence_duplicate_selection_follows_strategy() {
759        let rt = evidence_runtime();
760        let id = Uuid::from_u128(1);
761        for (strategy, keyword_raw) in [
762            (FusionStrategy::rrf(), 1_i64 << 30),
763            (FusionStrategy::Union, 3_i64 << 30),
764            (FusionStrategy::weighted(vec![0.5, 0.5]), 3_i64 << 30),
765        ] {
766            let hits = rt
767                .fuse_with_strategy(
768                    vec![text_hit(id, 0.25, "first"), text_hit(id, 0.75, "second")],
769                    vec![vector_hit(id, 0.5)],
770                    &strategy,
771                    10,
772                )
773                .await
774                .unwrap();
775            assert_eq!(hits.len(), 1);
776            assert_eq!(
777                hits[0].signals,
778                SearchSignals {
779                    vector_similarity: Some(DeterministicScore::from_raw(1_i64 << 31)),
780                    keyword_score: Some(DeterministicScore::from_raw(keyword_raw)),
781                }
782            );
783            assert_eq!(hits[0].title.as_deref(), Some("first"));
784        }
785    }
786
787    #[tokio::test]
788    async fn custom_fusion_evidence_uses_declared_kind() {
789        let rt = evidence_runtime();
790        let id = Uuid::from_u128(1);
791        rt.register_fusion_strategy("invert", Arc::new(InvertScoreExecutor));
792        let strategy =
793            FusionStrategy::try_custom("invert".into(), serde_json::Value::Null).unwrap();
794        let hits = rt
795            .fuse_with_strategy(vec![text_hit(id, 0.75, "candidate")], vec![], &strategy, 10)
796            .await
797            .unwrap();
798        assert_eq!(hits.len(), 1);
799        assert_eq!(hits[0].score.to_raw(), 1_i64 << 30);
800        assert_eq!(hits[0].rank_score_kind, RankScoreKind::Weighted);
801        assert_eq!(
802            hits[0].signals,
803            SearchSignals {
804                vector_similarity: None,
805                keyword_score: Some(DeterministicScore::from_raw(3_i64 << 30)),
806            }
807        );
808    }
809
810    fn cosine_fixture_vector(dimensions: usize, x: f32, y: f32) -> Vec<f32> {
811        let mut vector = vec![0.0; dimensions];
812        vector[0] = x;
813        vector[1] = y;
814        vector
815    }
816
817    async fn stale_full_prefix_fixture() -> (
818        KhiveRuntime,
819        NamespaceToken,
820        &'static str,
821        Vec<f32>,
822        Vec<TextSearchHit>,
823        Vec<VectorSearchHit>,
824        HashSet<Uuid>,
825    ) {
826        let model = EmbeddingModel::AllMiniLmL6V2;
827        let dimensions = model.dimensions();
828        let rt = KhiveRuntime::new(RuntimeConfig {
829            db_path: None,
830            embedding_model: Some(model),
831            additional_embedding_models: vec![],
832            ..RuntimeConfig::default()
833        })
834        .unwrap();
835        let tok = NamespaceToken::local();
836        let query_text = "fusionrefillterm";
837        let query_vector = cosine_fixture_vector(dimensions, 1.0, 0.0);
838
839        let common_stale_a = Uuid::from_u128(1);
840        let common_stale_b = Uuid::from_u128(2);
841        let text_only_stale = Uuid::from_u128(3);
842        let vector_only_stale = Uuid::from_u128(4);
843
844        let live_text = Entity::new("local", "concept", "live text candidate");
845        let live_vector = Entity::new("local", "concept", "live vector candidate");
846        rt.entities(&tok)
847            .unwrap()
848            .upsert_entities(vec![live_text.clone(), live_vector.clone()])
849            .await
850            .unwrap();
851
852        let document = |subject_id, repetitions: usize| TextDocument {
853            subject_id,
854            kind: SubstrateKind::Entity,
855            record_kind: None,
856            namespace: "local".to_string(),
857            title: None,
858            body: std::iter::repeat_n(query_text, repetitions)
859                .collect::<Vec<_>>()
860                .join(" "),
861            tags: vec![],
862            metadata: None,
863            updated_at: Utc::now(),
864        };
865        rt.text(&tok)
866            .unwrap()
867            .upsert_documents(vec![
868                document(common_stale_a, 12),
869                document(common_stale_b, 8),
870                document(text_only_stale, 4),
871                document(live_text.id, 1),
872            ])
873            .await
874            .unwrap();
875
876        let vectors = rt.vectors(&tok).unwrap();
877        for (id, vector) in [
878            (common_stale_a, cosine_fixture_vector(dimensions, 1.0, 0.0)),
879            (common_stale_b, cosine_fixture_vector(dimensions, 0.8, 0.6)),
880            (
881                vector_only_stale,
882                cosine_fixture_vector(dimensions, 0.5, 0.866_025_4),
883            ),
884            (live_vector.id, cosine_fixture_vector(dimensions, -1.0, 0.0)),
885        ] {
886            vectors
887                .insert(
888                    id,
889                    SubstrateKind::Entity,
890                    "local",
891                    "entity.body",
892                    vec![vector],
893                )
894                .await
895                .unwrap();
896        }
897
898        let text_hits = rt
899            .text(&tok)
900            .unwrap()
901            .search(TextSearchRequest {
902                query: query_text.to_string(),
903                mode: TextQueryMode::Plain,
904                filter: Some(TextFilter {
905                    namespaces: vec!["local".to_string()],
906                    ..TextFilter::default()
907                }),
908                top_k: CANDIDATE_MULTIPLIER,
909                snippet_chars: 0,
910            })
911            .await
912            .unwrap();
913        let vector_hits = vectors
914            .search(VectorSearchRequest {
915                query_vectors: vec![query_vector.clone()],
916                top_k: CANDIDATE_MULTIPLIER,
917                namespace: Some("local".to_string()),
918                kind: Some(SubstrateKind::Entity),
919                embedding_model: None,
920                filter: None,
921                backend_hints: None,
922            })
923            .await
924            .unwrap();
925
926        assert_eq!(text_hits.len(), CANDIDATE_MULTIPLIER as usize);
927        assert_eq!(vector_hits.len(), CANDIDATE_MULTIPLIER as usize);
928        let live = HashSet::from([live_text.id, live_vector.id]);
929        (
930            rt,
931            tok,
932            query_text,
933            query_vector,
934            text_hits,
935            vector_hits,
936            live,
937        )
938    }
939
940    // 1. RRF with custom k produces different ordering than k=60
941    #[tokio::test]
942    async fn rrf_custom_k_differs_from_k60() {
943        let rt = KhiveRuntime::memory().unwrap();
944        let a = Uuid::new_v4();
945        let b = Uuid::new_v4();
946        // Single-source input makes a and b tie in relative order at both k values,
947        // so assert on raw score magnitude (smaller k widens the rank-1-vs-rank-2 gap)
948        // rather than ordering.
949        let text = vec![text_hit(a, 0.9, "a"), text_hit(b, 0.1, "b")];
950        let hits_k1 = rt
951            .fuse_with_strategy(text.clone(), vec![], &FusionStrategy::Rrf { k: 1 }, 10)
952            .await
953            .unwrap();
954        let hits_k60 = rt
955            .fuse_with_strategy(text, vec![], &FusionStrategy::Rrf { k: 60 }, 10)
956            .await
957            .unwrap();
958        // Both should have a first (rank 1 always wins in single-source)
959        assert_eq!(hits_k1[0].entity_id, a);
960        assert_eq!(hits_k60[0].entity_id, a);
961        // k=1 produces higher raw score for rank 1 than k=60
962        assert!(hits_k1[0].score > hits_k60[0].score);
963    }
964
965    // 2. Canonical [vector, keyword] weights change ordering as documented.
966    #[tokio::test]
967    async fn weighted_ordering_depends_on_weights() {
968        let rt = KhiveRuntime::memory().unwrap();
969        let a = Uuid::new_v4();
970        let b = Uuid::new_v4();
971        // a scores high in text, b scores high in vector
972        let text = vec![text_hit(a, 0.9, "a"), text_hit(b, 0.1, "b")];
973        let vec_hits = vec![vector_hit(b, 0.9), vector_hit(a, 0.1)];
974
975        let heavy_vector = rt
976            .fuse_with_strategy(
977                text.clone(),
978                vec_hits.clone(),
979                &FusionStrategy::Weighted {
980                    weights: vec![0.7, 0.3],
981                },
982                10,
983            )
984            .await
985            .unwrap();
986        let heavy_keyword = rt
987            .fuse_with_strategy(
988                text,
989                vec_hits,
990                &FusionStrategy::Weighted {
991                    weights: vec![0.3, 0.7],
992                },
993                10,
994            )
995            .await
996            .unwrap();
997
998        assert_eq!(heavy_vector[0].entity_id, b);
999        assert_eq!(heavy_keyword[0].entity_id, a);
1000    }
1001
1002    // 3. Weighted [7.0, 3.0] = Weighted [0.7, 0.3] (normalization)
1003    #[tokio::test]
1004    async fn weighted_scale_invariant() {
1005        let rt = KhiveRuntime::memory().unwrap();
1006        let a = Uuid::new_v4();
1007        let b = Uuid::new_v4();
1008        let text = vec![text_hit(a, 0.9, "a"), text_hit(b, 0.1, "b")];
1009        let vec_hits = vec![vector_hit(b, 0.9), vector_hit(a, 0.1)];
1010
1011        let w1 = rt
1012            .fuse_with_strategy(
1013                text.clone(),
1014                vec_hits.clone(),
1015                &FusionStrategy::Weighted {
1016                    weights: vec![0.7, 0.3],
1017                },
1018                10,
1019            )
1020            .await
1021            .unwrap();
1022        let w2 = rt
1023            .fuse_with_strategy(
1024                text,
1025                vec_hits,
1026                &FusionStrategy::Weighted {
1027                    weights: vec![7.0, 3.0],
1028                },
1029                10,
1030            )
1031            .await
1032            .unwrap();
1033
1034        assert_eq!(w1[0].entity_id, w2[0].entity_id);
1035        assert_eq!(w1[1].entity_id, w2[1].entity_id);
1036        let diff = (w1[0].score.to_f64() - w2[0].score.to_f64()).abs();
1037        assert!(diff < 1e-9, "scores differ by {diff}");
1038    }
1039
1040    // 4. Weighted [0.0, 0.0] falls back to equal weights
1041    #[tokio::test]
1042    async fn weighted_zero_weights_equal_fallback() {
1043        let rt = KhiveRuntime::memory().unwrap();
1044        let a = Uuid::new_v4();
1045        let b = Uuid::new_v4();
1046        // Both sources agree: a > b
1047        let text = vec![text_hit(a, 0.9, "a"), text_hit(b, 0.1, "b")];
1048        let vec_hits = vec![vector_hit(a, 0.9), vector_hit(b, 0.1)];
1049
1050        let hits = rt
1051            .fuse_with_strategy(
1052                text,
1053                vec_hits,
1054                &FusionStrategy::Weighted {
1055                    weights: vec![0.0, 0.0],
1056                },
1057                10,
1058            )
1059            .await
1060            .unwrap();
1061        assert_eq!(hits[0].entity_id, a);
1062    }
1063
1064    // 5. Weighted with negative weight clamps to 0
1065    #[tokio::test]
1066    async fn weighted_negative_weight_clamped() {
1067        let rt = KhiveRuntime::memory().unwrap();
1068        let a = Uuid::new_v4();
1069        let text = vec![text_hit(a, 0.9, "a")];
1070        // Negative vector weight โ†’ only keyword/text contributes.
1071        let hits = rt
1072            .fuse_with_strategy(
1073                text,
1074                vec![],
1075                &FusionStrategy::Weighted {
1076                    weights: vec![-0.5, 1.0],
1077                },
1078                10,
1079            )
1080            .await
1081            .unwrap();
1082        assert_eq!(hits.len(), 1);
1083        assert_eq!(hits[0].entity_id, a);
1084    }
1085
1086    #[tokio::test]
1087    async fn weighted_empty_arm_keeps_canonical_position() {
1088        let rt = KhiveRuntime::memory().unwrap();
1089        let text_only = Uuid::new_v4();
1090        let hits = rt
1091            .fuse_with_strategy(
1092                vec![text_hit(text_only, 0.9, "text")],
1093                vec![],
1094                &FusionStrategy::Weighted {
1095                    // Canonical [vector, keyword]: the only non-empty arm has zero weight.
1096                    weights: vec![1.0, 0.0],
1097                },
1098                10,
1099            )
1100            .await
1101            .unwrap();
1102        assert!(
1103            hits.is_empty(),
1104            "dropping the empty vector arm would incorrectly rebind text to its weight"
1105        );
1106    }
1107
1108    // 6. Union returns max score per entity when same id appears in both lists
1109    #[tokio::test]
1110    async fn union_max_score_per_entity() {
1111        let rt = KhiveRuntime::memory().unwrap();
1112        let a = Uuid::new_v4();
1113        let text = vec![text_hit(a, 0.3, "a")];
1114        let vec_hits = vec![vector_hit(a, 0.9)];
1115
1116        let hits = rt
1117            .fuse_with_strategy(text, vec_hits, &FusionStrategy::Union, 10)
1118            .await
1119            .unwrap();
1120        assert_eq!(hits.len(), 1);
1121        assert!((hits[0].score.to_f64() - 0.9).abs() < 1e-6);
1122        assert_eq!(hits[0].source, SearchSource::Both);
1123    }
1124
1125    // 7. VectorOnly returns vector hits only (text hits dropped)
1126    #[tokio::test]
1127    async fn vector_only_drops_text() {
1128        let rt = KhiveRuntime::memory().unwrap();
1129        let a = Uuid::new_v4();
1130        let b = Uuid::new_v4();
1131        let text = vec![text_hit(b, 0.9, "b")];
1132        let vec_hits = vec![vector_hit(a, 0.8)];
1133
1134        let hits = rt
1135            .fuse_with_strategy(text, vec_hits, &FusionStrategy::VectorOnly, 10)
1136            .await
1137            .unwrap();
1138        assert_eq!(hits.len(), 1);
1139        assert_eq!(hits[0].entity_id, a);
1140        assert_eq!(hits[0].source, SearchSource::Vector);
1141        assert!(hits[0].title.is_none());
1142    }
1143
1144    #[tokio::test]
1145    async fn keyword_only_drops_vector() {
1146        let rt = KhiveRuntime::memory().unwrap();
1147        let text_id = Uuid::new_v4();
1148        let vector_id = Uuid::new_v4();
1149        let hits = rt
1150            .fuse_with_strategy(
1151                vec![text_hit(text_id, 0.8, "text")],
1152                vec![vector_hit(vector_id, 0.9)],
1153                &FusionStrategy::KeywordOnly,
1154                10,
1155            )
1156            .await
1157            .unwrap();
1158        assert_eq!(hits.len(), 1);
1159        assert_eq!(hits[0].entity_id, text_id);
1160        assert_eq!(hits[0].source, SearchSource::Text);
1161    }
1162
1163    /// Test-only executor: flattens all streams and reverses their order,
1164    /// keeping each candidate's original score.
1165    struct ReverseOrderExecutor;
1166
1167    #[async_trait::async_trait]
1168    impl FusionExecutor for ReverseOrderExecutor {
1169        fn rank_score_kind(&self) -> RankScoreKind {
1170            RankScoreKind::Union
1171        }
1172
1173        async fn fuse(
1174            &self,
1175            streams: Vec<CandidateStream>,
1176            _params: &serde_json::Value,
1177            _limit: usize,
1178        ) -> RuntimeResult<Vec<RankedHit>> {
1179            let mut flat: Vec<_> = streams.into_iter().flatten().collect();
1180            flat.reverse();
1181            Ok(flat)
1182        }
1183    }
1184
1185    /// Test-only executor: inverts each candidate's score (`1.0 - score`) so
1186    /// the fused ranking is the reverse of what score-descending built-ins
1187    /// (RRF, Union, Weighted) would produce on the same fixture -- unlike a
1188    /// mere insertion-order reversal, this survives the dispatch boundary's
1189    /// canonical re-sort, since the *scores* (not just the order) differ.
1190    struct InvertScoreExecutor;
1191
1192    #[async_trait::async_trait]
1193    impl FusionExecutor for InvertScoreExecutor {
1194        fn rank_score_kind(&self) -> RankScoreKind {
1195            RankScoreKind::Weighted
1196        }
1197
1198        async fn fuse(
1199            &self,
1200            streams: Vec<CandidateStream>,
1201            _params: &serde_json::Value,
1202            _limit: usize,
1203        ) -> RuntimeResult<Vec<RankedHit>> {
1204            Ok(streams
1205                .into_iter()
1206                .flatten()
1207                .map(|(id, score)| (id, DeterministicScore::from_f64(1.0 - score.to_f64())))
1208                .collect())
1209        }
1210    }
1211
1212    /// Test-only executor: returns every candidate at an identical score, in
1213    /// the arbitrary order the input streams happened to flatten to -- used
1214    /// to prove the dispatch boundary re-sorts by the canonical comparator
1215    /// rather than trusting executor output order.
1216    struct EqualScoreExecutor;
1217
1218    #[async_trait::async_trait]
1219    impl FusionExecutor for EqualScoreExecutor {
1220        fn rank_score_kind(&self) -> RankScoreKind {
1221            RankScoreKind::Weighted
1222        }
1223
1224        async fn fuse(
1225            &self,
1226            streams: Vec<CandidateStream>,
1227            _params: &serde_json::Value,
1228            _limit: usize,
1229        ) -> RuntimeResult<Vec<RankedHit>> {
1230            Ok(streams
1231                .into_iter()
1232                .flatten()
1233                .map(|(id, _)| (id, DeterministicScore::from_f64(1.0)))
1234                .collect())
1235        }
1236    }
1237
1238    // 7b. A registered custom executor dispatches and yields a different
1239    // ranking than RRF on the same fixture.
1240    #[tokio::test]
1241    async fn custom_strategy_dispatches_through_executor_and_differs_from_rrf() {
1242        let rt = KhiveRuntime::memory().unwrap();
1243        let a = Uuid::new_v4();
1244        let b = Uuid::new_v4();
1245        let text = vec![text_hit(a, 0.9, "a"), text_hit(b, 0.5, "b")];
1246
1247        rt.register_fusion_strategy("invert", Arc::new(InvertScoreExecutor));
1248        let strategy =
1249            FusionStrategy::try_custom("invert".to_string(), serde_json::Value::Null).unwrap();
1250
1251        let custom = rt
1252            .fuse_with_strategy(text.clone(), vec![], &strategy, 10)
1253            .await
1254            .unwrap();
1255        let rrf = rt
1256            .fuse_with_strategy(text, vec![], &FusionStrategy::rrf(), 10)
1257            .await
1258            .unwrap();
1259
1260        let custom_ids: Vec<_> = custom.iter().map(|h| h.entity_id).collect();
1261        let rrf_ids: Vec<_> = rrf.iter().map(|h| h.entity_id).collect();
1262        assert_ne!(
1263            custom_ids, rrf_ids,
1264            "custom and RRF must yield different orderings on this fixture"
1265        );
1266    }
1267
1268    // 7c. An unregistered Custom name fails closed rather than falling back to RRF.
1269    #[tokio::test]
1270    async fn custom_strategy_unknown_name_fails_closed() {
1271        let rt = KhiveRuntime::memory().unwrap();
1272        let a = Uuid::new_v4();
1273        let text = vec![text_hit(a, 0.9, "a")];
1274        let strategy =
1275            FusionStrategy::try_custom("nonexistent".to_string(), serde_json::Value::Null).unwrap();
1276
1277        let result = rt.fuse_with_strategy(text, vec![], &strategy, 10).await;
1278        assert!(matches!(
1279            result,
1280            Err(RuntimeError::UnknownFusionStrategy(name)) if name == "nonexistent"
1281        ));
1282    }
1283
1284    // 7d. Unknown name errors even with empty sources -- it must not be
1285    // indistinguishable from a valid empty result.
1286    #[tokio::test]
1287    async fn custom_strategy_unknown_name_fails_closed_even_on_empty_input() {
1288        let rt = KhiveRuntime::memory().unwrap();
1289        let strategy =
1290            FusionStrategy::try_custom("nonexistent".to_string(), serde_json::Value::Null).unwrap();
1291
1292        let result = rt.fuse_with_strategy(vec![], vec![], &strategy, 10).await;
1293        assert!(matches!(
1294            result,
1295            Err(RuntimeError::UnknownFusionStrategy(name)) if name == "nonexistent"
1296        ));
1297    }
1298
1299    // 7e. Empty sources with a *registered* name is a valid empty result, not
1300    // an error -- distinguishing "misconfigured" from "genuinely nothing".
1301    #[tokio::test]
1302    async fn custom_strategy_registered_name_empty_input_returns_ok_empty() {
1303        let rt = KhiveRuntime::memory().unwrap();
1304        rt.register_fusion_strategy("reverse", Arc::new(ReverseOrderExecutor));
1305        let strategy =
1306            FusionStrategy::try_custom("reverse".to_string(), serde_json::Value::Null).unwrap();
1307
1308        let result = rt
1309            .fuse_with_strategy(vec![], vec![], &strategy, 10)
1310            .await
1311            .unwrap();
1312        assert!(result.is_empty());
1313    }
1314
1315    // 7f. Registering a custom strategy never perturbs the default (non-Custom) path.
1316    #[tokio::test]
1317    async fn registered_custom_strategy_leaves_default_path_unaffected() {
1318        let rt = KhiveRuntime::memory().unwrap();
1319        let a = Uuid::new_v4();
1320        let b = Uuid::new_v4();
1321        let text = vec![text_hit(a, 0.9, "a"), text_hit(b, 0.5, "b")];
1322
1323        rt.register_fusion_strategy("reverse", Arc::new(ReverseOrderExecutor));
1324
1325        let via_rt_with_registration = rt
1326            .fuse_with_strategy(text.clone(), vec![], &FusionStrategy::rrf(), 10)
1327            .await
1328            .unwrap();
1329        let rt2 = KhiveRuntime::memory().unwrap();
1330        let via_rt_without_registration = rt2
1331            .fuse_with_strategy(text, vec![], &FusionStrategy::rrf(), 10)
1332            .await
1333            .unwrap();
1334
1335        let ids_with: Vec<_> = via_rt_with_registration
1336            .iter()
1337            .map(|h| h.entity_id)
1338            .collect();
1339        let ids_without: Vec<_> = via_rt_without_registration
1340            .iter()
1341            .map(|h| h.entity_id)
1342            .collect();
1343        assert_eq!(ids_with, ids_without);
1344    }
1345
1346    // 7g. A custom executor returning equal-score IDs in arbitrary/reversed
1347    // order still yields the crate's canonical score-desc/id-asc order --
1348    // the dispatch boundary re-sorts rather than trusting executor output.
1349    #[tokio::test]
1350    async fn custom_executor_output_is_sorted_by_canonical_comparator() {
1351        let rt = KhiveRuntime::memory().unwrap();
1352        // Deliberately not in ID order, so a passthrough bug would be visible.
1353        let ids: Vec<Uuid> = vec![Uuid::from_u128(3), Uuid::from_u128(1), Uuid::from_u128(2)];
1354        let text: Vec<TextSearchHit> = ids.iter().map(|&id| text_hit(id, 0.5, "tied")).collect();
1355
1356        rt.register_fusion_strategy("equal_score", Arc::new(EqualScoreExecutor));
1357        let strategy =
1358            FusionStrategy::try_custom("equal_score".to_string(), serde_json::Value::Null).unwrap();
1359
1360        let hits = rt
1361            .fuse_with_strategy(text, vec![], &strategy, 10)
1362            .await
1363            .unwrap();
1364
1365        let mut expected = ids.clone();
1366        expected.sort();
1367        let actual: Vec<_> = hits.iter().map(|h| h.entity_id).collect();
1368        assert_eq!(
1369            actual, expected,
1370            "equal-score executor output must be tie-broken by ascending ID"
1371        );
1372    }
1373
1374    // 8. Default strategy is Rrf{k:60}
1375    #[test]
1376    fn default_strategy_is_rrf_k60() {
1377        assert_eq!(FusionStrategy::default(), FusionStrategy::Rrf { k: 60 });
1378    }
1379
1380    #[tokio::test]
1381    async fn hybrid_union_alive_filter_refills_below_complete_four_x_prefix() {
1382        let (rt, tok, query_text, query_vector, text_hits, vector_hits, live) =
1383            stale_full_prefix_fixture().await;
1384        let truncated = rt
1385            .fuse_with_strategy(
1386                text_hits,
1387                vector_hits,
1388                &FusionStrategy::Union,
1389                CANDIDATE_MULTIPLIER as usize,
1390            )
1391            .await
1392            .unwrap();
1393        assert!(truncated.iter().all(|hit| !live.contains(&hit.entity_id)));
1394
1395        let hits = rt
1396            .hybrid_search_with_strategy(
1397                &tok,
1398                query_text,
1399                Some(query_vector),
1400                FusionStrategy::Union,
1401                1,
1402            )
1403            .await
1404            .unwrap();
1405
1406        assert_eq!(hits.len(), 1);
1407        assert!(live.contains(&hits[0].entity_id));
1408    }
1409
1410    #[tokio::test]
1411    async fn hybrid_rrf_alive_filter_refills_below_complete_four_x_prefix() {
1412        let (rt, tok, query_text, query_vector, text_hits, vector_hits, live) =
1413            stale_full_prefix_fixture().await;
1414        let strategy = FusionStrategy::Rrf { k: 60 };
1415        let truncated = rt
1416            .fuse_with_strategy(
1417                text_hits,
1418                vector_hits,
1419                &strategy,
1420                CANDIDATE_MULTIPLIER as usize,
1421            )
1422            .await
1423            .unwrap();
1424        assert!(truncated.iter().all(|hit| !live.contains(&hit.entity_id)));
1425
1426        let hits = rt
1427            .hybrid_search_with_strategy(&tok, query_text, Some(query_vector), strategy, 1)
1428            .await
1429            .unwrap();
1430
1431        assert_eq!(hits.len(), 1);
1432        assert!(live.contains(&hits[0].entity_id));
1433    }
1434
1435    #[tokio::test]
1436    async fn hybrid_default_rrf_alive_filter_refills_below_complete_four_x_prefix() {
1437        let (rt, tok, query_text, query_vector, _text_hits, _vector_hits, live) =
1438            stale_full_prefix_fixture().await;
1439
1440        let hits = rt
1441            .hybrid_search(
1442                &tok,
1443                query_text,
1444                Some(query_vector),
1445                1,
1446                None,
1447                None,
1448                &[],
1449                None,
1450            )
1451            .await
1452            .unwrap();
1453
1454        assert_eq!(hits.len(), 1);
1455        assert!(live.contains(&hits[0].entity_id));
1456    }
1457
1458    // 9. Roundtrip serde preserves variant
1459    #[test]
1460    fn serde_roundtrip() {
1461        let cases = vec![
1462            FusionStrategy::Rrf { k: 60 },
1463            FusionStrategy::Rrf { k: 20 },
1464            FusionStrategy::Weighted {
1465                weights: vec![0.7, 0.3],
1466            },
1467            FusionStrategy::Union,
1468            FusionStrategy::VectorOnly,
1469            FusionStrategy::KeywordOnly,
1470        ];
1471        for strategy in cases {
1472            let json = serde_json::to_string(&strategy).expect("serialize");
1473            let back: FusionStrategy = serde_json::from_str(&json).expect("deserialize");
1474            assert_eq!(strategy, back, "roundtrip failed for {json}");
1475        }
1476    }
1477
1478    // 10. hybrid_search_with_strategy must not hard-fail on a query containing FTS5
1479    // metacharacters like `$`, since sanitize_fts5_query strips them before the query
1480    // reaches SQLite. This covers the sanitizer path; test 11 covers the fail-loud
1481    // path for characters the sanitizer does not strip.
1482    #[tokio::test]
1483    async fn hybrid_search_with_strategy_dollar_sign_query_does_not_error() {
1484        let rt = KhiveRuntime::memory().unwrap();
1485        let tok = NamespaceToken::local();
1486        rt.create_entity(
1487            &tok,
1488            "concept",
1489            None,
1490            "DSL docs",
1491            Some("use $prev.id to chain calls"),
1492            None,
1493            vec![],
1494        )
1495        .await
1496        .unwrap();
1497
1498        let result = rt
1499            .hybrid_search_with_strategy(&tok, "$prev.id", None, FusionStrategy::default(), 10)
1500            .await;
1501
1502        assert!(
1503            result.is_ok(),
1504            "#388 hybrid_search_with_strategy must not hard-fail on a '$'-bearing query, got: {:?}",
1505            result.err()
1506        );
1507    }
1508
1509    // 11. #916: `@` used to reach SQLite FTS5's bareword parser raw and error,
1510    // surfacing as RuntimeError::InvalidInput per #569's fail-loud policy.
1511    // sanitize_fts5_token_group's bareword-safety gate now routes it through the
1512    // quoted-phrase alternative instead, so the query succeeds and the fail-loud
1513    // arm is no longer reached for ordinary punctuation.
1514    #[tokio::test]
1515    async fn hybrid_search_with_strategy_residual_fts5_char_now_sanitized() {
1516        let rt = KhiveRuntime::memory().unwrap();
1517        let tok = NamespaceToken::local();
1518        rt.create_entity(
1519            &tok,
1520            "concept",
1521            None,
1522            "DSL docs",
1523            Some("use foo@bar to chain calls"),
1524            None,
1525            vec![],
1526        )
1527        .await
1528        .unwrap();
1529
1530        let result = rt
1531            .hybrid_search_with_strategy(&tok, "foo@bar", None, FusionStrategy::default(), 10)
1532            .await;
1533
1534        let hits = result.unwrap_or_else(|e| {
1535            panic!(
1536                "#916 hybrid_search_with_strategy must not fail on an '@'-bearing query, got: {e:?}"
1537            )
1538        });
1539        assert!(
1540            !hits.is_empty(),
1541            "#916 '@'-bearing query must still find the seeded 'foo@bar' content via the \
1542             quoted-phrase alternative"
1543        );
1544    }
1545}