From 0340b5a9ddcf738672e5795802e0251f026813db Mon Sep 17 00:00:00 2001 From: Thor Whalen <1906276+thorwhalen@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:55:32 +0000 Subject: [PATCH 1/2] store: stamp every record write so a packed set built before it is never published or loaded A reader that rebuilt the packed matrix while another process was writing could publish a set missing the writer's later records; the writer clears the cache only once per session, so every fresh process then served that partial set (#86). Every put_record/delete_record now replaces matrix/write-stamp with a fresh random token, after the record files. matrix() reads the token before listing records, stamps it into sig.json, skips publishing if it changed during the build, and _load_packed treats a sig whose stamp differs from the current one as a miss. A token, not the meta/ directory mtime, so a write in the same clock tick cannot slip through and no quiet period is needed. Sigs from an older ir carry no stamp and stay valid until the first stamped write. Closes #86 Co-Authored-By: Claude Opus 5 --- ir/store.py | 109 +++++++++++++++++++++++++++++--- tests/test_store_concurrency.py | 68 ++++++++++++++++++-- tests/test_store_packed.py | 11 ++-- 3 files changed, 169 insertions(+), 19 deletions(-) diff --git a/ir/store.py b/ir/store.py index 4ff696e..3fb5423 100644 --- a/ir/store.py +++ b/ir/store.py @@ -49,6 +49,18 @@ #: older ``ir`` is treated as invalid (rebuilt) rather than mis-read. _PACKED_FORMAT = 1 +#: File in the packed-cache directory that every record write replaces with a +#: fresh random token (i2mint/ir#86). A matrix build notes the token *before* +#: listing records and stamps it into ``sig.json``; a packed set is only loaded +#: while the token is unchanged. Any write that landed after the build started +#: therefore turns the set into a miss, however close together the two were: a +#: token comparison has no clock resolution to fall through. +_WRITE_STAMP_FILE = "write-stamp" + +#: Default of ``CorpusStore._save_packed(write_stamp=...)``: "the result was +#: built from the corpus as it is now", i.e. stamp it with the current token. +_CURRENT_STAMP = object() + def _write_durably(path: Path, write) -> None: """Create ``path``, let ``write(f)`` fill it, then flush and fsync it.""" @@ -484,7 +496,10 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: normalized-matrix ``.npy`` plus its ids/metas, written once and reloaded 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. + file reads; it is cleared by a writer's first record write, and every + record write replaces a *write stamp* that a packed set must match to be + published or loaded, so a set built before any later write -- by this + process or another -- is never served (i2mint/ir#86). Another process may be writing or deleting records while this one rebuilds. A record that vanishes or is only half-written when read is @@ -498,6 +513,9 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: if packed is not None: self._matrix_cache = packed return packed + # Noted before listing, so any write the build might have missed + # changes it (put/delete replace it after their record files; ir#86). + write_stamp = self._read_write_stamp() try: result = self._build_matrix() except _IncompleteRead as incomplete: @@ -507,7 +525,7 @@ def matrix(self) -> tuple[list[str], np.ndarray, list[dict]]: # of the set every other process loads. self._matrix_cache = incomplete.result return incomplete.result - self._save_packed(result) + self._save_packed(result, write_stamp=write_stamp) self._matrix_cache = result return result @@ -596,17 +614,59 @@ def add(rid, meta, row): # ----- packed-matrix disk cache --------------------------------------- # def _invalidate_matrix(self) -> None: - """Drop the in-process matrix and clear the on-disk packed cache once. - - Called on every record write. The on-disk clear happens at most once per - rebuild (guarded by ``_packed_stale``) so a bulk build's thousands of - ``put_record`` calls don't each touch the filesystem. + """Drop the in-process matrix, stamp the write, clear the packed cache once. + + Called on every record write, *after* the record's files are written. + The on-disk clear happens at most once per rebuild (guarded by + ``_packed_stale``) so a bulk build's thousands of ``put_record`` calls + don't each sweep the cache directory. The write stamp is replaced on + every call: it is what stops another process from publishing, or + loading, a packed set built before this write (i2mint/ir#86) -- the + once-per-session clear cannot, since that process may publish after it. """ self._matrix_cache = None - if self._packed_dir is not None and not self._packed_stale: + if self._packed_dir is None: + return + self._stamp_write() + if not self._packed_stale: self._clear_packed() self._packed_stale = True + def _write_stamp_path(self) -> Path: + return self._packed_dir / _WRITE_STAMP_FILE + + def _stamp_write(self) -> None: + """Replace the write stamp with a fresh token (atomically; one small file). + + If the stamp can't be written, the packed set is removed instead + (``sig.json`` first), so no set built before this write stays loadable. + """ + try: + self._packed_dir.mkdir(parents=True, exist_ok=True) + _write_atomically( + self._write_stamp_path(), uuid.uuid4().hex.encode("ascii") + ) + except OSError as error: + logger.warning( + "ir: could not update the packed-cache write stamp (%s); " + "dropping the packed cache instead", + error, + ) + self._clear_packed() + + def _read_write_stamp(self) -> str | None: + """The current write stamp, or ``None`` if no stamped write happened yet.""" + if self._packed_dir is None: + return None + try: + return self._write_stamp_path().read_text(encoding="ascii") + except FileNotFoundError: + return None + except (OSError, ValueError): + # Unreadable (e.g. mid-replace on Windows): a value no sig carries, + # so a load misses and a build does not publish. + return "" + # Legacy (pre-generation) flat file names. A cache written by an older # ``ir`` is still read through these; new writes never use them. _LEGACY_PACKED_FILES = { @@ -722,11 +782,30 @@ def _load_packed(self): shape = sig.get("shape") if shape is not None and list(mat.shape) != list(shape): return None + # Checked last, after the data files are read: a set built before the + # latest record write may be missing it (i2mint/ir#86). A sig from an + # ``ir`` predating the stamp has none, which matches only a corpus no + # stamping ``ir`` has written to since. + if sig.get("write_stamp") != self._read_write_stamp(): + return None return (ids, mat, metas) - def _save_packed(self, result: tuple[list[str], np.ndarray, list[dict]]) -> None: + def _save_packed( + self, + result: tuple[list[str], np.ndarray, list[dict]], + *, + write_stamp: str | None | object = _CURRENT_STAMP, + ) -> None: """Persist a freshly built matrix to the packed cache (best-effort). + ``write_stamp`` is the write stamp read *before* the build listed its + records. It goes into ``sig.json``, and a set whose stamp is no longer + current is not published at all: a record was written or deleted + while it was being built, so it may not reflect that write + (i2mint/ir#86). Loading checks the same stamp again, which covers a + write that lands between this check and the publish. Left out, the + current stamp is used: the caller vouches the result is up to date. + Skips empty corpora. The matrix, ids and metas go to files named by a fresh generation token that only this call writes, and ``sig.json`` (naming that generation, plus a :func:`_packed_content_sig` digest and @@ -744,6 +823,8 @@ def _save_packed(self, result: tuple[list[str], np.ndarray, list[dict]]) -> None ids, mat, metas = result if not ids: return + if write_stamp is _CURRENT_STAMP: + write_stamp = self._read_write_stamp() generation = uuid.uuid4().hex try: self._packed_dir.mkdir(parents=True, exist_ok=True) @@ -761,10 +842,20 @@ def _save_packed(self, result: tuple[list[str], np.ndarray, list[dict]]) -> None "generation": generation, "content_sig": _packed_content_sig(ids_json, metas_json), "shape": list(arr.shape), + "write_stamp": write_stamp, } ).encode("utf-8") sig_tmp = self._packed_dir / f"sig-{generation}.tmp" _write_durably(sig_tmp, lambda f: f.write(sig_json)) + if self._read_write_stamp() != write_stamp: + # Written to while we built: don't replace a possibly-current + # set with one that may be missing that write. + for path in (sig_tmp, paths["matrix"], paths["ids"], paths["metas"]): + try: + path.unlink() + except OSError: + pass + return os.replace(sig_tmp, paths["sig"]) _fsync_dir(self._packed_dir) self._packed_stale = False diff --git a/tests/test_store_concurrency.py b/tests/test_store_concurrency.py index 922005b..c026328 100644 --- a/tests/test_store_concurrency.py +++ b/tests/test_store_concurrency.py @@ -259,19 +259,17 @@ def test_reads_while_another_process_writes_never_raise(tmp_path, round_): 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.) + # Every record the writer wrote is whole on disk, and a fresh process's + # matrix -- packed set or rebuild -- includes all of them (i2mint/ir#86). 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 == (_DIM,) + assert sorted(_file_store(root).matrix()[0]) == sorted(store.meta) -@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): + """i2mint/ir#86: a set built before a later write is never published or served.""" root = tmp_path / "corpus" writer = _file_store(root) writer.put_record(_rec("r1")) @@ -287,3 +285,61 @@ def build_then_writer_continues(): reader.matrix() # publishes a set without r2 ids, _mat, _metas = _file_store(root).matrix() assert sorted(ids) == ["r1", "r2"] + + +def test_a_set_built_across_a_write_is_not_published(tmp_path): + """The build noted the stamp before the write landed: no publish at all.""" + 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() + writer.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) + return result + + reader._build_matrix = build_then_writer_continues + reader.matrix() + assert not (root / "matrix" / "sig.json").exists() + assert list((root / "matrix").glob("matrix-*.npy")) == [] + + +def test_a_set_published_before_a_write_is_not_loaded(tmp_path): + """Publish, then write from another store with its once-per-session clear + already spent: the write stamp alone must turn the set into a miss.""" + root = tmp_path / "corpus" + writer = _file_store(root) + writer.put_record(_rec("r1")) # spends the writer's one clear + _file_store(root).matrix() # another process publishes {r1} + assert (root / "matrix" / "sig.json").exists() + writer.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) # no clear this time + assert (root / "matrix" / "sig.json").exists() + assert sorted(_file_store(root).matrix()[0]) == ["r1", "r2"] + + +def test_an_unchanged_corpus_keeps_serving_its_packed_set(tmp_path): + """The stamp must not turn every load into a rebuild.""" + root = tmp_path / "corpus" + writer = _file_store(root) + writer.put_record(_rec("r1")) + writer.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) + _file_store(root).matrix() + fresh = _file_store(root) + + def _no_rebuild(): + raise AssertionError("the packed set should have been loaded") + + fresh._build_matrix = _no_rebuild + assert sorted(fresh.matrix()[0]) == ["r1", "r2"] + + +def test_delete_record_also_invalidates_a_published_set(tmp_path): + root = tmp_path / "corpus" + writer = _file_store(root) + writer.put_record(_rec("r1")) + writer.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) + _file_store(root).matrix() + writer.delete_record("r2") + assert sorted(_file_store(root).matrix()[0]) == ["r1"] diff --git a/tests/test_store_packed.py b/tests/test_store_packed.py index aedf15e..3e9b66a 100644 --- a/tests/test_store_packed.py +++ b/tests/test_store_packed.py @@ -242,6 +242,8 @@ def test_legacy_flat_packed_layout_still_loads(tmp_path): (packed / "sig.json").write_text( json.dumps({"format": 1, "count": len(ids)}), encoding="utf-8" ) + # Both written by an older ``ir``, which kept no write stamp (ir#86). + (packed / "write-stamp").unlink() store2, _ = _file_store(tmp_path) @@ -261,10 +263,11 @@ def test_republishing_sweeps_older_generations(tmp_path): for _ in range(3): store._save_packed(result) names = sorted(p.name for p in (root / "matrix").iterdir()) - assert len(names) == 4 and "sig.json" in names + # One generation (3 data files), sig.json, and the write stamp (ir#86). + assert len(names) == 5 and {"sig.json", "write-stamp"} <= set(names) - store.put_record(_rec("r2")) # a write clears the cache entirely - assert list((root / "matrix").iterdir()) == [] + store.put_record(_rec("r2")) # a write clears the cache, keeping the stamp + assert [p.name for p in (root / "matrix").iterdir()] == ["write-stamp"] def test_writer_sees_its_own_write_despite_another_process_publishing(tmp_path): @@ -299,7 +302,7 @@ def test_empty_or_truncated_packed_matrix_is_a_miss_not_an_error(tmp_path): for corrupt in (b"", full[:20], full[:-4]): matrix_file.write_bytes(corrupt) store2, _ = _file_store(tmp_path) - store2._save_packed = lambda result: None # keep the corrupt set on disk + store2._save_packed = lambda result, **kw: None # keep the corrupt set on disk ids2, mat2, _metas2 = store2.matrix() assert sorted(ids2) == ["r1", "r2"] assert np.asarray(mat2).shape == (2, 3) From e5f560016c3dda2aa0eca205fb41af9d0a6421ca Mon Sep 17 00:00:00 2001 From: Thor Whalen <1906276+thorwhalen@users.noreply.github.com> Date: Tue, 22 Sep 2026 16:03:48 +0000 Subject: [PATCH 2/2] store: an unreadable or unwritable write stamp never lets a stale set through - An empty/unreadable stamp (e.g. mid in-place rewrite on Windows) reads as a fresh unique value, so it matches no sig and a build does not publish. - A failed stamp write removes the stamp instead of leaving the old token a concurrent build may already hold. - The packed dir is created once per store, not on every record write. Co-Authored-By: Claude Opus 5 --- ir/store.py | 29 +++++++++++++++++------- tests/test_store_concurrency.py | 39 +++++++++++++++++++++++++++++++++ 2 files changed, 60 insertions(+), 8 deletions(-) diff --git a/ir/store.py b/ir/store.py index 3fb5423..44e9682 100644 --- a/ir/store.py +++ b/ir/store.py @@ -303,6 +303,7 @@ def __init__( # is a *write-invalidated read cache*: any record write clears it. self._packed_dir = Path(packed_dir) if packed_dir is not None else None self._packed_stale = False + self._packed_dir_ready = False # ----- factories ------------------------------------------------------ # @@ -638,20 +639,28 @@ def _write_stamp_path(self) -> Path: def _stamp_write(self) -> None: """Replace the write stamp with a fresh token (atomically; one small file). - If the stamp can't be written, the packed set is removed instead - (``sig.json`` first), so no set built before this write stays loadable. + If the stamp can't be written, it is removed instead (a value no build + in flight can hold, unless it began before any stamped write), and the + packed set with it (``sig.json`` first). """ try: - self._packed_dir.mkdir(parents=True, exist_ok=True) + if not self._packed_dir_ready: + self._packed_dir.mkdir(parents=True, exist_ok=True) + self._packed_dir_ready = True _write_atomically( self._write_stamp_path(), uuid.uuid4().hex.encode("ascii") ) except OSError as error: + self._packed_dir_ready = False logger.warning( "ir: could not update the packed-cache write stamp (%s); " - "dropping the packed cache instead", + "dropping the stamp and the packed cache instead", error, ) + try: + self._write_stamp_path().unlink() + except OSError: + pass self._clear_packed() def _read_write_stamp(self) -> str | None: @@ -659,13 +668,17 @@ def _read_write_stamp(self) -> str | None: if self._packed_dir is None: return None try: - return self._write_stamp_path().read_text(encoding="ascii") + stamp = self._write_stamp_path().read_text(encoding="ascii") except FileNotFoundError: return None except (OSError, ValueError): - # Unreadable (e.g. mid-replace on Windows): a value no sig carries, - # so a load misses and a build does not publish. - return "" + stamp = "" + if not stamp: + # Unreadable, or empty mid-rewrite (the Windows in-place fallback of + # ``_write_atomically``): a fresh value that equals nothing, so a + # load misses and a build does not publish. + return f"unreadable-{uuid.uuid4().hex}" + return stamp # Legacy (pre-generation) flat file names. A cache written by an older # ``ir`` is still read through these; new writes never use them. diff --git a/tests/test_store_concurrency.py b/tests/test_store_concurrency.py index c026328..a45bbec 100644 --- a/tests/test_store_concurrency.py +++ b/tests/test_store_concurrency.py @@ -343,3 +343,42 @@ def test_delete_record_also_invalidates_a_published_set(tmp_path): _file_store(root).matrix() writer.delete_record("r2") assert sorted(_file_store(root).matrix()[0]) == ["r1"] + + +def test_an_unreadable_write_stamp_never_matches(tmp_path): + """An empty stamp (mid-rewrite on Windows) must neither publish nor load.""" + root = tmp_path / "corpus" + writer = _file_store(root) + writer.put_record(_rec("r1")) + stamp = root / "matrix" / "write-stamp" + stamp.write_bytes(b"") + _file_store(root).matrix() + assert not (root / "matrix" / "sig.json").exists() + + +def test_a_failed_stamp_write_drops_the_stamp(tmp_path, monkeypatch): + """A build that read the old stamp must not publish after a failed stamp.""" + root = tmp_path / "corpus" + writer = _file_store(root) + writer.put_record(_rec("r1")) + reader = _file_store(root) + real_build = reader._build_matrix + + real_write = ir_store._write_atomically + + def refuse(path, data): + if path.name == "write-stamp": + raise OSError("read-only") + real_write(path, data) + + def build_then_writer_continues(): + result = real_build() + monkeypatch.setattr(ir_store, "_write_atomically", refuse) + writer.put_record(_rec("r2", vec=(0.0, 1.0, 0.0))) + monkeypatch.undo() + return result + + reader._build_matrix = build_then_writer_continues + reader.matrix() + assert not (root / "matrix" / "write-stamp").exists() + assert sorted(_file_store(root).matrix()[0]) == ["r1", "r2"]