diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-10-08 01:27:58 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-10-08 01:27:58 +0200 |
| commit | c27e04c88557716bbe8e42b9174ba9b07f90facf (patch) | |
| tree | a1a56edbede2bba9bb0cc99b81b98570b8313915 /packages | |
| parent | cdd5fd52e981c4c59643e7dee705b6c55acae68e (diff) | |
| download | meshbay-c27e04c88557716bbe8e42b9174ba9b07f90facf.tar.gz | |
perf(node): sample 9 MB with the size above 9 MB, keep known ids0.19
hash_version 3: size + first 4 MB + last 4 MB + 1 MB at the middle,
5.5x faster cold on a USB disk than the 45 MB sample. The cache now
serves a hit under whatever version it holds, so no existing id moves.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages')
5 files changed, 84 insertions, 54 deletions
diff --git a/packages/meshbay-common/src/meshbay_common/protocol.py b/packages/meshbay-common/src/meshbay_common/protocol.py index f8e56ec..c048b8d 100644 --- a/packages/meshbay-common/src/meshbay_common/protocol.py +++ b/packages/meshbay-common/src/meshbay_common/protocol.py @@ -281,7 +281,7 @@ class IndexEntry: track_no: int | None = None # tag or parsed, Music app taken_at: int | None = None # unix timestamp, EXIF DateTimeOriginal — Photos app camera: str | None = None # "Make Model", when both present — Photos app - hash_version: int = 1 # 1 = full-file blake3, 2 = partial-read (45 MB sample) + hash_version: int = 1 # 1 = full-file blake3, 2 = 45 MB sample, 3 = size + 9 MB sample def index_entry_wire(e: IndexEntry) -> dict: diff --git a/packages/meshbay-node/src/meshbay_node/indexer/cache.py b/packages/meshbay-node/src/meshbay_node/indexer/cache.py index 5e328ef..afa4f85 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/cache.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/cache.py @@ -112,17 +112,16 @@ class IndexCache: async def __aexit__(self, *_): await self.close() - async def lookup(self, path: str, size: int, mtime: float, - hash_version: int = 1) -> CachedEntry | None: + async def lookup(self, path: str, size: int, mtime: float) -> CachedEntry | None: """ - A cache hit requires an EXACT match on size, mtime AND hash_version. - A v1 cached hash won't serve a v2 lookup for the same path — the file - is re-hashed with the new algorithm instead. + A cache hit requires an EXACT match on size and mtime, and serves the + hash under whatever hash_version it was computed: a scheme change must + not change the id of a file the node already knows. """ async with self._db.execute( "SELECT hash, type, added_at, hash_version FROM files " - "WHERE path = ? AND size = ? AND mtime = ? AND hash_version = ?", - (path, size, mtime, hash_version)) as cur: + "WHERE path = ? AND size = ? AND mtime = ?", + (path, size, mtime)) as cur: row = await cur.fetchone() return CachedEntry(hash=row[0], type=row[1], added_at=row[2], hash_version=row[3]) if row else None diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py index f06d30f..4eeeec0 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py @@ -97,10 +97,14 @@ def _is_indexable_size(path: Path, size: int) -> bool: _HASH_CHUNK = 8 * 1024 * 1024 # 8 MB streaming hash chunks -_PARTIAL_THRESHOLD = 40 * 1024 * 1024 # files above this use partial-read hashing -_PARTIAL_HEAD = 20 * 1024 * 1024 -_PARTIAL_TAIL = 20 * 1024 * 1024 -_PARTIAL_MID = 5 * 1024 * 1024 +# On a spinning disk the cost of a sample is its three seeks, not its bytes: +# measured cold on USB, 45 MB took 516 ms a file, 9 MB 94 ms, 3 MB 79 ms. The +# threshold is the sample size, so no file costs more to read than a sample. +_PARTIAL_THRESHOLD = 9 * 1024 * 1024 # files above this use partial-read hashing +_PARTIAL_HEAD = 4 * 1024 * 1024 +_PARTIAL_TAIL = 4 * 1024 * 1024 +_PARTIAL_MID = 1 * 1024 * 1024 +_PARTIAL_VERSION = 3 def _feed(hasher, f, nbytes: int) -> None: @@ -115,6 +119,7 @@ def _feed(hasher, f, nbytes: int) -> None: def _partial_hash(file_path: Path, size: int) -> str: hasher = blake3.blake3() + hasher.update(size.to_bytes(8, "little")) with open(long_path(file_path), "rb") as f: _feed(hasher, f, _PARTIAL_HEAD) f.seek(size - _PARTIAL_TAIL) @@ -204,9 +209,9 @@ def _size_files(files: list[Path]) -> list[tuple[Path, int]]: def _scan_file(root: Root, file_path: Path) -> IndexEntry | None: """Compute IndexEntry for a file. Blocking — run in executor. - Files <= 40 MB are hashed in full (hash_version 1). Files > 40 MB use a - 45 MB partial read — first 20 MB, last 20 MB, 5 MB at 50% — for - hash_version 2.""" + Files <= 9 MB are hashed in full (hash_version 1). Larger files hash their + size, first 4 MB, last 4 MB and 1 MB at 50% (hash_version 3). Version 2, + the earlier 45 MB sample, is still served from the cache, never computed.""" if not _is_indexable(file_path): return None try: @@ -215,7 +220,7 @@ def _scan_file(root: Root, file_path: Path) -> IndexEntry | None: return None if stat.st_size > _PARTIAL_THRESHOLD: hex_hash = _partial_hash(file_path, stat.st_size) - hv = 2 + hv = _PARTIAL_VERSION else: hasher = blake3.blake3() with open(long_path(file_path), "rb") as f: @@ -518,8 +523,12 @@ class DirectoryIndexer: async def _hash_or_cached(self, root: Root, file_path: Path) -> IndexEntry | None: """ Cache-aware replacement for a bare _scan_file() call: skips the - content read entirely when this path's (size, mtime, hash_version) - still match what was hashed last time. + content read entirely when this path's (size, mtime) still match what + was hashed last time. + + Whatever hash_version the cache holds is kept: TMDB matches, manual + corrections, thumbnails and members' playlists are keyed by the id, so + a file the node already knows never changes id because the scheme did. """ if not _is_indexable(file_path): return None @@ -530,11 +539,9 @@ class DirectoryIndexer: if not _is_indexable_size(file_path, st.st_size): return None - expected_hv = 2 if st.st_size > _PARTIAL_THRESHOLD else 1 - if self._cache is not None: cached = await self._cache.lookup( - str(file_path), st.st_size, st.st_mtime, expected_hv) + str(file_path), st.st_size, st.st_mtime) if cached is not None: return await self._attribute(IndexEntry( id=cached.hash, diff --git a/packages/meshbay-node/tests/test_index_cache.py b/packages/meshbay-node/tests/test_index_cache.py index 8f0ff25..e4424d2 100644 --- a/packages/meshbay-node/tests/test_index_cache.py +++ b/packages/meshbay-node/tests/test_index_cache.py @@ -127,43 +127,33 @@ async def test_cache_survives_reopen(tmp_path): @pytest.mark.asyncio -async def test_put_v2_then_lookup_hits(cache): +@pytest.mark.parametrize("hv", [1, 2, 3]) +async def test_a_hit_keeps_the_version_it_was_hashed_under(cache, hv): + """An older scheme is served as is: the id is what TMDB matches, + thumbnails and playlists hang on.""" await cache.put("/lib/big.mkv", size=50_000_000, mtime=111.0, - hash="partial_abc", type="video", added_at=42, - hash_version=2) + hash="old_id", type="video", added_at=42, + hash_version=hv) - hit = await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0, - hash_version=2) + hit = await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0) assert hit is not None - assert hit.hash == "partial_abc" - assert hit.hash_version == 2 + assert hit.hash == "old_id" + assert hit.hash_version == hv @pytest.mark.asyncio -async def test_lookup_misses_on_wrong_hash_version(cache): - await cache.put("/lib/big.mkv", size=50_000_000, mtime=111.0, - hash="full_hash", type="video", added_at=42, - hash_version=1) - - assert await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0, - hash_version=2) is None - - -@pytest.mark.asyncio -async def test_put_v2_overwrites_v1_for_same_path(cache): +async def test_a_newer_hash_overwrites_the_older_for_same_path(cache): await cache.put("/lib/big.mkv", size=50_000_000, mtime=111.0, hash="full_hash", type="video", added_at=42, hash_version=1) await cache.put("/lib/big.mkv", size=50_000_000, mtime=111.0, hash="partial_hash", type="video", added_at=42, - hash_version=2) + hash_version=3) - assert await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0, - hash_version=1) is None - hit = await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0, - hash_version=2) + hit = await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0) assert hit is not None assert hit.hash == "partial_hash" + assert hit.hash_version == 3 @pytest.mark.asyncio @@ -191,8 +181,7 @@ async def test_v1_schema_auto_migrates(tmp_path): cache = IndexCache(db_path=db_path) await cache.open() - hit = await cache.lookup("/lib/old.mkv", size=1000, mtime=111.0, - hash_version=1) + hit = await cache.lookup("/lib/old.mkv", size=1000, mtime=111.0) await cache.close() assert hit is not None diff --git a/packages/meshbay-node/tests/test_indexer.py b/packages/meshbay-node/tests/test_indexer.py index c60600f..fda40c8 100644 --- a/packages/meshbay-node/tests/test_indexer.py +++ b/packages/meshbay-node/tests/test_indexer.py @@ -678,7 +678,7 @@ def test_small_file_gets_hash_version_1(tmp_path): assert entry.hash_version == 1 -def test_large_file_gets_hash_version_2(tmp_path): +def test_large_file_gets_hash_version_3(tmp_path): from meshbay_node.indexer.indexer import _PARTIAL_THRESHOLD, _scan_file d = tmp_path / "root" d.mkdir() @@ -688,7 +688,7 @@ def test_large_file_gets_hash_version_2(tmp_path): entry = _scan_file(one_root(d).roots[0], f) assert entry is not None - assert entry.hash_version == 2 + assert entry.hash_version == 3 assert entry.size == size @@ -719,7 +719,7 @@ def test_partial_hash_differs_from_full_hash(tmp_path): full_hash = b3.blake3(content).hexdigest() assert entry.id != full_hash - assert entry.hash_version == 2 + assert entry.hash_version == 3 def test_partial_hash_is_deterministic(tmp_path): @@ -771,16 +771,51 @@ async def test_cache_aware_scan_uses_hash_version(tmp_path, sk_node, gek): by_name = {e.name: e for e in indexer.index.entries} assert by_name["small.mp4"].hash_version == 1 - assert by_name["big.mkv"].hash_version == 2 + assert by_name["big.mkv"].hash_version == 3 hit_small = await cache.lookup( - str(small), small.stat().st_size, small.stat().st_mtime, - hash_version=1) + str(small), small.stat().st_size, small.stat().st_mtime) assert hit_small is not None hit_big = await cache.lookup( - str(big), big.stat().st_size, big.stat().st_mtime, - hash_version=2) + str(big), big.stat().st_size, big.stat().st_mtime) assert hit_big is not None await cache.close() + + +def test_the_size_is_part_of_a_sampled_hash(tmp_path): + """Two preallocated downloads, all zeros, sample the same bytes.""" + from meshbay_node.indexer.indexer import _PARTIAL_THRESHOLD, _scan_file + d = tmp_path / "root" + d.mkdir() + for name, extra in (("a.mkv", 1), ("b.mkv", 2)): + with open(d / name, "wb") as f: + f.truncate(_PARTIAL_THRESHOLD + extra) + + root = one_root(d).roots[0] + assert _scan_file(root, d / "a.mkv").id != _scan_file(root, d / "b.mkv").id + + +@pytest.mark.asyncio +async def test_a_file_hashed_under_an_older_scheme_keeps_its_id(tmp_path, sk_node, gek): + """The id is what TMDB matches, thumbnails and playlists are keyed by.""" + from meshbay_node.indexer.indexer import _PARTIAL_THRESHOLD + d = tmp_path / "root" + d.mkdir() + big = d / "big.mkv" + big.write_bytes(os.urandom(_PARTIAL_THRESHOLD + 1)) + st = big.stat() + + cache = IndexCache(db_path=tmp_path / "cache.db") + await cache.open() + await cache.put(str(big), st.st_size, st.st_mtime, "v2_id", "video", 0, + hash_version=2) + try: + indexer = DirectoryIndexer(roots=one_root(d), group_id="g", + sk_node=sk_node, gek=gek, cache=cache) + await indexer.initial_scan() + [entry] = indexer.index.entries + assert (entry.id, entry.hash_version) == ("v2_id", 2) + finally: + await cache.close() |