diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/indexer')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/group_index.py | 8 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/indexer.py | 51 |
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: |