diff --git a/ir/store.py b/ir/store.py index 4ff696e..44e9682 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.""" @@ -291,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 ------------------------------------------------------ # @@ -484,7 +497,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 +514,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 +526,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 +615,71 @@ 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, 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: + 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 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: + """The current write stamp, or ``None`` if no stamped write happened yet.""" + if self._packed_dir is None: + return None + try: + stamp = self._write_stamp_path().read_text(encoding="ascii") + except FileNotFoundError: + return None + except (OSError, ValueError): + 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. _LEGACY_PACKED_FILES = { @@ -722,11 +795,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 +836,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 +855,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..a45bbec 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,100 @@ 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"] + + +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"] 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)