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/group_index.py8
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py51
2 files changed, 40 insertions, 19 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
index 4edeeab..a69429c 100644
--- a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
+++ b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py
@@ -25,14 +25,16 @@ import zstandard as zstd
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
from meshbay_common.crypto import (
- chunk_key as derive_chunk_key,
- encrypt_chunk,
- decrypt_chunk,
sign_chunk,
verify_chunk_signature,
pk_to_b64,
generate_gek,
)
+from meshbay_common.webcrypto import (
+ chunk_key_aes as derive_chunk_key,
+ encrypt_chunk_aes as encrypt_chunk,
+ decrypt_chunk_aes as decrypt_chunk,
+)
from meshbay_common.protocol import IndexEntry, IndexDelta
log = logging.getLogger(__name__)
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
index 4bb1527..60dc04b 100644
--- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
+++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py
@@ -169,30 +169,49 @@ class DirectoryIndexer:
# ── Internal update ───────────────────────────────────────────────────────
+ _DEBOUNCE_SECS = 2.0
+
def _schedule_update(self, file_path: Path, deleted: bool = False) -> None:
- """Called from watchdog thread — schedule async update on the event loop."""
- if self._loop:
- self._loop.call_soon_threadsafe(
- lambda: asyncio.ensure_future(
- self._update_entry(file_path, deleted)))
+ """Called from watchdog thread — schedule debounced async update."""
+ if not self._loop:
+ return
+ key = str(file_path.resolve())
+ self._loop.call_soon_threadsafe(
+ self._debounce, key, file_path, deleted)
+
+ def _debounce(self, key: str, file_path: Path, deleted: bool) -> None:
+ if not hasattr(self, "_pending_timers"):
+ self._pending_timers: dict[str, asyncio.TimerHandle] = {}
+ old = self._pending_timers.pop(key, None)
+ if old:
+ old.cancel()
+ handle = self._loop.call_later(
+ self._DEBOUNCE_SECS,
+ lambda: asyncio.ensure_future(self._update_entry(file_path, deleted)),
+ )
+ self._pending_timers[key] = handle
+
+ def _remove_by_path(self, file_path: Path) -> None:
+ """Remove any existing entries that match this file's path + name."""
+ resolved = file_path.resolve()
+ to_remove = [
+ e.id for e in self._index.entries
+ if (self.root / e.path / e.name).resolve() == resolved
+ ]
+ for fid in to_remove:
+ self._index.remove_entry(fid)
async def _update_entry(self, file_path: Path, deleted: bool) -> None:
- if deleted:
- # Remove by matching path (hash not available after deletion)
- to_remove = [
- e.id for e in self._index.entries
- if (self.root / e.path / e.name).resolve() == file_path.resolve()
- ]
- for fid in to_remove:
- self._index.remove_entry(fid)
- log.debug("Removed from index: %s", file_path.name)
- else:
+ self._remove_by_path(file_path)
+
+ if not deleted:
loop = asyncio.get_event_loop()
entry = await loop.run_in_executor(
self._executor, _scan_file, self.root, file_path)
if entry:
self._index.add_entry(entry)
- log.debug("Indexed: %s (%s)", file_path.name, entry.id[:8])
+ log.debug("Indexed: %s (%s, %d bytes)",
+ file_path.name, entry.id[:8], entry.size)
self._index.version = int(time.time())
if self.on_change: