From 62070fdc4f5fc1f4d06fef4aead13ac678ab1812 Mon Sep 17 00:00:00 2001 From: Yaroslav Date: Mon, 6 Jul 2026 15:04:38 +0000 Subject: [PATCH] fix: delete async upsert_chunks/delete_file paths; repoint E2E harness to sync The async upsert_chunks and delete_file methods diverged from the sync path, lacking the text-size guard and poison/empty-vector skip. Their only caller was the E2E test harness, making deletion the cleaner fix. Changes: - Delete async upsert_chunks, delete_file and their embed helpers (_embed_documents_local, _embed_sparse, _embed_colbert) from qdrant_writer.py - Remove unused threading import, _dense_semaphore, and try_update_ingestion_trace - Repoint live_harness.py index_fixture_documents to upsert_chunks_sync via asyncio.to_thread (harness is async) - Update contract tests: remove upsert_chunks from parametrize (deleted) - Remove test_embed_colbert tests (tested deleted async-only helper) - Clean up test helper (remove _dense_semaphore injection) Fixes card_9aef142cdfb8 --- src/ingestion/unified/qdrant_writer.py | 199 ------------------ ...t_qdrant_writer_atomic_replace_contract.py | 6 +- tests/e2e_core/live_harness.py | 21 +- tests/unit/ingestion/test_payload_contract.py | 26 --- 4 files changed, 13 insertions(+), 239 deletions(-) diff --git a/src/ingestion/unified/qdrant_writer.py b/src/ingestion/unified/qdrant_writer.py index 3964fa0701..b67a7a5d55 100644 --- a/src/ingestion/unified/qdrant_writer.py +++ b/src/ingestion/unified/qdrant_writer.py @@ -15,7 +15,6 @@ import json as _json import logging -import threading import uuid from dataclasses import dataclass from typing import Any @@ -31,7 +30,6 @@ SparseVector, ) -from src.ingestion.unified.observability import try_update_ingestion_trace from src.retrieval.topic_classifier import classify_chunk_topic, classify_doc_type @@ -92,33 +90,11 @@ def __init__( timeout=bge_m3_timeout, batch_size=self.BGE_M3_BATCH_SIZE, ) - self._dense_semaphore = threading.Semaphore(max(1, bge_m3_concurrency)) logger.info("QdrantHybridWriter BGE-M3 URL: %s", self.bge_m3_url) logger.info("QdrantHybridWriter BGE-M3 timeout: %ss", bge_m3_timeout) logger.info("QdrantHybridWriter dense: local BGE-M3 (concurrency=%d)", bge_m3_concurrency) logger.info("QdrantHybridWriter sparse: BGE-M3 /encode/sparse") - def _embed_sparse(self, texts: list[str]) -> list[dict[str, list]]: - """Generate sparse embeddings via BGEM3SyncClient. - - Returns list of dicts with 'indices' and 'values' keys. - """ - if not texts: - return [] - result = self._bge_client.encode_sparse(texts) - return result.weights - - def _embed_colbert(self, texts: list[str]) -> list[list[list[float]]]: - """Generate ColBERT multivectors via BGEM3SyncClient. - - Returns list of token vector lists (one per document). - Each document's colbert is list[list[float]] (num_tokens x 1024). - """ - if not texts: - return [] - result = self._bge_client.encode_colbert(texts) - return result.colbert_vecs - @staticmethod def _to_sparse_vector(sparse_emb: Any) -> SparseVector: if isinstance(sparse_emb, dict): @@ -165,23 +141,6 @@ def get_chunk_location(chunk: Any, index: int) -> str: # Fallback return f"chunk_{index}" - def _embed_documents_local(self, texts: list[str]) -> list[list[float]]: - """Embed documents using BGEM3SyncClient. - - Args: - texts: List of document texts - - Returns: - List of 1024-dim dense embedding vectors - """ - if not texts: - return [] - # Semaphore covers entire call including internal batching. - # bge_m3_concurrency defaults to 1; ingestion runs sequentially. - with self._dense_semaphore: - result = self._bge_client.encode_dense(texts) - return result.vectors - @staticmethod def _infer_language(source_path: str, file_metadata: dict[str, Any]) -> str: value = file_metadata.get("language") @@ -372,164 +331,6 @@ def _upsert_points_in_batches( return total_upserted - async def delete_file(self, file_id: str, collection_name: str) -> int: - """Delete all points for a file. - - Uses metadata.file_id filter (more reliable than flat file_id). - """ - try_update_ingestion_trace( - command="qdrant-delete-file", - status="started", - metadata={"collection": collection_name}, - ) - # Count before delete - count_result = self.client.count( - collection_name=collection_name, - count_filter=Filter( - must=[FieldCondition(key="metadata.file_id", match=MatchValue(value=file_id))] - ), - ) - count = count_result.count - - if count > 0: - # Delete by metadata.file_id (canonical SDK shape: wrap Filter in FilterSelector) - self.client.delete( - collection_name=collection_name, - points_selector=FilterSelector( - filter=Filter( - must=[ - FieldCondition( - key="metadata.file_id", - match=MatchValue(value=file_id), - ) - ] - ) - ), - ) - logger.info(f"Deleted {count} points for file_id={file_id}") - - try_update_ingestion_trace( - command="qdrant-delete-file", - status="completed", - metadata={"collection": collection_name, "points_deleted": count}, - ) - return count - - async def upsert_chunks( - self, - chunks: list[Any], - file_id: str, - source_path: str, - file_metadata: dict[str, Any], - collection_name: str, - ) -> WriteStats: - """Upsert chunks with atomic replace semantics (#1602). - - 1. Generate embeddings and build replacement points first. - 2. Upsert replacement points with deterministic IDs (overwrites the - same chunk_locations in place). - 3. Only after a successful upsert, delete points whose IDs belong to - ``file_id`` but are not part of the new batch (stale orphans). - - If any step before the stale-id sweep fails, no destructive delete is - executed, so the previous good version of the document remains - searchable. - """ - stats = WriteStats() - try_update_ingestion_trace( - command="qdrant-upsert-chunks", - status="started", - metadata={"collection": collection_name, "chunks": len(chunks)}, - ) - - if not chunks: - try_update_ingestion_trace( - command="qdrant-upsert-chunks", - status="completed", - metadata={"collection": collection_name, "chunks": 0}, - ) - return stats - - try: - # Step 1: Extract texts - texts = [chunk.text for chunk in chunks] - - # Step 2: Generate embeddings (BGE-M3 dense + sparse + colbert). - # Any failure here exits via `except` BEFORE any destructive call. - dense_embeddings = self._embed_documents_local(texts) - sparse_embeddings = self._embed_sparse(texts) - colbert_embeddings = self._embed_colbert(texts) - - # Step 3: Build points - points: list[PointStruct] = [] - new_ids: list[str] = [] - for i, (chunk, dense_vec, sparse_emb) in enumerate( - zip(chunks, dense_embeddings, sparse_embeddings, strict=True) - ): - chunk_location = self.get_chunk_location(chunk, i) - point_id = self.generate_point_id(file_id, chunk_location) - payload = self.build_payload( - chunk, file_id, source_path, chunk_location, file_metadata - ) - - vector_dict: dict = { - "dense": dense_vec, - "bm42": self._to_sparse_vector(sparse_emb), - } - if colbert_embeddings: - vector_dict["colbert"] = colbert_embeddings[i] - - point = PointStruct( - id=point_id, - vector=vector_dict, - payload=payload, - ) - points.append(point) - new_ids.append(point_id) - - # Step 4: Upsert replacement points first (deterministic IDs - # overwrite same-location points, so readers always see a valid - # version of the document during the swap). - stats.points_upserted = self._upsert_points_in_batches( - collection_name=collection_name, - points=points, - source_path=source_path, - ) - - # Step 5: Sweep stale orphans (old points for this file_id that - # are not part of the new batch). Safe because Step 4 succeeded. - stats.points_deleted = self._delete_stale_points_sync( - file_id=file_id, - collection_name=collection_name, - new_ids=set(new_ids), - ) - - logger.info( - f"Upserted {stats.points_upserted} points for {source_path} " - f"(swept {stats.points_deleted} stale points)" - ) - - except Exception as e: - stats.errors = [str(e)] - logger.error(f"Error upserting chunks: {e}", exc_info=True) - try_update_ingestion_trace( - command="qdrant-upsert-chunks", - status="error", - metadata={"collection": collection_name, "error_type": type(e).__name__}, - ) - return stats - - try_update_ingestion_trace( - command="qdrant-upsert-chunks", - status="completed", - metadata={ - "collection": collection_name, - "points_deleted": stats.points_deleted, - "points_upserted": stats.points_upserted, - }, - ) - return stats - def delete_file_sync(self, file_id: str, collection_name: str) -> int: """Sync version of delete_file. diff --git a/tests/contract/test_qdrant_writer_atomic_replace_contract.py b/tests/contract/test_qdrant_writer_atomic_replace_contract.py index dc27021b75..892705a590 100644 --- a/tests/contract/test_qdrant_writer_atomic_replace_contract.py +++ b/tests/contract/test_qdrant_writer_atomic_replace_contract.py @@ -55,7 +55,7 @@ def _self_method_calls(func: ast.AST) -> list[str]: return calls -@pytest.mark.parametrize("func_name", ["upsert_chunks", "upsert_chunks_sync"]) +@pytest.mark.parametrize("func_name", ["upsert_chunks_sync"]) def test_upsert_does_not_call_destructive_delete_helpers(func_name: str) -> None: """``upsert_chunks*`` must not call ``delete_file`` / ``delete_file_sync``.""" func = _load_function(func_name) @@ -72,7 +72,7 @@ def test_upsert_does_not_call_destructive_delete_helpers(func_name: str) -> None ).strip() -@pytest.mark.parametrize("func_name", ["upsert_chunks", "upsert_chunks_sync"]) +@pytest.mark.parametrize("func_name", ["upsert_chunks_sync"]) def test_upsert_calls_post_upsert_stale_sweep(func_name: str) -> None: """``upsert_chunks*`` must perform the stale-id sweep after the upsert.""" func = _load_function(func_name) @@ -87,7 +87,7 @@ def test_upsert_calls_post_upsert_stale_sweep(func_name: str) -> None: ).strip() -@pytest.mark.parametrize("func_name", ["upsert_chunks", "upsert_chunks_sync"]) +@pytest.mark.parametrize("func_name", ["upsert_chunks_sync"]) def test_stale_sweep_runs_after_upsert(func_name: str) -> None: """Inside the function body, ``_upsert_points_in_batches`` must precede ``_delete_stale_points_sync``.""" diff --git a/tests/e2e_core/live_harness.py b/tests/e2e_core/live_harness.py index 503ce0c39c..094cf9e997 100644 --- a/tests/e2e_core/live_harness.py +++ b/tests/e2e_core/live_harness.py @@ -407,14 +407,12 @@ async def index_fixture_documents( ) -> int: """Index all Markdown fixture docs into the ephemeral Qdrant collection.""" - from src.ingestion.unified.qdrant_writer import QdrantHybridWriter + import asyncio - class _NoColbertQdrantHybridWriter(QdrantHybridWriter): - def _embed_colbert(self, texts: list[str]) -> list[list[list[float]]]: - return [] + from src.ingestion.unified.qdrant_writer import QdrantHybridWriter selected_ids = set(document_ids or []) - writer = _NoColbertQdrantHybridWriter( + writer = QdrantHybridWriter( qdrant_url=env.qdrant_url, qdrant_api_key=env.qdrant_api_key, bge_m3_url=env.bge_m3_url, @@ -431,17 +429,18 @@ def _embed_colbert(self, texts: list[str]) -> list[list[list[float]]]: document_name=path.name, extra_metadata={"headings": [_extract_title(path)]}, ) - stats = await writer.upsert_chunks( - chunks=[chunk], - file_id=path.stem, - source_path=str(path), - file_metadata={ + stats = await asyncio.to_thread( + writer.upsert_chunks_sync, + [chunk], + path.stem, + str(path), + { "file_name": path.name, "mime_type": "text/markdown", "language": "en", "source_type": "fixture", }, - collection_name=collection_name, + collection_name, ) total += stats.points_upserted finally: diff --git a/tests/unit/ingestion/test_payload_contract.py b/tests/unit/ingestion/test_payload_contract.py index f3f22f7158..2381293e9d 100644 --- a/tests/unit/ingestion/test_payload_contract.py +++ b/tests/unit/ingestion/test_payload_contract.py @@ -3,7 +3,6 @@ from __future__ import annotations -import threading from unittest.mock import MagicMock import pytest @@ -135,37 +134,12 @@ def _make_writer_with_mocks(mock_bge_client: MagicMock) -> QdrantHybridWriter: """Create QdrantHybridWriter bypassing __init__ and inject mock BGE client.""" writer = QdrantHybridWriter.__new__(QdrantHybridWriter) writer._bge_client = mock_bge_client - writer._dense_semaphore = threading.Semaphore(1) return writer class TestColbertVectorInUpsert: """Writer includes colbert multivectors in upserted points.""" - def test_embed_colbert_returns_nested_list(self): - """_embed_colbert returns list[list[list[float]]] from bge_client.""" - mock_bge = MagicMock() - colbert_result = MagicMock() - colbert_result.colbert_vecs = [[[0.1] * 1024, [0.2] * 1024]] - mock_bge.encode_colbert.return_value = colbert_result - - writer = _make_writer_with_mocks(mock_bge) - - result = writer._embed_colbert(["hello"]) - - mock_bge.encode_colbert.assert_called_once_with(["hello"]) - assert result == [[[0.1] * 1024, [0.2] * 1024]] - - def test_embed_colbert_empty_returns_empty(self): - """_embed_colbert with empty input returns [] without calling bge_client.""" - mock_bge = MagicMock() - writer = _make_writer_with_mocks(mock_bge) - - result = writer._embed_colbert([]) - - mock_bge.encode_colbert.assert_not_called() - assert result == [] - def test_upsert_chunks_sync_point_has_colbert_vector(self): """upsert_chunks_sync points include 'colbert' key when use_local_embeddings=True.""" from qdrant_client.models import PointStruct