aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-common/src/meshbay_common/protocol.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/cache.py13
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py33
-rw-r--r--packages/meshbay-node/tests/test_index_cache.py39
-rw-r--r--packages/meshbay-node/tests/test_indexer.py51
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()