diff --git a/README.md b/README.md index 9867c22..83960ee 100644 --- a/README.md +++ b/README.md @@ -47,6 +47,7 @@ Key dataclasses live in `libs/envector/config.py`: - A row deleted moments ago can still take a top-k slot briefly, so a search right after `delete` may return fewer than `k`; pass `fetch_k` to over-fetch. - Filtering happens client-side after the server returns `k` hits, so filtered results can be fewer than `k`; set `fetch_k` (or `IndexSettings.fetch_k`) to over-fetch. - One enVector endpoint per process; all stores in a process must point at the same server. +- `add_texts` / `add_documents` take `use_row_insert`: `True` uses EnVector's single insert, `False` batch insert. By default a call with one document uses single insert and a call with two or more uses batch insert. - Updates wait for that store's pending inserts to merge first; when interleaving inserts and updates, set `WriteSettings.await_insert=True` and use one store instance per index. - `update_documents` / `upsert_documents` above 10,000 items are sent in several batches; if a later batch fails, the earlier ones stay applied. @@ -106,6 +107,9 @@ store.add_texts( ) ``` +`use_row_insert=True` inserts with EnVector's single insert and `False` with batch insert. By default +a call with one document uses single insert and a call with two or more uses batch insert. + ### Similarity search ```python diff --git a/libs/envector/langchain_envector/config.py b/libs/envector/langchain_envector/config.py index f76dbbb..c392399 100644 --- a/libs/envector/langchain_envector/config.py +++ b/libs/envector/langchain_envector/config.py @@ -52,14 +52,19 @@ class WriteSettings: rebuilt, so those wait by default. Set a flag to ``False`` for fire-and-forget bulk work and wait once at the end. - Only the waiting is configured here; the SDK's per-call knobs - (``execute_until``, ``n_workers``, ``use_row_insert``, ...) reach - ``Index.insert`` through ``add_texts(**kwargs)``. + ``use_row_insert`` selects EnVector's insert path for ``add_texts``: + ``True`` single insert, ``False`` batch insert. ``None`` (default) uses + single insert for a call with one document and batch insert for two or + more. ``add_texts`` can override it per call. + + The SDK's other per-call knobs (``execute_until``, ``n_workers``, ...) + reach ``Index.insert`` through ``add_texts(**kwargs)``. """ await_insert: bool = False await_delete: bool = True await_update: bool = True + use_row_insert: Optional[bool] = None # Polling budget for the delete/update/upsert waits. Inserts keep the # SDK's own budget (a day) unless `timeout_s` is passed to add_texts. diff --git a/libs/envector/langchain_envector/vectorstore.py b/libs/envector/langchain_envector/vectorstore.py index 6e4b500..a8b64ea 100644 --- a/libs/envector/langchain_envector/vectorstore.py +++ b/libs/envector/langchain_envector/vectorstore.py @@ -169,6 +169,16 @@ def _loaded_index(self): index.load() return index + def _resolve_row_insert(self, use_row_insert: Optional[bool], n_rows: int) -> bool: + """Pick the insert path: an explicit argument wins, then the store + setting, then the default rule — one row takes single insert, two or + more take batch insert.""" + if use_row_insert is None: + use_row_insert = self.config.write.use_row_insert + if use_row_insert is None: + return n_rows == 1 + return bool(use_row_insert) + # ------------------------------- # VectorStore API # ------------------------------- @@ -181,6 +191,7 @@ def add_texts( vectors: Optional[List[List[float]]] = None, partition_name: Optional[str] = None, await_completion: Optional[bool] = None, + use_row_insert: Optional[bool] = None, **kwargs: Any, ) -> List[str]: """Add texts to the index and return their item IDs. @@ -193,6 +204,13 @@ def add_texts( ``Index.insert`` and apply to the rows this call inserts; the upsert arm below takes only ``timeout_s`` / ``poll_interval_s``. + ``use_row_insert`` selects EnVector's insert path: ``True`` uses single + insert, ``False`` uses batch insert. By default (``None``) a call with + one document uses single insert and a call with two or more documents + uses batch insert; ``WriteSettings.use_row_insert`` changes that + default for the store. Documents updated in place through ``ids`` are + not affected. + ``ids`` follows LangChain's add-or-update contract as far as enVector allows: an entry that is an item ID (int or numeric str, such as the ``Document.id`` search results carry) updates that item in place; an ID @@ -243,7 +261,7 @@ def add_texts( await_completion=await_completion, timeout_s=timeout_s, poll_interval_s=poll_interval_s, - insert_kwargs=kwargs, + insert_kwargs={**kwargs, "use_row_insert": use_row_insert}, ) # Prepare metadata JSON strings per item @@ -263,6 +281,7 @@ def add_texts( partition_name=partition_name, request_ids=request_ids, await_completion=awaited, + use_row_insert=self._resolve_row_insert(use_row_insert, len(texts)), **waits, **kwargs, ) diff --git a/tests/test_vectorstore.py b/tests/test_vectorstore.py index bf4c386..29a21e5 100644 --- a/tests/test_vectorstore.py +++ b/tests/test_vectorstore.py @@ -9,6 +9,7 @@ EnvectorConfig, IndexSettings, KeyConfig, + WriteSettings, ) from langchain_envector.vectorstore import Envector, Document as LC_Document @@ -575,8 +576,9 @@ def test_add_texts_does_not_wait_but_can_be_asked_to(): def test_add_texts_passes_sdk_tuning_knobs_through_kwargs(): - # execute_until / n_workers / use_row_insert are the SDK's own knobs; they - # are not mirrored in WriteSettings, they just travel through **kwargs. + # execute_until / n_workers are the SDK's own knobs; they are not mirrored + # in WriteSettings, they just travel through **kwargs. use_row_insert is a + # named parameter instead, because we decide it against dim. client = FakeClient() store = Envector(config=_cfg(), embeddings=FakeEmbeddings(dim=4), client=client) @@ -587,6 +589,67 @@ def test_add_texts_passes_sdk_tuning_knobs_through_kwargs(): assert call["use_row_insert"] is True +def test_default_is_single_insert_for_one_document_and_batch_for_more(): + client = FakeClient() + store = Envector(config=_cfg(), embeddings=FakeEmbeddings(dim=4), client=client) + + store.add_texts(["t1"]) + store.add_texts(["t1", "t2"]) + store.add_texts(["t1", "t2", "t3", "t4", "t5"]) + assert [c["use_row_insert"] for c in client.index.inserted] == [True, False, False] + + +def test_write_settings_sets_the_store_default_and_a_call_overrides_it(): + client = FakeClient() + off = _cfg() + off.write = WriteSettings(use_row_insert=False) + store = Envector(config=off, embeddings=FakeEmbeddings(dim=4), client=client) + store.add_texts(["t1"]) # store says batch insert, even for one document + assert client.index.inserted[0]["use_row_insert"] is False + store.add_texts(["t1"], use_row_insert=True) # the call wins + assert client.index.inserted[1]["use_row_insert"] is True + + on = _cfg() + on.write = WriteSettings(use_row_insert=True) + store = Envector(config=on, embeddings=FakeEmbeddings(dim=4), client=client) + store.add_texts(["t1", "t2", "t3"]) # store says single insert, even for three + assert client.index.inserted[2]["use_row_insert"] is True + store.add_texts(["t1", "t2", "t3"], use_row_insert=False) + assert client.index.inserted[3]["use_row_insert"] is False + + +def test_add_documents_forwards_use_row_insert(): + from langchain_core.documents import Document + + client = FakeClient() + store = Envector(config=_cfg(), embeddings=FakeEmbeddings(dim=4), client=client) + store.add_documents([Document(page_content="a")]) + store.add_documents([Document(page_content="a"), Document(page_content="b")]) + store.add_documents([Document(page_content="a")], use_row_insert=False) + assert [c["use_row_insert"] for c in client.index.inserted] == [True, False, False] + + +def test_row_insert_reaches_the_reinsert_arm_of_add_or_update(): + # ids= routes through upsert_documents; a row whose id matched nothing is + # re-inserted by add_texts, and that re-insert must keep the chosen path. + class _NotFoundIndex(FakeIndex): + def upsert(self, items, **kw): + result = super().upsert(items, **kw) + result["not_found_item_ids"] = [ + it.item_id for it in items if it.item_id == 99 + ] + return result + + index = _NotFoundIndex() + store = Envector( + config=_cfg(), embeddings=FakeEmbeddings(dim=4), client=FakeClient(index) + ) + with pytest.warns(UserWarning, match="match no live row"): + store.add_texts(["ghost"], ids=[99], use_row_insert=True) + + assert index.inserted[-1]["use_row_insert"] is True + + def test_partition_management_helpers(): client = FakeClient() store = Envector(config=_cfg(), embeddings=FakeEmbeddings(dim=4), client=client)