diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-06 01:33:16 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-06 01:33:37 +0200 |
| commit | f2da33a648f86e34ddbbb5f6bec124828ef2a847 (patch) | |
| tree | 670cbf62ff92598847e067b5835713dfd854dfb5 /packages/meshbay-node | |
| parent | fff1974edf19cf1186e0f49da5f8a4d237bcb13e (diff) | |
| download | meshbay-f2da33a648f86e34ddbbb5f6bec124828ef2a847.tar.gz | |
feat(node): indexing v2 — partial-read hashing for files above 40 MB
Files above 40 MB are no longer read in full. Instead, blake3 hashes
45 MB of samples (first 20 MB + last 20 MB + 5 MB at 50% offset).
Files at or below 40 MB are unchanged (full read, hash_version 1).
A new `hash_version` field on IndexEntry (default 1) travels on the
wire and through the cache so both versions coexist without breaking
existing nodes or clients.
The IndexCache auto-migrates its schema on open (ALTER TABLE), so no
manual step is required on upgrade. A standalone migration script is
available in QE/migration/ for operators who want to preview or force
a full re-hash.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/cache.py | 56 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/indexer.py | 67 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_index_cache.py | 77 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_indexer.py | 145 |
4 files changed, 305 insertions, 40 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/cache.py b/packages/meshbay-node/src/meshbay_node/indexer/cache.py index 31b9b15..c167db9 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/cache.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/cache.py @@ -32,21 +32,25 @@ log = logging.getLogger(__name__) _SCHEMA = """ CREATE TABLE IF NOT EXISTS files ( - path TEXT PRIMARY KEY, - mtime REAL NOT NULL, - size INTEGER NOT NULL, - hash TEXT NOT NULL, - type TEXT NOT NULL, - added_at INTEGER NOT NULL + path TEXT PRIMARY KEY, + mtime REAL NOT NULL, + size INTEGER NOT NULL, + hash TEXT NOT NULL, + type TEXT NOT NULL, + added_at INTEGER NOT NULL, + hash_version INTEGER NOT NULL DEFAULT 1 ); """ +_MIGRATE_V2 = "ALTER TABLE files ADD COLUMN hash_version INTEGER NOT NULL DEFAULT 1" + @dataclass class CachedEntry: hash: str type: str added_at: int + hash_version: int = 1 class IndexCache: @@ -60,6 +64,10 @@ class IndexCache: self._db_path.parent.mkdir(parents=True, exist_ok=True) self._db = await aiosqlite.connect(str(self._db_path)) await self._db.executescript(_SCHEMA) + try: + await self._db.execute(_MIGRATE_V2) + except Exception: + pass # column already exists await self._db.commit() async def close(self) -> None: @@ -74,34 +82,36 @@ class IndexCache: async def __aexit__(self, *_): await self.close() - async def lookup(self, path: str, size: int, mtime: float) -> CachedEntry | None: + async def lookup(self, path: str, size: int, mtime: float, + hash_version: int = 1) -> CachedEntry | None: """ - A cache hit requires an EXACT match on both size and mtime. A mtime - touched without a content change is a false negative (an unnecessary - rehash) — accepted, since the alternative (trusting a stale hash) is - a silent wrong answer instead of an occasional wasted read. + 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. """ async with self._db.execute( - "SELECT hash, type, added_at FROM files " - "WHERE path = ? AND size = ? AND mtime = ?", - (path, size, mtime)) as cur: + "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: row = await cur.fetchone() - return CachedEntry(hash=row[0], type=row[1], added_at=row[2]) if row else None + return CachedEntry(hash=row[0], type=row[1], added_at=row[2], + hash_version=row[3]) if row else None async def put(self, path: str, size: int, mtime: float, hash: str, - type: str, added_at: int) -> None: + type: str, added_at: int, hash_version: int = 1) -> None: """ - Written only once a file has been hashed in full — never partway - through — so a crash mid-hash leaves no stale/partial row behind: the - next scan simply finds no cache entry and hashes the file again. + Written only once a file has been hashed — never partway through — so + a crash mid-hash leaves no stale/partial row behind: the next scan + simply finds no cache entry and hashes the file again. """ await self._db.execute( - "INSERT INTO files (path, mtime, size, hash, type, added_at) " - "VALUES (?, ?, ?, ?, ?, ?) " + "INSERT INTO files (path, mtime, size, hash, type, added_at, hash_version) " + "VALUES (?, ?, ?, ?, ?, ?, ?) " "ON CONFLICT(path) DO UPDATE SET " "mtime = excluded.mtime, size = excluded.size, hash = excluded.hash, " - "type = excluded.type, added_at = excluded.added_at", - (path, mtime, size, hash, type, added_at)) + "type = excluded.type, added_at = excluded.added_at, " + "hash_version = excluded.hash_version", + (path, mtime, size, hash, type, added_at, hash_version)) await self._db.commit() # ── Maintenance (node admin UI "prune index cache") ────────────────────── diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py index f78465b..8376c23 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py @@ -94,6 +94,32 @@ 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 + + +def _feed(hasher, f, nbytes: int) -> None: + remaining = nbytes + while remaining > 0: + chunk = f.read(min(_HASH_CHUNK, remaining)) + if not chunk: + break + hasher.update(chunk) + remaining -= len(chunk) + + +def _partial_hash(file_path: Path, size: int) -> str: + hasher = blake3.blake3() + with open(long_path(file_path), "rb") as f: + _feed(hasher, f, _PARTIAL_HEAD) + f.seek(size - _PARTIAL_TAIL) + _feed(hasher, f, _PARTIAL_TAIL) + f.seek(size // 2) + _feed(hasher, f, _PARTIAL_MID) + return hasher.hexdigest() + @dataclass class IndexProgress: @@ -162,29 +188,33 @@ 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. - Uses streaming blake3 so arbitrarily large files (ISOs, VM images, etc.) - don't require loading the whole file into memory.""" + 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.""" if not _is_indexable(file_path): return None try: stat = file_path.stat() if not _is_indexable_size(file_path, stat.st_size): return None - hasher = blake3.blake3() - # long_path is a no-op off Windows; there it is what lets a deep media - # library past MAX_PATH. - with open(long_path(file_path), "rb") as f: - while chunk := f.read(_HASH_CHUNK): - hasher.update(chunk) + if stat.st_size > _PARTIAL_THRESHOLD: + hex_hash = _partial_hash(file_path, stat.st_size) + hv = 2 + else: + hasher = blake3.blake3() + with open(long_path(file_path), "rb") as f: + while chunk := f.read(_HASH_CHUNK): + hasher.update(chunk) + hex_hash = hasher.hexdigest() + hv = 1 return IndexEntry( - id=hasher.hexdigest(), - # Stored exactly as the filesystem gave it: this is the string that - # opens the file. Normalization is for comparison only. + id=hex_hash, name=file_path.name, path=_virtual_dir(root, file_path), size=stat.st_size, type=_detect_type(file_path), added_at=int(stat.st_mtime), + hash_version=hv, ) except (OSError, PermissionError, ValueError) as e: log.warning("Cannot index %s: %s", file_path, e) @@ -375,10 +405,8 @@ 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) still match - what was hashed last time — the difference between a redundant full - rehash of a 100+ GB library on every restart and a stat()-only pass. - The only place that decides to actually read a file's bytes. + content read entirely when this path's (size, mtime, hash_version) + still match what was hashed last time. """ if not _is_indexable(file_path): return None @@ -389,8 +417,11 @@ 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) + cached = await self._cache.lookup( + str(file_path), st.st_size, st.st_mtime, expected_hv) if cached is not None: return IndexEntry( id=cached.hash, @@ -399,13 +430,15 @@ class DirectoryIndexer: size=st.st_size, type=cached.type, added_at=cached.added_at, + hash_version=cached.hash_version, ) loop = asyncio.get_event_loop() entry = await loop.run_in_executor(self._executor, _scan_file, root, file_path) if entry and self._cache is not None: await self._cache.put(str(file_path), st.st_size, st.st_mtime, - entry.id, entry.type, entry.added_at) + entry.id, entry.type, entry.added_at, + entry.hash_version) return entry def _report_collisions(self) -> None: diff --git a/packages/meshbay-node/tests/test_index_cache.py b/packages/meshbay-node/tests/test_index_cache.py index 24e2f76..036b4d7 100644 --- a/packages/meshbay-node/tests/test_index_cache.py +++ b/packages/meshbay-node/tests/test_index_cache.py @@ -122,3 +122,80 @@ async def test_cache_survives_reopen(tmp_path): assert hit is not None assert hit.hash == "abc123" + + +# ── hash_version support (indexing v2) ───────────────────────────────────────── + + +@pytest.mark.asyncio +async def test_put_v2_then_lookup_hits(cache): + await cache.put("/lib/big.mkv", size=50_000_000, mtime=111.0, + hash="partial_abc", type="video", added_at=42, + hash_version=2) + + hit = await cache.lookup("/lib/big.mkv", size=50_000_000, mtime=111.0, + hash_version=2) + assert hit is not None + assert hit.hash == "partial_abc" + assert hit.hash_version == 2 + + +@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): + 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) + + 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) + assert hit is not None + assert hit.hash == "partial_hash" + + +@pytest.mark.asyncio +async def test_v1_schema_auto_migrates(tmp_path): + """An index_cache.db created by old code (no hash_version column) gains + the column on the next open(), and existing rows default to hash_version=1.""" + import aiosqlite + db_path = tmp_path / "old_cache.db" + async with aiosqlite.connect(str(db_path)) as db: + await db.executescript(""" + CREATE TABLE files ( + path TEXT PRIMARY KEY, + mtime REAL NOT NULL, + size INTEGER NOT NULL, + hash TEXT NOT NULL, + type TEXT NOT NULL, + added_at INTEGER NOT NULL + ); + """) + await db.execute( + "INSERT INTO files (path, mtime, size, hash, type, added_at) " + "VALUES (?, ?, ?, ?, ?, ?)", + ("/lib/old.mkv", 111.0, 1000, "oldhash", "video", 42)) + await db.commit() + + cache = IndexCache(db_path=db_path) + await cache.open() + hit = await cache.lookup("/lib/old.mkv", size=1000, mtime=111.0, + hash_version=1) + await cache.close() + + assert hit is not None + assert hit.hash == "oldhash" + assert hit.hash_version == 1 diff --git a/packages/meshbay-node/tests/test_indexer.py b/packages/meshbay-node/tests/test_indexer.py index 729dade..6aee1b5 100644 --- a/packages/meshbay-node/tests/test_indexer.py +++ b/packages/meshbay-node/tests/test_indexer.py @@ -723,3 +723,148 @@ async def test_reconcile_backoff_resets_when_something_actually_changes( await task except asyncio.CancelledError: pass + + +# ── Indexing v2 — partial-read hashing ──────────────────────────────────────── + + +def test_small_file_gets_hash_version_1(tmp_path): + from meshbay_node.indexer.indexer import _scan_file + d = tmp_path / "root" + d.mkdir() + f = d / "small.mp4" + f.write_bytes(os.urandom(1024)) + + entry = _scan_file(one_root(d).roots[0], f) + assert entry is not None + assert entry.hash_version == 1 + + +def test_large_file_gets_hash_version_2(tmp_path): + from meshbay_node.indexer.indexer import _scan_file, _PARTIAL_THRESHOLD + d = tmp_path / "root" + d.mkdir() + f = d / "big.mkv" + size = _PARTIAL_THRESHOLD + 1 + f.write_bytes(os.urandom(size)) + + entry = _scan_file(one_root(d).roots[0], f) + assert entry is not None + assert entry.hash_version == 2 + assert entry.size == size + + +def test_file_at_threshold_gets_hash_version_1(tmp_path): + from meshbay_node.indexer.indexer import _scan_file, _PARTIAL_THRESHOLD + d = tmp_path / "root" + d.mkdir() + f = d / "exact.mkv" + f.write_bytes(os.urandom(_PARTIAL_THRESHOLD)) + + entry = _scan_file(one_root(d).roots[0], f) + assert entry is not None + assert entry.hash_version == 1 + + +def test_partial_hash_differs_from_full_hash(tmp_path): + """For a file above the threshold, the partial hash must differ from what + a full-file blake3 would produce (they read different bytes).""" + import blake3 as b3 + from meshbay_node.indexer.indexer import _scan_file, _PARTIAL_THRESHOLD + d = tmp_path / "root" + d.mkdir() + f = d / "big.mkv" + content = os.urandom(_PARTIAL_THRESHOLD + 1024 * 1024) + f.write_bytes(content) + + entry = _scan_file(one_root(d).roots[0], f) + full_hash = b3.blake3(content).hexdigest() + + assert entry.id != full_hash + assert entry.hash_version == 2 + + +def test_partial_hash_is_deterministic(tmp_path): + from meshbay_node.indexer.indexer import _scan_file, _PARTIAL_THRESHOLD + d = tmp_path / "root" + d.mkdir() + f = d / "big.mkv" + f.write_bytes(os.urandom(_PARTIAL_THRESHOLD + 1)) + + e1 = _scan_file(one_root(d).roots[0], f) + e2 = _scan_file(one_root(d).roots[0], f) + assert e1.id == e2.id + + +def test_group_index_roundtrip_preserves_hash_version(sk_node, gek): + from meshbay_common.protocol import IndexEntry + idx = GroupIndex(group_id="hv-test", sk_node=sk_node, gek=gek) + idx.add_entry(IndexEntry( + id="aaa", name="small.mp4", path="root", size=1024, + type="video", added_at=100, hash_version=1)) + idx.add_entry(IndexEntry( + id="bbb", name="big.mkv", path="root", size=50_000_000, + type="video", added_at=200, hash_version=2)) + + wire = idx.serialize() + recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) + + by_id = {e.id: e for e in recovered.entries} + assert by_id["aaa"].hash_version == 1 + assert by_id["bbb"].hash_version == 2 + + +def test_deserialize_without_hash_version_defaults_to_1(sk_node, gek): + """Entries serialized by old code (no hash_version field) must deserialize + as hash_version=1.""" + from meshbay_common.protocol import IndexEntry + idx = GroupIndex(group_id="compat", sk_node=sk_node, gek=gek) + idx.add_entry(IndexEntry( + id="old", name="f.mp4", path="root", size=1024, + type="video", added_at=100)) + wire = idx.serialize() + recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) + assert recovered.entries[0].hash_version == 1 + + +def test_index_entry_wire_includes_hash_version(): + from meshbay_common.protocol import IndexEntry, index_entry_wire + e = IndexEntry(id="x", name="f.mp4", path="root", size=1, + type="video", added_at=0, hash_version=2) + w = index_entry_wire(e) + assert w["hash_version"] == 2 + + +@pytest.mark.asyncio +async def test_cache_aware_scan_uses_hash_version(tmp_path, sk_node, gek): + from meshbay_node.indexer.indexer import _PARTIAL_THRESHOLD + d = tmp_path / "root" + d.mkdir() + small = d / "small.mp4" + small.write_bytes(os.urandom(1024)) + big = d / "big.mkv" + big.write_bytes(os.urandom(_PARTIAL_THRESHOLD + 1)) + + cache = IndexCache(db_path=tmp_path / "cache.db") + await cache.open() + + indexer = DirectoryIndexer( + roots=one_root(d), group_id="g", sk_node=sk_node, gek=gek, + cache=cache) + await indexer.initial_scan() + + 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 + + hit_small = await cache.lookup( + str(small), small.stat().st_size, small.stat().st_mtime, + hash_version=1) + assert hit_small is not None + + hit_big = await cache.lookup( + str(big), big.stat().st_size, big.stat().st_mtime, + hash_version=2) + assert hit_big is not None + + await cache.close() |