1use 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
20pub type CandidateStream = Vec<(Uuid, DeterministicScore)>;
24
25pub type RankedHit = (Uuid, DeterministicScore);
27
28#[async_trait::async_trait]
36pub trait FusionExecutor: Send + Sync + 'static {
37 fn rank_score_kind(&self) -> RankScoreKind;
39
40 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
54pub(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 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 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 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 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 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 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 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 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 #[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 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 assert_eq!(hits_k1[0].entity_id, a);
960 assert_eq!(hits_k60[0].entity_id, a);
961 assert!(hits_k1[0].score > hits_k60[0].score);
963 }
964
965 #[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 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 #[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 #[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 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 #[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 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 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 #[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 #[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 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 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 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 #[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 #[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 #[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 #[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 #[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 #[tokio::test]
1350 async fn custom_executor_output_is_sorted_by_canonical_comparator() {
1351 let rt = KhiveRuntime::memory().unwrap();
1352 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 #[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 #[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 #[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 #[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}