From 79a4498487e762d30d0acf29ee1d558dd3e5bad3 Mon Sep 17 00:00:00 2001 From: Thor Whalen <1906276+thorwhalen@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:22:05 +0000 Subject: [PATCH 1/2] fix(store): matrix()/metas() no longer raise while another process writes Closes #85. - put_record writes the vector before the meta (delete_record already removes the meta first), so a listed id always has a vector. - Per-record files are written atomically (hidden temp file + os.replace, the same publish step as the packed cache's sig.json, without fsync). On Windows a replace blocked by a reader is retried, then falls back to the old in-place write. On-disk format is unchanged. - _build_matrix/metas skip a record that vanishes or cannot be decoded mid-read; a failed record is re-read once (a missing meta then means deleted, not a gap). A read that still left records out is returned but neither cached nor published as the packed set, and is logged. - delete_record tolerates a concurrent delete of the same record. Stale-cache follow-up found on the way: #86 (xfail test added). Co-Authored-By: Claude Opus 5 --- ir/store.py | 244 +++++++++++++++++++++++++++---- tests/test_store_concurrency.py | 250 ++++++++++++++++++++++++++++++++ 2 files changed, 463 insertions(+), 31 deletions(-) create mode 100644 tests/test_store_concurrency.py diff --git a/ir/store.py b/ir/store.py index 2c22da4..8587227 100644 --- a/ir/store.py +++ b/ir/store.py @@ -32,6 +32,7 @@ import hashlib import io import json +import logging import os import uuid from collections.abc import Iterator, Mapping, MutableMapping @@ -42,6 +43,8 @@ from .base import Record +logger = logging.getLogger(__name__) + #: Bump when the on-disk packed-matrix layout changes so a stale cache from an #: older ``ir`` is treated as invalid (rebuilt) rather than mis-read. _PACKED_FORMAT = 1 @@ -101,13 +104,117 @@ def _packed_content_sig(ids_json: bytes, metas_json: bytes) -> str: return digest.hexdigest() +#: What reading one per-record file can raise when another process is writing +#: or deleting that record right now: ``KeyError`` (the file vanished between +#: listing and reading), ``EOFError`` (``np.load`` of an empty/truncated +#: ``.npy``), ``ValueError`` (a torn ``.npy`` body, or a torn JSON meta -- +#: ``json.JSONDecodeError`` is a ``ValueError``). +_TORN_RECORD_ERRORS = (KeyError, EOFError, ValueError) + + +class _IncompleteRead(Exception): + """A matrix build skipped records that vanished or were torn mid-read. + + Carries the partial ``(ids, matrix, metas)`` so :meth:`CorpusStore.matrix` + can still answer the query that triggered it, without caching or + publishing a set it knows is missing records. + """ + + def __init__(self, result, skipped): + super().__init__(f"{len(skipped)} record(s) changed while being read") + self.result = result + self.skipped = skipped + + +#: How often a per-record ``os.replace`` is retried when the target is held +#: open (a Windows reader in another process), and the first back-off delay in +#: seconds (doubled on each retry). POSIX never needs a retry. +_REPLACE_ATTEMPTS = 6 +_REPLACE_FIRST_DELAY = 0.01 +_IS_WINDOWS = os.name == "nt" + + +def _write_atomically(path: Path, data: bytes) -> None: + """Write ``data`` to ``path`` so no reader ever sees a partial file. + + The bytes go to a hidden temp file in the same directory (hidden, so the + ``dol`` store listing that directory never shows it as a key), which then + replaces ``path`` in one ``os.replace``: a concurrent reader sees the old + file or the new one, never a torn mix. This is the same publish step as the + packed cache's ``sig.json``, minus the fsync -- per-record writes are + atomic against other processes, not durable against power loss, so a bulk + build does not pay an fsync per record. + + On Windows ``os.replace`` refuses (``PermissionError``) while another + process has the target open for reading; it is retried with a short + back-off, and if the target stays busy the bytes are written in place, as + before this function existed, rather than failing the write. + """ + import time + + tmp = path.with_name(f".{path.name}.{uuid.uuid4().hex}.tmp") + try: + tmp.write_bytes(data) + delay = _REPLACE_FIRST_DELAY + for attempt in range(_REPLACE_ATTEMPTS): + try: + os.replace(tmp, path) + return + except PermissionError: + if not _IS_WINDOWS: + raise + if attempt < _REPLACE_ATTEMPTS - 1: + time.sleep(delay) + delay *= 2 + path.write_bytes(data) + finally: + try: + tmp.unlink() + except OSError: + pass # already replaced into ``path`` (the normal case) + + +class _AtomicFiles(MutableMapping): + """``relative path -> bytes`` file store whose writes are atomic. + + Reads, listing and deletes go through ``dol.Files`` unchanged; only + ``__setitem__`` differs, publishing through :func:`_write_atomically` (and + creating missing sub-directories, as ``dol.mk_dirs_if_missing`` did). + """ + + def __init__(self, rootdir): + import dol + + self.rootdir = str(rootdir) + os.makedirs(self.rootdir, exist_ok=True) + self._files = dol.Files(self.rootdir) + + def __getitem__(self, k): + return self._files[k] + + def __setitem__(self, k, v): + path = Path(self.rootdir, k) + path.parent.mkdir(parents=True, exist_ok=True) + _write_atomically(path, v) + + def __delitem__(self, k): + del self._files[k] + + def __iter__(self): + return iter(self._files) + + def __len__(self): + return len(self._files) + + def __contains__(self, k): + return k in self._files + + def _ndarray_store(rootdir) -> MutableMapping[str, np.ndarray]: - """A ``dol`` file store whose values are float32 ``ndarray``s.""" + """A file store whose values are float32 ``ndarray``s (atomic writes).""" import dol - rootdir = str(rootdir) - os.makedirs(rootdir, exist_ok=True) - files = dol.mk_dirs_if_missing(dol.Files(rootdir)) + files = _AtomicFiles(rootdir) def encode(arr: np.ndarray) -> bytes: buf = io.BytesIO() @@ -121,12 +228,29 @@ def decode(data: bytes) -> np.ndarray: def _json_store(rootdir) -> MutableMapping[str, Any]: - """A ``dol`` file store whose values are JSON objects.""" + """A file store whose values are JSON objects (atomic writes). + + Same on-disk format as ``dol.JsonFiles`` (UTF-8, ``indent=4``), so existing + corpora read unchanged. + """ import dol - rootdir = str(rootdir) - os.makedirs(rootdir, exist_ok=True) - return dol.mk_dirs_if_missing(dol.JsonFiles(rootdir)) + def encode(obj) -> bytes: + return json.dumps(obj, indent=4).encode("utf-8") + + return dol.wrap_kvs( + _AtomicFiles(rootdir), obj_of_data=json.loads, data_of_obj=encode + ) + + +def _normalized_matrix(ids, rows, metas): + """``(ids, row-L2-normalized matrix, metas)``; a ``(0, 0)`` matrix when empty.""" + if not ids: + return ([], np.zeros((0, 0), dtype=np.float32), []) + mat = np.vstack(rows) + norms = np.linalg.norm(mat, axis=1, keepdims=True) + norms[norms == 0] = 1.0 + return (ids, mat / norms, metas) class CorpusStore: @@ -187,7 +311,13 @@ def memory(cls) -> "CorpusStore": # ----- record CRUD ---------------------------------------------------- # def put_record(self, record: Record) -> None: - """Persist *record*'s metadata + vector, invalidating the search matrix.""" + """Persist *record*'s metadata + vector, invalidating the search matrix. + + The vector is written **before** the meta: record ids are listed from + the meta view, so a reader in another process that lists an id always + finds its vector (``delete_record`` removes in the reverse order). + """ + self.vectors[record.id] = np.asarray(record.vector, dtype=np.float32) self.meta[record.id] = { "artifact_id": record.artifact_id, "surface_kind": record.surface_kind, @@ -195,16 +325,20 @@ def put_record(self, record: Record) -> None: "text": record.text, "metadata": dict(record.metadata), } - self.vectors[record.id] = np.asarray(record.vector, dtype=np.float32) self._invalidate_matrix() def delete_record(self, record_id: str) -> None: - """Remove a record's metadata + vector; a missing id is tolerated.""" - self.meta.pop(record_id, None) - try: - del self.vectors[record_id] - except KeyError: - pass + """Remove a record's metadata + vector; a missing id is tolerated. + + The meta goes first, so the id stops being listed before its vector + disappears. Each removal tolerates the file being gone already (another + process may be deleting the same record). + """ + for view in (self.meta, self.vectors): + try: + del view[record_id] + except KeyError: + pass self._invalidate_matrix() def record_ids(self) -> Iterator[str]: @@ -344,6 +478,12 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: with a single memory-mapped read. The packed cache turns a cold reopen from a per-record vector-file storm (thousands of tiny reads) into three file reads; it is cleared by any record write, so it never goes stale. + + Another process may be writing or deleting records while this one + rebuilds. A record that vanishes or is only half-written when read is + left out of the result (it is "not yet written"), and such a partial + result is returned but neither cached nor published, so the next call + reads again. """ if self._matrix_cache is not None: return self._matrix_cache @@ -351,7 +491,10 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: if packed is not None: self._matrix_cache = packed return packed - result = self._build_matrix() + try: + result = self._build_matrix() + except _IncompleteRead as incomplete: + return incomplete.result self._save_packed(result) self._matrix_cache = result return result @@ -363,7 +506,8 @@ def metas(self) -> tuple[list[str], list[dict]]: score on text alone (``mode="lexical"``): they need candidate metadata (text + filter fields) but never the embedding matrix, so they must not pay its I/O. Reuses the in-process or packed cache when present; else - reads only the ``meta`` view (not ``vectors``). + reads only the ``meta`` view (not ``vectors``), skipping a record that + vanishes or is half-written mid-read (as :meth:`matrix` does). """ if self._matrix_cache is not None: ids, _mat, metas = self._matrix_cache @@ -372,8 +516,15 @@ def metas(self) -> tuple[list[str], list[dict]]: if packed is not None: self._matrix_cache = packed return packed[0], packed[2] - ids = list(self.meta) - metas = [self.meta[rid] for rid in ids] + ids: list[str] = [] + metas: list[dict] = [] + for rid in list(self.meta): + try: + meta = self.meta[rid] + except _TORN_RECORD_ERRORS: + continue + ids.append(rid) + metas.append(meta) return ids, metas def _build_matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: @@ -381,20 +532,51 @@ def _build_matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: One pass over the ids reads each record's meta and vector together (the previous implementation iterated the meta view three times). + + Raises :class:`_IncompleteRead` (carrying the partial result) when a + listed record vanished or could not be decoded because another process + was writing or deleting it (see ``_TORN_RECORD_ERRORS``). """ - ids = list(self.meta) - if not ids: - return ([], np.zeros((0, 0), dtype=np.float32), []) + ids: list[str] = [] metas: list[dict] = [] rows: list[np.ndarray] = [] - for rid in ids: - metas.append(self.meta[rid]) - rows.append(np.asarray(self.vectors[rid], dtype=np.float32)) - mat = np.vstack(rows) - norms = np.linalg.norm(mat, axis=1, keepdims=True) - norms[norms == 0] = 1.0 - mat = mat / norms - return (ids, mat, metas) + + def read(rid): + meta = self.meta[rid] + return meta, np.asarray(self.vectors[rid], dtype=np.float32) + + def add(rid, meta, row): + ids.append(rid) + metas.append(meta) + rows.append(row) + + retry: list[str] = [] + for rid in list(self.meta): + try: + add(rid, *read(rid)) + except _TORN_RECORD_ERRORS: + retry.append(rid) + # Second look at what failed: a writer has usually finished by now, and a + # record whose meta is gone was deleted -- consistent with the rest of + # the read, so not a gap. Only what still can't be read is missing. + skipped: list[str] = [] + for rid in retry: + if rid not in self.meta: + continue + try: + add(rid, *read(rid)) + except _TORN_RECORD_ERRORS: + skipped.append(rid) + result = _normalized_matrix(ids, rows, metas) + if skipped: + logger.warning( + "ir: %d record(s) could not be read (being written by another " + "process, or corrupt) and were left out of this search: %s", + len(skipped), + ", ".join(skipped[:5]) + (", ..." if len(skipped) > 5 else ""), + ) + raise _IncompleteRead(result, skipped) + return result # ----- packed-matrix disk cache --------------------------------------- # diff --git a/tests/test_store_concurrency.py b/tests/test_store_concurrency.py new file mode 100644 index 0000000..37cfc28 --- /dev/null +++ b/tests/test_store_concurrency.py @@ -0,0 +1,250 @@ +"""Reading a file-backed corpus while another process writes it (i2mint/ir#85). + +``matrix()``/``metas()`` rebuild from the per-record ``meta``/``vectors`` files +whenever the packed cache is missing, which is exactly while a writer is adding +records (its first write clears the cache). These tests pin down that such a +read never raises: a record is written atomically and vector-first, a record +that vanishes mid-read is treated as deleted, and a record that still can't be +read is left out of that one answer -- which is then neither cached nor +published as the packed set. +""" + +import json +import logging +import subprocess +import sys +import textwrap + +import numpy as np +import pytest + +import ir.store as ir_store +from ir.base import Record +from ir.store import CorpusStore, _json_store, _ndarray_store + + +def _rec(rid, *, vec=(1.0, 0.0, 0.0)): + return Record( + id=rid, + artifact_id=f"art_{rid}", + surface_kind="document", + surface_index=0, + text=f"text of {rid}", + vector=np.asarray(vec, dtype=np.float32), + metadata={}, + ) + + +def _file_store(root): + return CorpusStore( + meta=_json_store(root / "meta"), + vectors=_ndarray_store(root / "vectors"), + ledger=_json_store(root / "ledger"), + config=_json_store(root / "config"), + packed_dir=root / "matrix", + ) + + +def test_put_record_writes_vector_before_meta(): + """A reader lists ids from ``meta``, so a listed id must already have a vector.""" + order = [] + + class Recording(dict): + def __init__(self, name): + super().__init__() + self.name = name + + def __setitem__(self, k, v): + order.append(self.name) + super().__setitem__(k, v) + + store = CorpusStore( + meta=Recording("meta"), vectors=Recording("vectors"), ledger={}, config={} + ) + store.put_record(_rec("r1")) + assert order == ["vectors", "meta"] + + +def test_record_with_meta_but_no_vector_is_left_out_not_published(tmp_path, caplog): + root = tmp_path / "corpus" + store = _file_store(root) + store.put_record(_rec("r1")) + store.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) + (root / "vectors" / "r2").unlink() # meta without vector: unreadable record + + reader = _file_store(root) + with caplog.at_level(logging.WARNING, logger="ir.store"): + ids, mat, metas = reader.matrix() + assert ids == ["r1"] and mat.shape == (1, 3) and len(metas) == 1 + assert "r2" in caplog.text + assert not (root / "matrix" / "sig.json").exists() # partial set not published + assert reader._matrix_cache is None # nor cached: the next call reads again + + store.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) # the writer finishes + ids, _mat, _metas = _file_store(root).matrix() + assert sorted(ids) == ["r1", "r2"] + assert (root / "matrix" / "sig.json").exists() + + +@pytest.mark.parametrize("kind", ["meta", "vectors"]) +def test_torn_record_file_is_skipped_not_raised(tmp_path, kind): + """An empty or truncated file (a pre-atomic writer, a crash) is a miss.""" + root = tmp_path / "corpus" + store = _file_store(root) + store.put_record(_rec("r1")) + store.put_record(_rec("r2")) + path = root / kind / "r2" + path.write_bytes(path.read_bytes()[:7]) + + ids, _mat, metas = _file_store(root).matrix() + assert ids == ["r1"] and len(metas) == 1 + if kind == "meta": + assert _file_store(root).metas()[0] == ["r1"] + + +def test_record_deleted_mid_read_is_not_a_gap(): + """An id listed but gone by the time it is read was deleted: the read is whole.""" + + class VanishingMeta(dict): + def __iter__(self): + return iter([*super().__iter__(), "gone"]) + + store = CorpusStore(meta=VanishingMeta(), vectors={}, ledger={}, config={}) + store.put_record(_rec("r1")) + ids, _mat, _metas = store.matrix() + assert ids == ["r1"] + assert store._matrix_cache is not None # complete, so cached + assert store.metas()[0] == ["r1"] + + +def test_writes_leave_no_temp_files_and_temp_files_are_not_keys(tmp_path): + root = tmp_path / "corpus" + store = _file_store(root) + store.put_record(_rec("r1")) + store.put_record(_rec("r1", vec=(0.0, 1.0, 0.0))) # overwrite + for kind in ("meta", "vectors"): + assert sorted(p.name for p in (root / kind).iterdir()) == ["r1"] + # A temp file left by a writer that crashed before its replace is invisible. + (root / "meta" / ".r2.deadbeef.tmp").write_bytes(b"{") + assert list(store.meta) == ["r1"] + + +def test_json_store_format_is_unchanged(tmp_path): + """Corpora written by ``dol.JsonFiles`` (indent=4, UTF-8) read the same.""" + import dol + + old = dol.JsonFiles(str(tmp_path) + "/") + old["k"] = {"text": "café", "n": [1, 2]} + new = _json_store(tmp_path) + assert new["k"] == {"text": "café", "n": [1, 2]} + raw = (tmp_path / "k").read_bytes() + new["k"] = {"text": "café", "n": [1, 2]} + assert json.loads((tmp_path / "k").read_bytes()) == json.loads(raw) + + +def test_windows_busy_target_retries_then_writes_in_place(tmp_path, monkeypatch): + """On Windows a reader holding the file makes ``os.replace`` refuse; don't fail.""" + calls = [] + + def busy_replace(src, dst): + calls.append(dst) + raise PermissionError("target is open in another process") + + monkeypatch.setattr(ir_store, "_IS_WINDOWS", True) + monkeypatch.setattr(ir_store, "_REPLACE_FIRST_DELAY", 0) + monkeypatch.setattr(ir_store.os, "replace", busy_replace) + target = tmp_path / "f" + ir_store._write_atomically(target, b"payload") + assert target.read_bytes() == b"payload" + assert len(calls) == ir_store._REPLACE_ATTEMPTS + assert [p.name for p in tmp_path.iterdir()] == ["f"] # temp file cleaned up + + +def test_permission_error_off_windows_is_raised(tmp_path, monkeypatch): + def denied(src, dst): + raise PermissionError("denied") + + monkeypatch.setattr(ir_store, "_IS_WINDOWS", False) + monkeypatch.setattr(ir_store.os, "replace", denied) + with pytest.raises(PermissionError): + ir_store._write_atomically(tmp_path / "f", b"x") + assert list(tmp_path.iterdir()) == [] + + +_WRITER = textwrap.dedent( + """ + import sys + from pathlib import Path + import numpy as np + from ir.base import Record + from ir.store import CorpusStore, _json_store, _ndarray_store + + root = Path(sys.argv[1]) + store = CorpusStore( + meta=_json_store(root / "meta"), + vectors=_ndarray_store(root / "vectors"), + ledger=_json_store(root / "ledger"), + config=_json_store(root / "config"), + packed_dir=root / "matrix", + ) + rng = np.random.default_rng(0) + for i in range(int(sys.argv[2])): + rid = f"r{i % 40}" # new records, then overwrites of existing ones + store.put_record(Record( + id=rid, artifact_id="a", surface_kind="document", surface_index=0, + text="x" * int(rng.integers(1, 4000)), metadata={}, + vector=rng.random(128).astype(np.float32), + )) + """ +) + + +def test_reads_while_another_process_writes_never_raise(tmp_path): + """The #85 repro: query a corpus while another process builds it.""" + root = tmp_path / "corpus" + writer = subprocess.Popen( + [sys.executable, "-c", _WRITER, str(root), "400"], + stderr=subprocess.PIPE, + ) + reads = 0 + try: + while writer.poll() is None or reads == 0: + reader = _file_store(root) + ids, mat, metas = reader.matrix() + assert len(ids) == mat.shape[0] == len(metas) + if len(ids): + np.testing.assert_allclose(np.linalg.norm(mat, axis=1), 1, atol=1e-5) + assert len(set(_file_store(root).metas()[0])) <= 40 + reads += 1 + finally: + _out, err = writer.communicate(timeout=120) + assert writer.returncode == 0, err.decode() + assert reads > 0 + # Every record the writer wrote is whole on disk. (Whether a *packed* set a + # reader published mid-build includes them is i2mint/ir#86, below.) + store = _file_store(root) + assert sorted(store.meta) == sorted(f"r{i}" for i in range(40)) + for rid in store.meta: + assert store.get_record(rid).vector.shape == (128,) + + +@pytest.mark.xfail( + strict=True, + reason="i2mint/ir#86: a packed set published mid-build hides later writes", +) +def test_reader_publishing_mid_build_does_not_hide_later_writes(tmp_path): + root = tmp_path / "corpus" + writer = _file_store(root) + writer.put_record(_rec("r1")) + reader = _file_store(root) + real_build = reader._build_matrix + + def build_then_writer_continues(): + result = real_build() # the reader's listing predates r2 + writer.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) + return result + + reader._build_matrix = build_then_writer_continues + reader.matrix() # publishes a set without r2 + ids, _mat, _metas = _file_store(root).matrix() + assert sorted(ids) == ["r1", "r2"] From 04f489bb8769faa08379d4d0a918ac03ddfda29b Mon Sep 17 00:00:00 2001 From: Thor Whalen <1906276+thorwhalen@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:37:38 +0000 Subject: [PATCH 2/2] review: keep writes inside the store root; cache a partial read in-process - _AtomicFiles built the write path with Path(root, key), which lets an absolute key (an artifact id can be an absolute path, e.g. links) replace the root and overwrite a file outside the store, while reads still looked under the root. The path now comes from dol.Files itself (root + key, as reads and deletes use), and a key climbing out of the root is refused. - A partial matrix (records unreadable after the retry) is now cached in-process, still never published, so one damaged record no longer makes every search in that process re-read every record file. The warning says how to clear it. - metas() logs skipped records at debug level. - The subprocess stress test uses two writers and larger records, in three rounds; on master it now fails in about 5 of 6 runs (was 1 of 5). Co-Authored-By: Claude Opus 5 --- ir/store.py | 25 +++++++++++--- tests/test_store_concurrency.py | 61 +++++++++++++++++++++++++-------- 2 files changed, 66 insertions(+), 20 deletions(-) diff --git a/ir/store.py b/ir/store.py index 8587227..4ff696e 100644 --- a/ir/store.py +++ b/ir/store.py @@ -193,7 +193,14 @@ def __getitem__(self, k): return self._files[k] def __setitem__(self, k, v): - path = Path(self.rootdir, k) + # The file path comes from ``dol.Files`` itself (root prefix + key, as a + # string), so a write lands exactly where reads and deletes look. Never + # ``Path(root, k)``: pathlib lets an absolute key replace the root, which + # would write outside the store (an artifact id can be an absolute path). + path = Path(self._files._id_of_key(k)) + root = os.path.normpath(self.rootdir) + if os.path.commonpath([root, os.path.normpath(path)]) != root: + raise KeyError(f"key {k!r} would write outside the store at {root}") path.parent.mkdir(parents=True, exist_ok=True) _write_atomically(path, v) @@ -481,9 +488,9 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: Another process may be writing or deleting records while this one rebuilds. A record that vanishes or is only half-written when read is - left out of the result (it is "not yet written"), and such a partial - result is returned but neither cached nor published, so the next call - reads again. + left out of the result (it is "not yet written"), with a warning; such + a partial result is cached in-process but never published as the + packed set, so a fresh process reads the records again. """ if self._matrix_cache is not None: return self._matrix_cache @@ -494,6 +501,11 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: try: result = self._build_matrix() except _IncompleteRead as incomplete: + # Kept in-process like any build (the in-process cache never sees + # other processes' writes anyway), but not published: a record that + # can't be read -- torn, or damaged for good -- must not become part + # of the set every other process loads. + self._matrix_cache = incomplete.result return incomplete.result self._save_packed(result) self._matrix_cache = result @@ -522,6 +534,7 @@ def metas(self) -> tuple[list[str], list[dict]]: try: meta = self.meta[rid] except _TORN_RECORD_ERRORS: + logger.debug("ir: skipped unreadable meta for record %s", rid) continue ids.append(rid) metas.append(meta) @@ -571,7 +584,9 @@ def add(rid, meta, row): if skipped: logger.warning( "ir: %d record(s) could not be read (being written by another " - "process, or corrupt) and were left out of this search: %s", + "process, or damaged) and were left out of this search: %s. If " + "this persists, re-index or delete those records; until then the " + "packed matrix cache is not written for this corpus.", len(skipped), ", ".join(skipped[:5]) + (", ..." if len(skipped) > 5 else ""), ) diff --git a/tests/test_store_concurrency.py b/tests/test_store_concurrency.py index 37cfc28..5e4c7f8 100644 --- a/tests/test_store_concurrency.py +++ b/tests/test_store_concurrency.py @@ -78,7 +78,7 @@ def test_record_with_meta_but_no_vector_is_left_out_not_published(tmp_path, capl assert ids == ["r1"] and mat.shape == (1, 3) and len(metas) == 1 assert "r2" in caplog.text assert not (root / "matrix" / "sig.json").exists() # partial set not published - assert reader._matrix_cache is None # nor cached: the next call reads again + assert _file_store(root).matrix()[0] == ["r1"] # a fresh process reads again store.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) # the writer finishes ids, _mat, _metas = _file_store(root).matrix() @@ -129,6 +129,25 @@ def test_writes_leave_no_temp_files_and_temp_files_are_not_keys(tmp_path): assert list(store.meta) == ["r1"] +def test_absolute_key_is_written_inside_the_root(tmp_path): + """An artifact id can be an absolute path: it must never replace the root.""" + victim = tmp_path / "victim.md" + victim.write_text("# my source file") + root = tmp_path / "store" + store = _json_store(root) + store[str(victim)] = {"cites": ["x"]} + assert victim.read_text() == "# my source file" + assert store[str(victim)] == {"cites": ["x"]} + assert [p for p in root.rglob("*") if p.is_file()] # stored under the root + + +def test_key_climbing_out_of_the_root_is_refused(tmp_path): + store = _json_store(tmp_path / "store") + with pytest.raises(KeyError, match="outside the store"): + store["a/../../escape"] = {} + assert not (tmp_path / "escape").exists() + + def test_json_store_format_is_unchanged(tmp_path): """Corpora written by ``dol.JsonFiles`` (indent=4, UTF-8) read the same.""" import dol @@ -171,6 +190,11 @@ def denied(src, dst): assert list(tmp_path.iterdir()) == [] +# Large records (~80 KB vector, up to ~200 KB of text) keep each write in +# flight long enough that a reader on the old, non-atomic stores reliably hit a +# torn or half-listed record. +_DIM = 20_000 + _WRITER = textwrap.dedent( """ import sys @@ -180,6 +204,7 @@ def denied(src, dst): from ir.store import CorpusStore, _json_store, _ndarray_store root = Path(sys.argv[1]) + DIM = int(sys.argv[3]) store = CorpusStore( meta=_json_store(root / "meta"), vectors=_ndarray_store(root / "vectors"), @@ -187,28 +212,33 @@ def denied(src, dst): config=_json_store(root / "config"), packed_dir=root / "matrix", ) - rng = np.random.default_rng(0) + first = int(sys.argv[4]) # this writer's ids: r .. r + rng = np.random.default_rng(first) for i in range(int(sys.argv[2])): - rid = f"r{i % 40}" # new records, then overwrites of existing ones + rid = f"r{first + i % 20}" # new records, then overwrites of existing ones store.put_record(Record( id=rid, artifact_id="a", surface_kind="document", surface_index=0, - text="x" * int(rng.integers(1, 4000)), metadata={}, - vector=rng.random(128).astype(np.float32), + text="x" * int(rng.integers(1, 200_000)), metadata={}, + vector=rng.random(DIM).astype(np.float32), )) """ ) -def test_reads_while_another_process_writes_never_raise(tmp_path): - """The #85 repro: query a corpus while another process builds it.""" +@pytest.mark.parametrize("round_", range(3)) # each round alone caught master ~1 in 2 +def test_reads_while_another_process_writes_never_raise(tmp_path, round_): + """The #85 repro: query a corpus while other processes build it.""" root = tmp_path / "corpus" - writer = subprocess.Popen( - [sys.executable, "-c", _WRITER, str(root), "400"], - stderr=subprocess.PIPE, - ) + writers = [ # two writers on disjoint ids, as when a build and a maintain overlap + subprocess.Popen( + [sys.executable, "-c", _WRITER, str(root), "150", str(_DIM), str(first)], + stderr=subprocess.PIPE, + ) + for first in (0, 20) + ] reads = 0 try: - while writer.poll() is None or reads == 0: + while any(w.poll() is None for w in writers) or reads == 0: reader = _file_store(root) ids, mat, metas = reader.matrix() assert len(ids) == mat.shape[0] == len(metas) @@ -217,15 +247,16 @@ def test_reads_while_another_process_writes_never_raise(tmp_path): assert len(set(_file_store(root).metas()[0])) <= 40 reads += 1 finally: - _out, err = writer.communicate(timeout=120) - assert writer.returncode == 0, err.decode() + errs = [w.communicate(timeout=120)[1] for w in writers] + for w, err in zip(writers, errs, strict=True): + assert w.returncode == 0, err.decode() assert reads > 0 # Every record the writer wrote is whole on disk. (Whether a *packed* set a # reader published mid-build includes them is i2mint/ir#86, below.) store = _file_store(root) assert sorted(store.meta) == sorted(f"r{i}" for i in range(40)) for rid in store.meta: - assert store.get_record(rid).vector.shape == (128,) + assert store.get_record(rid).vector.shape == (_DIM,) @pytest.mark.xfail(