aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-13 15:40:34 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-13 15:40:34 +0200
commitd917bb61e42336c38782b22da604d7ca923d484a (patch)
tree3ad99072f02d621ca4a54f4fcce65a2c530c4b79 /packages/meshbay-node/src
parentc5fff4ce8366b08669c0c8b6d30b99b94b9fefca (diff)
downloadmeshbay-d917bb61e42336c38782b22da604d7ca923d484a.tar.gz
fix(node): an uploaded file records who sent it
`_register_uploader` walked the index for the entry it had just written, at a moment when no such entry can exist: the file was a `.part` until the rename on the line above, which is not indexable, and the watchdog that will index it debounces for two seconds and then hashes. The walk matched nothing, silently, so every uploaded file in every group was owned by nobody — and `file_delete` refuses a caller with no admin authority when the entry records no uploader, so a member could not delete what they had just sent. MESHBAY_DESIGN.md §5.4 grants that to any non-revoked device of the uploading account. The record is now written when the last chunk lands (`indexer.record_upload`) and the entry is stamped from it in `_hash_or_cached`, the one funnel every entry passes through — initial scan, watchdog, reconcile and replug alike. It lives in the index cache rather than on the entry alone, because the index is rebuilt from disk at every start and an owner the node forgets on restart is a right quietly taken away. It is validated against a live `stat()`, so whatever later occupies that path inherits nothing; and `_rescan_root`'s carry-over no longer copies over it, or memory would beat the durable record. §5.4 also claimed ownership was *provable* — a transcript the uploader signs, stored with the entry. No such signature has ever existed; `meshbay:upload:v1` in the code is the groupbox purpose that seals the envelope. The section now states what the code does, and the transcript is an open item in §15.3. `test_upload_attribution.py` drives the real handler and a real indexer across that seam. Against the previous source its two positive cases fail on the property, not on a missing method — an upload, then a rebuild from disk, then a different file at the same path inheriting nothing. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UMxEQadpzPkYLFf5CYKhpW
Diffstat (limited to 'packages/meshbay-node/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py6
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/cache.py68
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py57
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py25
4 files changed, 142 insertions, 14 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py
index dcfa8f8..9428ddd 100644
--- a/packages/meshbay-node/src/meshbay_node/daemon.py
+++ b/packages/meshbay-node/src/meshbay_node/daemon.py
@@ -393,6 +393,11 @@ class NodeDaemon:
# _reconcile_loop) so the backstop is prompt again now
# that someone is actually looking.
"note_activity": indexer.note_activity,
+ # Bound method, called when an upload finishes. The entry
+ # it belongs to does not exist yet (see
+ # webrtc_server._register_uploader), so the indexer keeps
+ # the record and stamps the entry when it creates it.
+ "record_upload": indexer.record_upload,
# Shown to the operator in Settings, and kept current in
# place by set_scan_settings (ops.py) — same reasoning as
# enabled_apps below.
@@ -879,6 +884,7 @@ class NodeDaemon:
"index": indexer.index,
"progress": indexer.progress,
"note_activity": indexer.note_activity,
+ "record_upload": indexer.record_upload,
"reconcile_interval_secs": scan_settings["reconcile_interval_secs"],
"debounce_secs": scan_settings["debounce_secs"],
"visibility": group_cfg.visibility,
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)
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
index 29d6e7f..8df88d8 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -4914,26 +4914,29 @@ class WebRTCPeerSession:
log.info("Upload complete: %s (%d chunks, %d bytes)",
stored_name, total_chunks, state.bytes)
self._audit("file_upload", f"{rel_dir}/{stored_name}")
- self._register_uploader(ctx, rel_dir, stored_name)
+ self._register_uploader(ctx, final_path)
- def _register_uploader(self, ctx: dict, rel_dir: str, filename: str) -> None:
+ def _register_uploader(self, ctx: dict, file_path: Path) -> None:
"""
- Tag the index entry with the uploader's identity after upload completes.
+ Record who sent this file, for the index entry that does not exist yet.
- The key recorded here is the one this node pinned, not the one the token
+ The key recorded is the one this node pinned, not the one the token
carried. `pk_user` was a hub-chosen claim, and it decided who could later
delete the file: a hub issuing a token naming its own key could delete
anyone's uploads on any node. Deletion is supposed to be authorized by the
node, and this closes the last place where it was not.
+
+ **The entry is not here to be tagged.** This used to walk `ctx["index"]`
+ for the name just written and set the fields on it; at this point the
+ watchdog has not fired (it debounces for two seconds and then hashes)
+ and the file was a `.part` until the line above, which is not indexable
+ — so the walk matched nothing, every time, and said nothing about it.
+ The indexer stamps the entry from this record when it creates it.
"""
- idx = ctx.get("index")
- if not idx:
+ record = ctx.get("record_upload")
+ if record is None:
return
- for entry in idx.entries:
- if entry.name == filename and entry.path == rel_dir:
- entry.uploader_id = self._user_id
- entry.uploader_pk = self._pinned_pk
- return
+ self._spawn(record(file_path, self._user_id or "", self._pinned_pk or ""))
def _do_file_delete(self, msg: dict) -> None:
ctx = self._group_ctx()