aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/indexer
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/indexer')
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/cache.py68
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py57
2 files changed, 122 insertions, 3 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/cache.py b/packages/meshbay-node/src/meshbay_node/indexer/cache.py
index c167db9..6c87f40 100644
--- a/packages/meshbay-node/src/meshbay_node/indexer/cache.py
+++ b/packages/meshbay-node/src/meshbay_node/indexer/cache.py
@@ -23,6 +23,7 @@ only daemon.py's wiring changed.
"""
import logging
+import time
from dataclasses import dataclass
from pathlib import Path
@@ -40,6 +41,25 @@ CREATE TABLE IF NOT EXISTS files (
added_at INTEGER NOT NULL,
hash_version INTEGER NOT NULL DEFAULT 1
);
+
+-- Who sent a file. Written when an upload finishes, read when the file is
+-- indexed — two different moments, and the second is much later: the watchdog
+-- debounces for two seconds and then hashes, so the entry does not exist yet
+-- when the last chunk lands. Recording it in memory would also lose it at every
+-- restart, where the index is rebuilt from disk, and an owner the node forgets
+-- is an owner who cannot delete their own file tomorrow.
+--
+-- Validated against a live stat() exactly as `files` is: a row whose size or
+-- mtime no longer match is a different file at that path, and attributes
+-- nothing. That is what makes a row left behind by a deleted file harmless.
+CREATE TABLE IF NOT EXISTS uploads (
+ path TEXT PRIMARY KEY,
+ size INTEGER NOT NULL,
+ mtime REAL NOT NULL,
+ user_id TEXT NOT NULL,
+ pk_ed25519 TEXT NOT NULL DEFAULT '',
+ at INTEGER NOT NULL
+);
"""
_MIGRATE_V2 = "ALTER TABLE files ADD COLUMN hash_version INTEGER NOT NULL DEFAULT 1"
@@ -59,6 +79,11 @@ class IndexCache:
def __init__(self, db_path: Path):
self._db_path = db_path
self._db: aiosqlite.Connection | None = None
+ # Whether `uploads` holds anything at all. Every entry the indexer
+ # builds asks this cache who sent the file, and on a node that has
+ # never received an upload — most of them, most of the time — that is
+ # one query per file per scan for an answer that is always None.
+ self._has_uploads = False
async def open(self) -> None:
self._db_path.parent.mkdir(parents=True, exist_ok=True)
@@ -69,6 +94,10 @@ class IndexCache:
except Exception:
pass # column already exists
await self._db.commit()
+ async with self._db.execute(
+ "SELECT EXISTS(SELECT 1 FROM uploads)") as cur:
+ row = await cur.fetchone()
+ self._has_uploads = bool(row and row[0])
async def close(self) -> None:
if self._db:
@@ -114,6 +143,45 @@ class IndexCache:
(path, mtime, size, hash, type, added_at, hash_version))
await self._db.commit()
+ # ── Who sent a file ──────────────────────────────────────────────────────
+
+ async def record_upload(self, path: str, size: int, mtime: float,
+ user_id: str, pk_ed25519: str) -> None:
+ """Remember that this account put this file here.
+
+ Written once the upload is complete and the file is at its final name,
+ never partway through — a `.part` is not indexable and would key a row
+ to a path that is about to change.
+ """
+ await self._db.execute(
+ "INSERT INTO uploads (path, size, mtime, user_id, pk_ed25519, at) "
+ "VALUES (?, ?, ?, ?, ?, ?) "
+ "ON CONFLICT(path) DO UPDATE SET "
+ "size = excluded.size, mtime = excluded.mtime, "
+ "user_id = excluded.user_id, pk_ed25519 = excluded.pk_ed25519, "
+ "at = excluded.at",
+ (path, size, mtime, user_id, pk_ed25519, int(time.time())))
+ await self._db.commit()
+ self._has_uploads = True
+
+ async def uploader(self, path: str, size: int,
+ mtime: float) -> tuple[str, str] | None:
+ """Who sent the file currently at this path, or None.
+
+ The size and mtime are matched exactly, like `lookup`: the path alone
+ would credit whoever last uploaded *a* file of that name for whatever
+ occupies the name now — including something the operator put there
+ themselves afterwards, which would hand a member the right to delete it.
+ """
+ if not self._has_uploads:
+ return None
+ async with self._db.execute(
+ "SELECT user_id, pk_ed25519 FROM uploads "
+ "WHERE path = ? AND size = ? AND mtime = ?",
+ (path, size, mtime)) as cur:
+ row = await cur.fetchone()
+ return (row[0], row[1]) if row else None
+
# ── Maintenance (node admin UI "prune index cache") ──────────────────────
async def count(self) -> int:
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
index 33e7210..852c061 100644
--- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
+++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
@@ -436,7 +436,7 @@ class DirectoryIndexer:
cached = await self._cache.lookup(
str(file_path), st.st_size, st.st_mtime, expected_hv)
if cached is not None:
- return IndexEntry(
+ return await self._attribute(IndexEntry(
id=cached.hash,
name=file_path.name,
path=_virtual_dir(root, file_path),
@@ -444,7 +444,7 @@ class DirectoryIndexer:
type=cached.type,
added_at=cached.added_at,
hash_version=cached.hash_version,
- )
+ ), file_path, st)
loop = asyncio.get_event_loop()
entry = await loop.run_in_executor(self._executor, _scan_file, root, file_path)
@@ -452,8 +452,51 @@ class DirectoryIndexer:
await self._cache.put(str(file_path), st.st_size, st.st_mtime,
entry.id, entry.type, entry.added_at,
entry.hash_version)
+ return await self._attribute(entry, file_path, st)
+
+ async def _attribute(self, entry: IndexEntry | None, file_path: Path,
+ st) -> IndexEntry | None:
+ """Stamp an entry with whoever sent the file, if a member did.
+
+ Here, rather than beside each `add_entry`, because this is the one
+ funnel every entry passes through: the initial scan, the watchdog,
+ reconcile and a replug all build theirs from `_hash_or_cached`.
+
+ The attribution used to be written at the end of the *upload* instead,
+ by walking the index for an entry that by construction did not exist
+ yet — the watchdog has not fired, and the `.part` the file was until the
+ rename is not indexable. It matched nothing, silently, so every uploaded
+ file was owned by nobody and `file_delete` refused everyone but the
+ operator, where MESHBAY_DESIGN.md §5.4 grants it to any non-revoked
+ device of the uploading account.
+ """
+ if entry is None or self._cache is None:
+ return entry
+ who = await self._cache.uploader(str(file_path), st.st_size, st.st_mtime)
+ if who is not None:
+ entry.uploader_id, entry.uploader_pk = who
return entry
+ async def record_upload(self, file_path: Path, user_id: str,
+ pk_ed25519: str) -> None:
+ """Remember who sent this file, for the entry that does not exist yet.
+
+ Called by the transport once the last chunk has landed and the file is
+ at its final name. Durable rather than in-memory: the index is rebuilt
+ from disk at every start, and an owner the node forgets on restart is an
+ owner who cannot delete their own file tomorrow.
+ """
+ if self._cache is None or not user_id:
+ return
+ try:
+ st = file_path.stat()
+ except OSError:
+ # Gone between the rename and here. Nothing to attribute, and
+ # nothing for anyone to delete either.
+ return
+ await self._cache.record_upload(
+ str(file_path), st.st_size, st.st_mtime, user_id, pk_ed25519 or "")
+
def _report_collisions(self) -> None:
"""
Names that are the same file on a case-insensitive filesystem.
@@ -749,7 +792,15 @@ class DirectoryIndexer:
if old is None:
continue
for field in self._ENRICHED_FIELDS:
- setattr(entry, field, getattr(old, field))
+ # What the rescan itself established wins. `uploader_id` and
+ # `uploader_pk` now come off the durable record (`_attribute`),
+ # and the entry being replaced is memory this process happens
+ # to still hold — so copying over them would let a stale blank
+ # beat the thing that survives a restart. Every other field is
+ # None on a freshly scanned entry, so for those this is exactly
+ # the carry-over it has always been.
+ if getattr(entry, field) is None:
+ setattr(entry, field, getattr(old, field))
# It came back intact, so it is not one of the entries the daemon
# needs to enrich again.
self.rescanned_ids.discard(entry.id)