From c359fcbcd8d0e4379039945b225afe9fa5d351c6 Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Mon, 15 Dec 2025 17:49:31 +0100 Subject: [PATCH] feat: two-stage RRF for fair wiki vs web ranking MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Merge vector+graph into single wiki source before RRF with web - Wiki pages no longer get 2x advantage from dual retrieval - Add vector similarity threshold (0.7 default) - Skip synonyms in graph search to reduce noise - Fix duplicate entity links bug in graph search 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 --- CHANGELOG.md | 23 ++++ pyproject.toml | 2 +- src/config.py | 1 + src/services/graph_service.py | 14 ++- src/services/hybrid_rag_service.py | 179 ++++++++++++++++++++++------- tests/test_hybrid_rag.py | 75 ++++++------ 6 files changed, 213 insertions(+), 81 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 82e166c..ca5d87c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,29 @@ All notable changes to Library Desk will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [1.3.0] - 2025-12-15 + +### Changed + +- **Two-Stage RRF Architecture** - Major refactor to level the playing field between wiki and web results + - Stage 1: Vector and graph results merged into single "wiki" ranking using mini-RRF + - Stage 2: Final RRF between wiki (single source) and web (single source) + - Wiki pages no longer get 2x advantage from appearing in both vector and graph searches + - Multi-source confirmation still determines wiki internal ranking + +- **Skip synonyms in graph search** - LLM-generated synonyms (e.g., "author") no longer match unrelated graph entities (e.g., "author2000") + - Vector search still uses synonyms for semantic similarity + - Graph search uses only core keywords for exact entity matching + +### Added + +- `VECTOR_SIMILARITY_THRESHOLD` config setting (default: 0.7) to filter weak vector matches +- Deduplication in graph search to prevent same document appearing multiple times + +### Fixed + +- Graph search duplicate entity bug where same document could appear twice if entity linked multiple times + ## [1.2.1] - 2025-12-15 ### Fixed diff --git a/pyproject.toml b/pyproject.toml index 72b85a5..5e88845 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "library-desk" -version = "1.2.1" +version = "1.3.0" description = "Coordination service for The Library system - HybridRAG queries, document ingestion, entity extraction, and knowledge consolidation" readme = "README.md" requires-python = ">=3.12" diff --git a/src/config.py b/src/config.py index 1a69c60..1c846be 100644 --- a/src/config.py +++ b/src/config.py @@ -72,6 +72,7 @@ class Settings(BaseSettings): hybrid_rag_vector_limit: int = Field(default=10, ge=1, le=50, description="Vector search limit") hybrid_rag_graph_limit: int = Field(default=10, ge=1, le=50, description="Graph search limit") hybrid_rag_web_limit: int = Field(default=5, ge=1, le=20, description="Web search limit") + vector_similarity_threshold: float = Field(default=0.7, ge=0.0, le=1.0, description="Minimum similarity score for vector results") # Entity Linking Fuzzy Matching Configuration entity_linking_min_confidence: float = Field(default=0.70, ge=0.0, le=1.0, description="Minimum confidence for entity-document matching") diff --git a/src/services/graph_service.py b/src/services/graph_service.py index 8358ec4..34f8bc2 100644 --- a/src/services/graph_service.py +++ b/src/services/graph_service.py @@ -1077,8 +1077,18 @@ Feel free to expand it with more details! search_query, {"terms": all_terms, "limit": limit} ) - logger.info(f"Graph search found {len(results)} documents") - return results + + # Deduplicate by page_id (safety net for any edge cases) + seen_page_ids = set() + unique_results = [] + for r in results: + page_id = r.get("page_id") + if page_id and page_id not in seen_page_ids: + seen_page_ids.add(page_id) + unique_results.append(r) + + logger.info(f"Graph search found {len(unique_results)} unique documents (raw: {len(results)})") + return unique_results except Exception as e: logger.error(f"Graph document search failed: {e}", exc_info=True) return [] diff --git a/src/services/hybrid_rag_service.py b/src/services/hybrid_rag_service.py index 13beb75..331798a 100644 --- a/src/services/hybrid_rag_service.py +++ b/src/services/hybrid_rag_service.py @@ -105,14 +105,20 @@ class HybridRAGService: timing["graph_ms"] = raw_results.get("timing", {}).get("graph_ms", 0) timing["web_ms"] = raw_results.get("timing", {}).get("web_ms", 0) - # Phase 2: RRF Fusion + # Phase 2: Two-Stage RRF Fusion phase2_start = time.time() + + # Stage 1: Merge wiki sources (vector + graph) into single ranking + wiki_merged = self._merge_wiki_sources( + vector_results=raw_results.get("vector", []), + graph_results=raw_results.get("graph", []), + k=config.rrf_k + ) + + # Stage 2: Final RRF between wiki and web (equal footing) fused_results = self._reciprocal_rank_fusion( - results_by_source={ - "vector": raw_results.get("vector", []), - "graph": raw_results.get("graph", []), - "web": raw_results.get("web", []) - }, + wiki_results=wiki_merged, + web_results=raw_results.get("web", []), k=config.rrf_k ) timing["fusion_ms"] = (time.time() - phase2_start) * 1000 @@ -288,7 +294,8 @@ JSON:""" response = await self.vector.search( query=query, user=user, - limit=config.vector_limit + limit=config.vector_limit, + score_threshold=self.settings.vector_similarity_threshold ) results = [ { @@ -313,11 +320,19 @@ JSON:""" async def graph_search(): start = time.time() try: + # Skip synonyms for graph search - only use core keywords + # Synonyms like "author" can match unrelated entities like "author2000" + graph_keywords = { + "core_keywords": keywords_data.get("core_keywords", []), + "entities": keywords_data.get("entities", []), + "synonyms": {}, # No synonyms for exact entity matching + "expansions": {} + } results = await self.graph.search_documents( query=query, user=user, limit=config.graph_limit, - keywords_data=keywords_data + keywords_data=graph_keywords ) formatted = [ { @@ -394,52 +409,134 @@ JSON:""" return output - def _reciprocal_rank_fusion( + def _merge_wiki_sources( self, - results_by_source: Dict[str, List], + vector_results: List[Dict], + graph_results: List[Dict], k: int = 60 ) -> List[Dict[str, Any]]: """ - Phase 2: Merge results using Reciprocal Rank Fusion. + Stage 1: Merge vector and graph into single wiki ranking using RRF. - RRF formula: score = sum(1 / (k + rank)) for each source + Both sources search the same wiki pool, so we combine them before + final RRF with web to avoid double-counting wiki pages. Args: - results_by_source: Results from each source + vector_results: Results from vector search + graph_results: Results from graph search k: RRF constant (default 60) Returns: - Merged and sorted results + Merged wiki results sorted by wiki RRF score + """ + wiki_scores = {} + + # Process vector results + for rank, result in enumerate(vector_results, start=1): + page_id = result.get("page_id") + if not page_id: + continue + result_id = f"page_{page_id}" + + if result_id not in wiki_scores: + wiki_scores[result_id] = { + "result": dict(result), # Copy to avoid mutation + "wiki_rrf_score": 0.0, + "found_by": [] + } + + wiki_scores[result_id]["wiki_rrf_score"] += 1 / (k + rank) + wiki_scores[result_id]["found_by"].append("vector") + + # Process graph results + for rank, result in enumerate(graph_results, start=1): + page_id = result.get("page_id") + if not page_id: + continue + result_id = f"page_{page_id}" + + if result_id not in wiki_scores: + wiki_scores[result_id] = { + "result": dict(result), + "wiki_rrf_score": 0.0, + "found_by": [] + } + + wiki_scores[result_id]["wiki_rrf_score"] += 1 / (k + rank) + wiki_scores[result_id]["found_by"].append("graph") + + # Add graph metadata to existing result + wiki_scores[result_id]["result"]["entity_matches"] = result.get("entity_matches") + wiki_scores[result_id]["result"]["matched_entities"] = result.get("matched_entities") + + # Sort by wiki RRF score + sorted_wiki = sorted( + wiki_scores.values(), + key=lambda x: x["wiki_rrf_score"], + reverse=True + ) + + # Return merged results with wiki ranking + merged = [] + for wiki_rank, item in enumerate(sorted_wiki, start=1): + merged.append({ + **item["result"], + "wiki_rank": wiki_rank, + "wiki_rrf_score": item["wiki_rrf_score"], + "found_by": item["found_by"], + "source": "wiki" + }) + + logger.info(f"Wiki merge: {len(merged)} unique pages from vector+graph") + return merged + + def _reciprocal_rank_fusion( + self, + wiki_results: List[Dict], + web_results: List[Dict], + k: int = 60 + ) -> List[Dict[str, Any]]: + """ + Stage 2: Final RRF between wiki (single source) and web. + + Wiki results are pre-merged from vector+graph, so wiki and web + now compete on equal footing. + + Args: + wiki_results: Pre-merged wiki results from _merge_wiki_sources() + web_results: Results from web search + k: RRF constant (default 60) + + Returns: + Final merged and sorted results """ rrf_scores = {} - for source, results in results_by_source.items(): - for rank, result in enumerate(results, start=1): - # Use page_id for wiki results, url hash for web results - if result.get("page_id"): - result_id = f"page_{result['page_id']}" - elif result.get("url"): - result_id = f"url_{hash(result['url'])}" - else: - continue # Skip results without ID + # Wiki results (single source, already merged) + for rank, result in enumerate(wiki_results, start=1): + page_id = result.get("page_id") + if not page_id: + continue + result_id = f"page_{page_id}" + rrf_scores[result_id] = { + "result": result, + "rrf_score": 1 / (k + rank), + "sources": result.get("found_by", ["wiki"]), + "source_type": "wiki" + } - if result_id not in rrf_scores: - rrf_scores[result_id] = { - "result": result, - "rrf_score": 0.0, - "sources": [], - "source_type": source - } - - # RRF formula: sum of 1/(k + rank) across sources - rrf_scores[result_id]["rrf_score"] += 1 / (k + rank) - rrf_scores[result_id]["sources"].append(source) - - # If result appears in multiple sources, update source_type - if len(rrf_scores[result_id]["sources"]) > 1: - rrf_scores[result_id]["source_type"] = "+".join( - sorted(set(rrf_scores[result_id]["sources"])) - ) + # Web results (single source) + for rank, result in enumerate(web_results, start=1): + url = result.get("url") + if not url: + continue + result_id = f"url_{hash(url)}" + rrf_scores[result_id] = { + "result": result, + "rrf_score": 1 / (k + rank), + "sources": ["web"], + "source_type": "web" + } # Sort by RRF score descending sorted_results = sorted( @@ -448,7 +545,7 @@ JSON:""" reverse=True ) - logger.info(f"RRF fusion: {len(sorted_results)} unique results from {len(results_by_source)} sources") + logger.info(f"Final RRF: {len(sorted_results)} results (wiki + web)") return sorted_results diff --git a/tests/test_hybrid_rag.py b/tests/test_hybrid_rag.py index 75e5d88..19a517b 100644 --- a/tests/test_hybrid_rag.py +++ b/tests/test_hybrid_rag.py @@ -219,53 +219,54 @@ async def test_vector_data(vector_service, test_wiki_page): # ============================================================================ class TestRRFFusion: - """Test Reciprocal Rank Fusion algorithm.""" + """Test two-stage Reciprocal Rank Fusion algorithm.""" - def test_rrf_single_source(self, hybrid_rag_service): - """Test RRF with single source.""" - results_by_source = { - "vector": [ - {"page_id": 1, "title": "Doc 1", "content": "test"}, - {"page_id": 2, "title": "Doc 2", "content": "test"} - ] - } + def test_wiki_merge_single_source(self, hybrid_rag_service): + """Test wiki merge with single source (vector only).""" + vector_results = [ + {"page_id": 1, "title": "Doc 1", "content": "test"}, + {"page_id": 2, "title": "Doc 2", "content": "test"} + ] - fused = hybrid_rag_service._reciprocal_rank_fusion(results_by_source, k=60) + merged = hybrid_rag_service._merge_wiki_sources(vector_results, [], k=60) - assert len(fused) == 2 - assert fused[0]["rrf_score"] > fused[1]["rrf_score"] # Rank 1 > Rank 2 - assert fused[0]["sources"] == ["vector"] + assert len(merged) == 2 + assert merged[0]["wiki_rrf_score"] > merged[1]["wiki_rrf_score"] # Rank 1 > Rank 2 + assert merged[0]["found_by"] == ["vector"] - def test_rrf_multiple_sources_same_doc(self, hybrid_rag_service): - """Test RRF with same document from multiple sources.""" - results_by_source = { - "vector": [{"page_id": 1, "title": "Doc 1", "content": "test"}], - "graph": [{"page_id": 1, "title": "Doc 1", "content": ""}], - } + def test_wiki_merge_multiple_sources_same_doc(self, hybrid_rag_service): + """Test wiki merge with same document from vector and graph.""" + vector_results = [{"page_id": 1, "title": "Doc 1", "content": "test"}] + graph_results = [{"page_id": 1, "title": "Doc 1", "content": ""}] - fused = hybrid_rag_service._reciprocal_rank_fusion(results_by_source, k=60) + merged = hybrid_rag_service._merge_wiki_sources(vector_results, graph_results, k=60) - assert len(fused) == 1 # Deduplicated - assert len(fused[0]["sources"]) == 2 # Both sources - assert "vector" in fused[0]["sources"] - assert "graph" in fused[0]["sources"] - # RRF score should be sum: 1/(60+1) + 1/(60+1) + assert len(merged) == 1 # Deduplicated + assert len(merged[0]["found_by"]) == 2 # Both sources + assert "vector" in merged[0]["found_by"] + assert "graph" in merged[0]["found_by"] + # Wiki RRF score should be sum: 1/(60+1) + 1/(60+1) expected_score = 1/61 + 1/61 - assert abs(fused[0]["rrf_score"] - expected_score) < 0.001 + assert abs(merged[0]["wiki_rrf_score"] - expected_score) < 0.001 - def test_rrf_web_results(self, hybrid_rag_service): - """Test RRF with web results (URL-based).""" - results_by_source = { - "web": [ - {"url": "https://example.com/1", "title": "Web 1", "content": "test"}, - {"url": "https://example.com/2", "title": "Web 2", "content": "test"} - ] - } + def test_final_rrf_wiki_and_web(self, hybrid_rag_service): + """Test final RRF between wiki and web results.""" + # Pre-merged wiki results + wiki_results = [ + {"page_id": 1, "title": "Wiki 1", "content": "test", "found_by": ["vector"]} + ] + web_results = [ + {"url": "https://example.com/1", "title": "Web 1", "content": "test"}, + {"url": "https://example.com/2", "title": "Web 2", "content": "test"} + ] - fused = hybrid_rag_service._reciprocal_rank_fusion(results_by_source, k=60) + fused = hybrid_rag_service._reciprocal_rank_fusion(wiki_results, web_results, k=60) - assert len(fused) == 2 - assert fused[0]["result"]["url"] == "https://example.com/1" + assert len(fused) == 3 + # Wiki rank 1 and web rank 1 should have same RRF score + wiki_score = next(r["rrf_score"] for r in fused if r["source_type"] == "wiki") + web_score = next(r["rrf_score"] for r in fused if r["source_type"] == "web") + assert abs(wiki_score - web_score) < 0.001 # Equal footing class TestContextFormatting: