diff --git a/.claude/skills/pre-push-review/SKILL.md b/.claude/skills/pre-push-review/SKILL.md index efc906b0..26883d49 100644 --- a/.claude/skills/pre-push-review/SKILL.md +++ b/.claude/skills/pre-push-review/SKILL.md @@ -7,7 +7,6 @@ description: | in this repo's automated PR reviews. Use when the user is about to push, says "ready to push", "review my work", "check before PR", or invokes /pre-push-review. Report-only — does not modify code. -model: sonnet allowed-tools: - Bash - Read diff --git a/nextcloud_mcp_server/search/context.py b/nextcloud_mcp_server/search/context.py index f818e278..c9baf084 100644 --- a/nextcloud_mcp_server/search/context.py +++ b/nextcloud_mcp_server/search/context.py @@ -369,21 +369,20 @@ async def get_chunk_with_context( # Prefer chunk_index lookup (always-indexed field) when caller supplied it; # fall back to (chunk_start, chunk_end) lookup otherwise. chunk_text: str | None = None - if doc_id: - if chunk_index is not None: - chunk_text = await _get_chunk_by_index_from_qdrant( - user_id, doc_id, doc_type, chunk_index - ) - # Skip the offset fallback for files when the indexed chunk_index - # lookup already ran: chunk_start/end_offset aren't indexed in Qdrant - # Cloud strict mode, so the call returns 400 and surfaces a misleading - # logger.error. The file fast-fail below correctly handles the miss - # without it. - skip_offset_lookup = chunk_index is not None and doc_type == "file" - if chunk_text is None and not skip_offset_lookup: - chunk_text = await _get_chunk_from_qdrant( - user_id, doc_id, doc_type, chunk_start, chunk_end - ) + if chunk_index is not None: + chunk_text = await _get_chunk_by_index_from_qdrant( + user_id, doc_id, doc_type, chunk_index + ) + # Skip the offset fallback for files when the indexed chunk_index + # lookup already ran: chunk_start/end_offset aren't indexed in Qdrant + # Cloud strict mode, so the call returns 400 and surfaces a misleading + # logger.error. The file fast-fail below correctly handles the miss + # without it. + skip_offset_lookup = chunk_index is not None and doc_type == "file" + if chunk_text is None and not skip_offset_lookup: + chunk_text = await _get_chunk_from_qdrant( + user_id, doc_id, doc_type, chunk_start, chunk_end + ) if chunk_text: logger.info( diff --git a/nextcloud_mcp_server/vector/qdrant_client.py b/nextcloud_mcp_server/vector/qdrant_client.py index bd3a3e54..4cdb8434 100644 --- a/nextcloud_mcp_server/vector/qdrant_client.py +++ b/nextcloud_mcp_server/vector/qdrant_client.py @@ -3,6 +3,7 @@ import logging from typing import Any +import anyio from qdrant_client import AsyncQdrantClient, models from qdrant_client.http.exceptions import UnexpectedResponse from qdrant_client.models import ( @@ -43,8 +44,13 @@ _PAYLOAD_INDEX_FIELDS: dict[str, PayloadSchemaType] = { _DOC_ID_BACKFILL_SENTINEL_ID: str = "00000000-0000-0000-0000-d0c1d0d1d0c1" _DOC_ID_BACKFILL_SENTINEL_PAYLOAD: dict[str, str] = {"_migration_marker": "doc_id_v1"} -# Singleton instance +# Singleton instance + init lock. The lock serialises concurrent first +# callers so the idempotent-but-expensive startup migration +# (``_backfill_doc_id_to_string`` + ``_ensure_payload_indexes``) only runs +# once per process. Steady-state callers hit the fast path above the lock +# and never acquire it. _qdrant_client: AsyncQdrantClient | None = None +_qdrant_init_lock: anyio.Lock = anyio.Lock() async def _ensure_payload_indexes( @@ -392,131 +398,151 @@ async def get_qdrant_client() -> AsyncQdrantClient: """ global _qdrant_client - if _qdrant_client is None: - settings = get_settings() + # Fast path: already initialized — skip lock acquisition for the + # steady-state hot path (every MCP tool call after first start). + if _qdrant_client is not None: + return _qdrant_client - # Detect mode and initialize client accordingly - if settings.qdrant_url: - # Network mode - logger.info(f"Using Qdrant network mode: {settings.qdrant_url}") - _qdrant_client = AsyncQdrantClient( - url=settings.qdrant_url, - api_key=settings.qdrant_api_key, - timeout=30, - ) - elif settings.qdrant_location: - # Local mode (either :memory: or persistent path) - if settings.qdrant_location == ":memory:": - logger.info("Using Qdrant in-memory mode: :memory:") + # Slow path: serialise concurrent first-callers so the idempotent-but- + # expensive startup migration (``_backfill_doc_id_to_string`` + + # ``_ensure_payload_indexes``) runs exactly once. Without this lock, + # parallel cold-start callers would all enter the init block, run the + # migration N times, and emit duplicate "skip-because-exists" warnings + # from the index helper — annoying log noise but not data corruption. + async with _qdrant_init_lock: + # Double-checked: another waiter may have initialized while we + # blocked on the lock. + if _qdrant_client is None: + settings = get_settings() + + # Detect mode and initialize client accordingly + if settings.qdrant_url: + # Network mode + logger.info(f"Using Qdrant network mode: {settings.qdrant_url}") + _qdrant_client = AsyncQdrantClient( + url=settings.qdrant_url, + api_key=settings.qdrant_api_key, + timeout=30, + ) + elif settings.qdrant_location: + # Local mode (either :memory: or persistent path) + if settings.qdrant_location == ":memory:": + logger.info("Using Qdrant in-memory mode: :memory:") + _qdrant_client = AsyncQdrantClient(":memory:") + else: + # Persistent local mode - use path parameter + logger.info( + f"Using Qdrant persistent mode: {settings.qdrant_location}" + ) + _qdrant_client = AsyncQdrantClient(path=settings.qdrant_location) + else: + # Should not happen due to __post_init__ validation, but handle gracefully + logger.warning("No Qdrant mode configured, defaulting to :memory:") _qdrant_client = AsyncQdrantClient(":memory:") - else: - # Persistent local mode - use path parameter - logger.info(f"Using Qdrant persistent mode: {settings.qdrant_location}") - _qdrant_client = AsyncQdrantClient(path=settings.qdrant_location) - else: - # Should not happen due to __post_init__ validation, but handle gracefully - logger.warning("No Qdrant mode configured, defaulting to :memory:") - _qdrant_client = AsyncQdrantClient(":memory:") - # Get collection name (auto-generated from deployment ID + model) - collection_name = settings.get_collection_name() + # Get collection name (auto-generated from deployment ID + model) + collection_name = settings.get_collection_name() - embedding_service = get_embedding_service() + embedding_service = get_embedding_service() - # Detect dimension dynamically (for OllamaEmbeddingProvider) - if hasattr(embedding_service.provider, "_detect_dimension"): - await embedding_service.provider._detect_dimension() # type: ignore[call-non-callable] + # Detect dimension dynamically (for OllamaEmbeddingProvider) + if hasattr(embedding_service.provider, "_detect_dimension"): + await embedding_service.provider._detect_dimension() # type: ignore[call-non-callable] - expected_dimension = embedding_service.get_dimension() + expected_dimension = embedding_service.get_dimension() - # Explicitly check if collection exists - logger.debug(f"Checking if collection '{collection_name}' exists...") - collections = await _qdrant_client.get_collections() - collection_names = [c.name for c in collections.collections] + # Explicitly check if collection exists + logger.debug(f"Checking if collection '{collection_name}' exists...") + collections = await _qdrant_client.get_collections() + collection_names = [c.name for c in collections.collections] - if collection_name in collection_names: - # Collection exists - validate dimensions - logger.debug( - f"Collection '{collection_name}' found, validating dimensions..." - ) - collection_info = await _qdrant_client.get_collection(collection_name) - # Handle both named vectors (dict) and legacy single vector - vectors = collection_info.config.params.vectors - if isinstance(vectors, dict): - actual_dimension = vectors["dense"].size - else: - # Type narrowing: vectors must be VectorParams if not dict - assert isinstance(vectors, VectorParams) - actual_dimension = vectors.size + if collection_name in collection_names: + # Collection exists - validate dimensions + logger.debug( + f"Collection '{collection_name}' found, validating dimensions..." + ) + collection_info = await _qdrant_client.get_collection(collection_name) + # Handle both named vectors (dict) and legacy single vector + vectors = collection_info.config.params.vectors + if isinstance(vectors, dict): + actual_dimension = vectors["dense"].size + else: + # Type narrowing: vectors must be VectorParams if not dict + assert isinstance(vectors, VectorParams) + actual_dimension = vectors.size - # Validate dimension matches - if actual_dimension != expected_dimension: - embedding_model = settings.get_embedding_model_name() - raise ValueError( - f"Dimension mismatch for collection '{collection_name}':\n" - f" Expected: {expected_dimension} (from embedding model '{embedding_model}')\n" - f" Found: {actual_dimension}\n" - f"This usually means you changed the embedding model.\n" - f"Solutions:\n" - f" 1. Delete the old collection: Collection will be recreated with new dimensions\n" - f" 2. Set QDRANT_COLLECTION to use a different collection name\n" - f" 3. Revert to the original embedding model" + # Validate dimension matches + if actual_dimension != expected_dimension: + embedding_model = settings.get_embedding_model_name() + raise ValueError( + f"Dimension mismatch for collection '{collection_name}':\n" + f" Expected: {expected_dimension} (from embedding model '{embedding_model}')\n" + f" Found: {actual_dimension}\n" + f"This usually means you changed the embedding model.\n" + f"Solutions:\n" + f" 1. Delete the old collection: Collection will be recreated with new dimensions\n" + f" 2. Set QDRANT_COLLECTION to use a different collection name\n" + f" 3. Revert to the original embedding model" + ) + + logger.info( + f"Using existing Qdrant collection: {collection_name} " + f"(dimension={actual_dimension}, model={settings.get_embedding_model_name()})" ) - logger.info( - f"Using existing Qdrant collection: {collection_name} " - f"(dimension={actual_dimension}, model={settings.get_embedding_model_name()})" - ) + # Existing collections may pre-date the doc_id normalization / + # payload-index work. Backfill before creating the index so the + # index covers every point. Pass the already-fetched + # collection_info.payload_schema through to avoid a redundant + # get_collection round-trip on every restart. + await _backfill_doc_id_to_string( + _qdrant_client, collection_name, expected_dimension + ) + await _ensure_payload_indexes( + _qdrant_client, + collection_name, + existing_schema=collection_info.payload_schema or {}, + ) - # Existing collections may pre-date the doc_id normalization / - # payload-index work. Backfill before creating the index so the - # index covers every point. Pass the already-fetched - # collection_info.payload_schema through to avoid a redundant - # get_collection round-trip on every restart. - await _backfill_doc_id_to_string( - _qdrant_client, collection_name, expected_dimension - ) - await _ensure_payload_indexes( - _qdrant_client, - collection_name, - existing_schema=collection_info.payload_schema or {}, - ) - - else: - # Collection doesn't exist - create it - embedding_model = settings.get_embedding_model_name() - logger.info( - f"Collection '{collection_name}' not found, creating with " - f"dimension={expected_dimension}, model={embedding_model}..." - ) - await _qdrant_client.create_collection( - collection_name=collection_name, - vectors_config={ - "dense": VectorParams( - size=expected_dimension, - distance=Distance.COSINE, - ), - }, - sparse_vectors_config={ - "sparse": models.SparseVectorParams( - index=models.SparseIndexParams( - on_disk=False, - ) - ), - }, - ) - logger.info( - f"Created Qdrant collection: {collection_name}\n" - f" Dense vector dimension: {expected_dimension}\n" - f" Dense embedding model: {embedding_model}\n" - f" Sparse vectors: BM25 (for hybrid search)\n" - f" Distance: COSINE\n" - f"Background sync will index all documents with dense + sparse vectors." - ) - # Freshly created collection has no payload schema yet; pass {} - # explicitly to skip the otherwise-redundant get_collection call. - await _ensure_payload_indexes( - _qdrant_client, collection_name, existing_schema={} - ) + else: + # Collection doesn't exist - create it + embedding_model = settings.get_embedding_model_name() + logger.info( + f"Collection '{collection_name}' not found, creating with " + f"dimension={expected_dimension}, model={embedding_model}..." + ) + await _qdrant_client.create_collection( + collection_name=collection_name, + vectors_config={ + "dense": VectorParams( + size=expected_dimension, + distance=Distance.COSINE, + ), + }, + sparse_vectors_config={ + "sparse": models.SparseVectorParams( + index=models.SparseIndexParams( + on_disk=False, + ) + ), + }, + ) + logger.info( + f"Created Qdrant collection: {collection_name}\n" + f" Dense vector dimension: {expected_dimension}\n" + f" Dense embedding model: {embedding_model}\n" + f" Sparse vectors: BM25 (for hybrid search)\n" + f" Distance: COSINE\n" + f"Background sync will index all documents with dense + sparse vectors." + ) + # Freshly created collection has no payload schema yet; pass {} + # explicitly to skip the otherwise-redundant get_collection call. + await _ensure_payload_indexes( + _qdrant_client, collection_name, existing_schema={} + ) + # Lock released. ``_qdrant_client`` is guaranteed non-None here: + # either the fast path returned earlier, the lock-protected branch + # set it, or a sibling waiter set it before we got the lock. + assert _qdrant_client is not None return _qdrant_client diff --git a/tests/unit/vector/test_qdrant_client.py b/tests/unit/vector/test_qdrant_client.py index 9479acdc..4743648a 100644 --- a/tests/unit/vector/test_qdrant_client.py +++ b/tests/unit/vector/test_qdrant_client.py @@ -28,6 +28,7 @@ from nextcloud_mcp_server.vector.qdrant_client import ( _PAYLOAD_INDEX_FIELDS, _backfill_doc_id_to_string, _ensure_payload_indexes, + _group_int_doc_ids, ) @@ -597,6 +598,83 @@ async def test_backfill_emits_progress_log_every_20_batches(mocker, caplog): assert "test-collection" in progress_messages[0] +# --------------------------------------------------------------------------- +# _group_int_doc_ids +# --------------------------------------------------------------------------- + + +@pytest.mark.unit +def test_group_int_doc_ids_skips_float_and_warns(caplog): + """A float doc_id is not stringified; it logs WARNING and is skipped. + + Producers always write int or str. A float would round-trip to e.g. + ``"3.0"``, which the keyword index and verification path + (``int(doc_id)``) would never match. Skipping with a loud warning is + the only safe choice. + """ + float_point = SimpleNamespace(id=99, payload={"doc_id": 3.0}) + int_point = SimpleNamespace(id=42, payload={"doc_id": 7}) + + with caplog.at_level("WARNING", logger="nextcloud_mcp_server.vector.qdrant_client"): + by_value, scanned = _group_int_doc_ids([float_point, int_point]) + + # Only the int point made it into by_value; float was dropped. + assert by_value == {"7": [42]} + # Both points still count toward the scanned total — the warning + # should not hide them from progress logs. + assert scanned == 2 + + warnings = [r for r in caplog.records if r.levelname == "WARNING"] + assert len(warnings) == 1 + msg = warnings[0].getMessage() + assert "float" in msg + assert "99" in msg + + +@pytest.mark.unit +def test_group_int_doc_ids_handles_str_and_missing_silently(caplog): + """str / missing doc_id payloads are skipped without warning. + + These are the steady-state paths — already-migrated str values and + sentinel-style points without a doc_id key. Neither should noise up + the log on every restart. + """ + str_point = SimpleNamespace(id=1, payload={"doc_id": "abc"}) + none_payload_point = SimpleNamespace(id=2, payload=None) + missing_key_point = SimpleNamespace(id=3, payload={"other": "value"}) + explicit_none_point = SimpleNamespace(id=4, payload={"doc_id": None}) + + with caplog.at_level("WARNING", logger="nextcloud_mcp_server.vector.qdrant_client"): + by_value, scanned = _group_int_doc_ids( + [str_point, none_payload_point, missing_key_point, explicit_none_point] + ) + + assert by_value == {} + assert scanned == 4 + # No warnings — these paths are expected and silent. + assert not [r for r in caplog.records if r.levelname == "WARNING"] + + +@pytest.mark.unit +def test_group_int_doc_ids_groups_ints_by_str_value(): + """Multiple int-doc_id points sharing a value collapse into one entry. + + Pins the chunk-batching contract: all chunks of one document share its + doc_id, so the helper hands ``_apply_backfill_writes`` a single key + with all chunk point-ids attached. + """ + by_value, scanned = _group_int_doc_ids( + [ + SimpleNamespace(id=10, payload={"doc_id": 42}), + SimpleNamespace(id=11, payload={"doc_id": 42}), + SimpleNamespace(id=12, payload={"doc_id": 7}), + ] + ) + + assert by_value == {"42": [10, 11], "7": [12]} + assert scanned == 3 + + @pytest.mark.unit async def test_ensure_payload_indexes_summarises_failed_fields(mocker, caplog): """A non-400 failure surfaces both as ERROR and a WARNING summary.