Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 113 additions & 9 deletions ir/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down Expand Up @@ -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 ------------------------------------------------------ #

Expand Down Expand Up @@ -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
Expand All @@ -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:
Expand All @@ -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

Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand All @@ -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
Expand Down
107 changes: 101 additions & 6 deletions tests/test_store_concurrency.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
Expand All @@ -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"]
11 changes: 7 additions & 4 deletions tests/test_store_packed.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -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):
Expand Down Expand Up @@ -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)
Expand Down
Loading