aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-06 01:33:16 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-06 01:33:37 +0200
commitf2da33a648f86e34ddbbb5f6bec124828ef2a847 (patch)
tree670cbf62ff92598847e067b5835713dfd854dfb5 /packages/meshbay-node/src
parentfff1974edf19cf1186e0f49da5f8a4d237bcb13e (diff)
downloadmeshbay-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/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/cache.py56
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py67
2 files changed, 83 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: