feat(vector_store): add iter_all enumeration and wire up qdrant migration

This commit is contained in:
Zohaib Hassnain
2026-09-02 17:19:19 +05:30
committed by GitHub
parent 3d32254b07
commit d175f894a4
6 changed files with 447 additions and 11 deletions
+23 -8
View File
@@ -5,6 +5,7 @@ This module provides the command-line interface for the Semantica framework,
enabling users to interact with the framework via terminal commands.
"""
import itertools
import json
import os
import sys
@@ -3714,7 +3715,7 @@ def store_stats(cli_ctx: CLIContext, backend: str, fmt: str, local_json: bool) -
_run_with_error_handling(_action)
_MIGRATE_SUPPORTED_BACKENDS = {"faiss", "sqlite", "pgvector"}
_MIGRATE_SUPPORTED_BACKENDS = {"faiss", "sqlite", "pgvector", "qdrant"}
_MIGRATE_BATCH_SIZE = 500
@@ -3759,12 +3760,10 @@ def store_migrate(cli_ctx: CLIContext, from_backend: str, to_backend: str,
namespace: Optional[str], local_dry: bool, local_json: bool) -> None:
"""Migrate data between backends.
Direct migration is only wired up between faiss, sqlite, and pgvector -
these are the backends whose storage contract supports paging through
every stored vector. Migrating to or from qdrant, pinecone, milvus, or
weaviate still needs the export/reindex workaround below, since each of
those needs its own enumeration design (Qdrant scroll, Pinecone list,
etc.) that hasn't been built yet.
Wired up between faiss, sqlite, pgvector and qdrant. Migrating to or from
pinecone, milvus or weaviate still needs the export/reindex workaround
below, since each of those needs its own enumeration design (Pinecone
list, Milvus query_iterator, Weaviate cursor) that is not built yet.
\b
Example:
@@ -3804,7 +3803,18 @@ def store_migrate(cli_ctx: CLIContext, from_backend: str, to_backend: str,
if source_index_path:
source._backend_store.load_index(source_index_path)
# Pull the first record before building the destination: qdrant, milvus
# and weaviate expose no `dimension` attribute, so without this the
# destination falls back to the default 768 and rejects every insert.
# The data itself is the reliable source of truth.
source_iter = source.iter_vectors(batch_size=_MIGRATE_BATCH_SIZE)
first_item = next(source_iter, None)
source_dimension = getattr(source._backend_store, "dimension", None)
if not source_dimension and first_item is not None:
first_vector = first_item.get("vector")
if first_vector is not None:
source_dimension = len(first_vector)
if source_dimension and "dimension" not in dest_cfg:
dest_cfg["dimension"] = source_dimension
@@ -3827,7 +3837,12 @@ def store_migrate(cli_ctx: CLIContext, from_backend: str, to_backend: str,
metadata_batch.clear()
ids_batch.clear()
for item in source.iter_vectors(batch_size=_MIGRATE_BATCH_SIZE):
items = (
source_iter
if first_item is None
else itertools.chain([first_item], source_iter)
)
for item in items:
meta = dict(item.get("metadata") or {})
if namespace and "namespace" not in meta:
meta["namespace"] = namespace
+56
View File
@@ -604,6 +604,62 @@ class QdrantStore:
self.logger.warning(f"Failed to scroll Qdrant points by metadata filter: {e}")
return []
def iter_all(self, batch_size: int = 500):
"""
Iterate over every stored point using Qdrant's native scroll cursor.
Qdrant paginates by point-ID cursor, not by row offset, so this is
exposed instead of scan_vectors(offset, limit). An integer passed to
scroll()'s offset is a point ID rather than a rank, so there is no way
to seek to "the Nth record" without walking from the start.
VectorStore.iter_vectors() prefers this method when it is present.
Assumes a single unnamed vector per point, matching how insert_vectors()
writes them and how get_vector() reads them back. Collections configured
with named or multi-vectors are not handled here.
Args:
batch_size: Points to request per scroll call
Yields:
Result dicts with 'id', 'metadata', and 'vector', in scroll order
Raises:
ProcessingError: If the collection or client is not initialized.
Errors are raised rather than swallowed because a scan that
silently yields nothing is indistinguishable from an empty
source, which would let a caller such as `store migrate`
report success having copied nothing (issue #1083).
"""
if self.collection is None or self.client is None or not QDRANT_AVAILABLE:
raise ProcessingError(
"Collection not initialized. Call create_collection() or get_collection() first."
)
next_offset = None
while True:
records, next_offset = self.client.scroll(
collection_name=self.collection.collection_name,
limit=batch_size,
offset=next_offset,
with_payload=True,
with_vectors=True,
)
for rec in records:
yield {
"id": str(rec.id),
"metadata": rec.payload or {},
"vector": np.array(rec.vector) if rec.vector is not None else None,
}
# The final page can carry records while already reporting no next
# cursor, so those records are yielded above before stopping here.
# Calling scroll() again with offset=None would restart from the
# beginning rather than continue past the end.
if next_offset is None or not records:
return
def delete_vectors(
self, point_ids: List[Union[str, int]], **options
) -> Dict[str, Any]:
+14 -1
View File
@@ -867,12 +867,25 @@ class VectorStore:
"""
Iterate over every stored vector, one page at a time.
Backends whose native pagination is cursor based (Qdrant, Pinecone,
Milvus, Weaviate) cannot honestly implement the positional
scan_vectors(offset, limit) contract, so they expose iter_all()
instead and it is preferred here when present. Backends with real
positional access (inmemory, FAISS, SQLite-vec, PgVector) fall
through to the offset loop below.
Args:
batch_size: Number of vectors to fetch per underlying scan_vectors() call
batch_size: Number of vectors to fetch per underlying call
Yields:
Result dicts with 'id', 'metadata', and 'vector', in scan order
"""
if self.backend != "inmemory" and self._backend_store is not None:
iter_all = getattr(self._backend_store, "iter_all", None)
if callable(iter_all):
yield from iter_all(batch_size=batch_size)
return
offset = 0
while True:
page = self.scan_vectors(offset=offset, limit=batch_size)
+94 -2
View File
@@ -1360,10 +1360,11 @@ class TestStore:
assert result.exit_code != 0
def test_migrate_refuses_unsupported_backend_pair(self, runner):
# qdrant is wired up now, so weaviate stands in as the unsupported side.
result = runner.invoke(cli_module.main, ["store", "migrate",
"--from", "faiss", "--to", "qdrant"])
"--from", "faiss", "--to", "weaviate"])
assert result.exit_code != 0
assert "faiss, pgvector, sqlite" in result.output
assert "faiss, pgvector, qdrant, sqlite" in result.output
def _fake_migrate_store_module(self, source_items, stored, dest_configs=None):
class _FakeBackendStore:
@@ -1445,6 +1446,97 @@ class TestStore:
assert result.exit_code != 0
assert "index_path" in result.output
def _fake_dimensionless_store_module(self, source_items, dest_configs, stored):
"""Fake VectorStore whose backend store has NO `dimension` attribute.
Matches qdrant/milvus/weaviate, none of which expose one.
"""
class _DimensionlessBackendStore:
pass
class _FakeStore:
def __init__(self, backend, config=None, **kw):
self.backend = backend
self._config = config or {}
self._backend_store = _DimensionlessBackendStore()
dest_configs[backend] = dict(self._config)
def iter_vectors(self, batch_size=500):
if self.backend == "qdrant":
yield from source_items
return
return
yield # pragma: no cover - makes this a generator
def store_vectors(self, vectors, metadata, ids=None):
for vec_id, vec in zip(ids, vectors):
stored[vec_id] = vec
return _fake_module(VectorStore=_FakeStore)
def test_migrate_infers_dimension_from_first_vector(self, runner, monkeypatch):
source_items = [
{"id": "a", "vector": [0.1, 0.2, 0.3], "metadata": {}},
{"id": "b", "vector": [0.4, 0.5, 0.6], "metadata": {}},
]
dest_configs, stored = {}, {}
fake_vs = self._fake_dimensionless_store_module(source_items, dest_configs, stored)
monkeypatch.setitem(__import__("sys").modules, "semantica.vector_store", fake_vs)
result = runner.invoke(cli_module.main, ["store", "migrate",
"--from", "qdrant", "--to", "sqlite", "--json"])
_ok(result)
assert dest_configs["sqlite"].get("dimension") == 3
def test_migrate_does_not_drop_the_peeked_first_record(self, runner, monkeypatch):
"""The first record is consumed to infer dimension, so it must be
chained back into the migration loop rather than lost."""
source_items = [
{"id": "a", "vector": [0.1, 0.2, 0.3], "metadata": {}},
{"id": "b", "vector": [0.4, 0.5, 0.6], "metadata": {}},
]
dest_configs, stored = {}, {}
fake_vs = self._fake_dimensionless_store_module(source_items, dest_configs, stored)
monkeypatch.setitem(__import__("sys").modules, "semantica.vector_store", fake_vs)
result = runner.invoke(cli_module.main, ["store", "migrate",
"--from", "qdrant", "--to", "sqlite", "--json"])
_ok(result)
assert _json_output(result)["migrated"] == 2
assert sorted(stored) == ["a", "b"]
def test_migrate_keeps_explicit_dest_dimension_over_inference(self, runner, monkeypatch):
source_items = [{"id": "a", "vector": [0.1, 0.2, 0.3], "metadata": {}}]
dest_configs, stored = {}, {}
fake_vs = self._fake_dimensionless_store_module(source_items, dest_configs, stored)
monkeypatch.setitem(__import__("sys").modules, "semantica.vector_store", fake_vs)
monkeypatch.setattr(
cli_module.Config, "to_dict",
lambda self: {"vector_store": {"qdrant": {}, "sqlite": {"dimension": 128}}},
)
result = runner.invoke(cli_module.main, ["store", "migrate",
"--from", "qdrant", "--to", "sqlite", "--json"])
_ok(result)
assert dest_configs["sqlite"]["dimension"] == 128
def test_migrate_qdrant_is_now_supported(self, runner, monkeypatch):
dest_configs, stored = {}, {}
fake_vs = self._fake_dimensionless_store_module([], dest_configs, stored)
monkeypatch.setitem(__import__("sys").modules, "semantica.vector_store", fake_vs)
result = runner.invoke(cli_module.main, ["store", "migrate",
"--from", "qdrant", "--to", "pgvector", "--json"])
_ok(result)
assert _json_output(result)["migrated"] == 0
def test_migrate_still_refuses_milvus(self, runner):
result = runner.invoke(cli_module.main, ["store", "migrate",
"--from", "qdrant", "--to", "milvus"])
assert result.exit_code != 0
assert "milvus" in result.output
def test_migrate_faiss_dest_requires_index_path(self, runner, monkeypatch):
fake_vs = _fake_module(VectorStore=lambda **kw: MagicMock())
monkeypatch.setitem(__import__("sys").modules, "semantica.vector_store", fake_vs)
+152
View File
@@ -0,0 +1,152 @@
"""Tests for QdrantStore.iter_all() cursor enumeration.
Qdrant is not installed in this environment, so these drive the real
QdrantStore against a MagicMock standing in for the qdrant_client, following
the pattern already used for qdrant in test_backend_metadata_filtering.py.
"""
from unittest.mock import MagicMock, patch
import numpy as np
import pytest
from semantica.utils.exceptions import ProcessingError
from semantica.vector_store.qdrant_store import QdrantStore
def _record(point_id, payload=None, vector=None):
"""Build a stand-in for a qdrant_client Record."""
rec = MagicMock()
rec.id = point_id
rec.payload = payload
rec.vector = vector
return rec
def _store_with_scroll(*pages):
"""QdrantStore whose client.scroll() returns the given (records, cursor) pages."""
store = QdrantStore()
store.client = MagicMock()
store.client.scroll.side_effect = list(pages)
store.collection = MagicMock()
store.collection.collection_name = "test_collection"
return store
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_threads_cursor_across_pages():
"""The next call must continue from the previous page's next_page_offset."""
store = _store_with_scroll(
([_record(1), _record(2)], "cursor-1"),
([_record(3)], None),
)
result = list(store.iter_all(batch_size=2))
assert [item["id"] for item in result] == ["1", "2", "3"]
calls = store.client.scroll.call_args_list
assert len(calls) == 2
assert calls[0][1]["offset"] is None
assert calls[0][1]["limit"] == 2
assert calls[1][1]["offset"] == "cursor-1"
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_yields_final_page_that_reports_no_next_cursor():
"""Qdrant can return records and a null cursor on the same page.
Those records must still be yielded. Treating a null cursor as "stop
before this page" would silently drop the tail of every scan.
"""
store = _store_with_scroll(([_record(1), _record(2)], None))
result = list(store.iter_all(batch_size=10))
assert [item["id"] for item in result] == ["1", "2"]
assert store.client.scroll.call_count == 1
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_converts_records_to_the_shared_result_shape():
store = _store_with_scroll(
([_record(7, payload={"tag": "x"}, vector=[0.1, 0.2, 0.3])], None),
)
item = list(store.iter_all())[0]
assert item["id"] == "7"
assert item["metadata"] == {"tag": "x"}
np.testing.assert_allclose(item["vector"], np.array([0.1, 0.2, 0.3]))
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_handles_missing_payload_and_vector():
store = _store_with_scroll(([_record(1, payload=None, vector=None)], None))
item = list(store.iter_all())[0]
assert item["metadata"] == {}
assert item["vector"] is None
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_empty_collection_yields_nothing():
store = _store_with_scroll(([], None))
assert list(store.iter_all()) == []
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_stops_on_empty_page_even_with_a_cursor():
"""Defensive: an empty page ends the scan rather than looping forever."""
store = _store_with_scroll(([], "cursor-that-never-clears"))
assert list(store.iter_all()) == []
assert store.client.scroll.call_count == 1
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_raises_when_collection_not_initialized():
"""Must fail loudly, not yield nothing.
An empty scan is indistinguishable from an empty source, which would let
`store migrate` report success having copied nothing (issue #1083).
"""
store = QdrantStore()
with pytest.raises(ProcessingError, match="Collection not initialized"):
list(store.iter_all())
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", False)
def test_iter_all_raises_when_qdrant_unavailable():
store = QdrantStore()
store.client = MagicMock()
store.collection = MagicMock()
with pytest.raises(ProcessingError):
list(store.iter_all())
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_propagates_scroll_errors():
store = QdrantStore()
store.client = MagicMock()
store.client.scroll.side_effect = RuntimeError("connection reset")
store.collection = MagicMock()
store.collection.collection_name = "test_collection"
with pytest.raises(RuntimeError, match="connection reset"):
list(store.iter_all())
@patch("semantica.vector_store.qdrant_store.QDRANT_AVAILABLE", True)
def test_iter_all_requests_payload_and_vectors():
store = _store_with_scroll(([], None))
list(store.iter_all())
kwargs = store.client.scroll.call_args[1]
assert kwargs["with_payload"] is True
assert kwargs["with_vectors"] is True
assert kwargs["collection_name"] == "test_collection"
@@ -26,6 +26,7 @@ from unittest.mock import MagicMock, patch
import numpy as np
from semantica.utils.exceptions import ProcessingError
from semantica.vector_store.vector_store import VectorStore, VectorManager
@@ -138,6 +139,38 @@ class _NonScanningBackendStore:
"""Fake persistent backend store without any scan capability."""
class _IterAllBackendStore:
"""Fake cursor-based backend store exposing iter_all() but not scan_vectors().
Mirrors qdrant/pinecone/milvus/weaviate, which cannot honour a positional
offset and therefore expose native iteration instead.
"""
def __init__(self, items):
self._items = items
self.batch_sizes = []
def iter_all(self, batch_size=500):
self.batch_sizes.append(batch_size)
for item in self._items:
yield item
def scan_vectors(self, offset=0, limit=100):
raise AssertionError("scan_vectors() must not be called when iter_all() exists")
class _MisShapedIterAllBackendStore:
"""Backend store whose ``iter_all`` attribute is not callable."""
iter_all = 42 # plain attribute, not a method
def __init__(self, items):
self._items = items
def scan_vectors(self, offset=0, limit=100):
return self._items[offset:offset + limit]
class VectorStoreScanVectorsTests(unittest.TestCase):
"""VectorStore.scan_vectors() / iter_vectors() backend-agnostic accessors."""
@@ -192,6 +225,81 @@ class VectorStoreScanVectorsTests(unittest.TestCase):
self.assertEqual(list(store.iter_vectors(batch_size=2)), [])
# ---------------------------------------------------------------------------
# VectorStore.iter_vectors() preference for a native iter_all()
# ---------------------------------------------------------------------------
class VectorStoreIterAllDispatchTests(unittest.TestCase):
"""iter_vectors() prefers a backend's native iter_all() when present.
Cursor-based backends cannot implement scan_vectors(offset, limit)
honestly, so they expose iter_all() instead and iter_vectors() routes to
it rather than walking offsets.
"""
def _persistent_store(self, backend_store, backend_name="qdrant"):
store = VectorStore(backend="inmemory", dimension=2)
store.backend = backend_name
store._backend_store = backend_store
return store
def test_iter_vectors_uses_iter_all_when_available(self):
items = [
{"id": "a", "vector": None, "metadata": {"n": 1}},
{"id": "b", "vector": None, "metadata": {"n": 2}},
]
backend = _IterAllBackendStore(items)
store = self._persistent_store(backend)
self.assertEqual(list(store.iter_vectors(batch_size=7)), items)
def test_iter_vectors_forwards_batch_size_to_iter_all(self):
backend = _IterAllBackendStore([])
store = self._persistent_store(backend)
list(store.iter_vectors(batch_size=32))
self.assertEqual(backend.batch_sizes, [32])
def test_iter_vectors_falls_back_to_scan_vectors_without_iter_all(self):
items = [{"id": "a", "vector": None, "metadata": {}}]
store = self._persistent_store(_ScanningBackendStore(items))
self.assertEqual(list(store.iter_vectors(batch_size=2)), items)
def test_iter_vectors_falls_back_when_iter_all_not_callable(self):
# A mis-shaped adapter exposing a non-callable ``iter_all`` must not be
# invoked; the offset path still has to work. Mirrors the count()
# precedent in _MisShapedBackendStore.
items = [{"id": "a", "vector": None, "metadata": {}}]
store = self._persistent_store(_MisShapedIterAllBackendStore(items))
self.assertEqual(list(store.iter_vectors(batch_size=2)), items)
def test_iter_vectors_inmemory_ignores_iter_all(self):
store = VectorStore(backend="inmemory", dimension=2)
store.store_vectors([np.array([1.0, 0.0])], [{"type": "a"}])
store._backend_store = _IterAllBackendStore([{"id": "wrong"}])
collected = list(store.iter_vectors(batch_size=2))
self.assertEqual([item["metadata"] for item in collected], [{"type": "a"}])
def test_iter_vectors_propagates_iter_all_errors(self):
# A scan that silently yields nothing is indistinguishable from an
# empty source, which would let `store migrate` report success having
# copied nothing (issue #1083).
class _FailingIterAll:
def iter_all(self, batch_size=500):
raise ProcessingError("backend unreachable")
yield # pragma: no cover - makes this a generator
store = self._persistent_store(_FailingIterAll())
with self.assertRaises(ProcessingError):
list(store.iter_vectors(batch_size=2))
# ---------------------------------------------------------------------------
# VectorManager tests — inmemory backend
# ---------------------------------------------------------------------------