From ddf9b9f8b43a348be586c1c0315b6b8f6dd8d81d Mon Sep 17 00:00:00 2001 From: sammywachtel Date: Tue, 6 Oct 2026 14:26:29 -0400 Subject: [PATCH] fix(core): index files in a new directory the watcher cannot see into On Linux, files copied into a running server in bulk could stay on disk and out of the index, whole folders at a time, with nothing logged. watchfiles (over notify's inotify backend) learns of a new directory from its parent and then adds a watch of its own. If the directory cannot be read at that instant the add fails and notify discards the error, so no file written into it ever produces an event. `install -d -o ` run as root creates directories this way, and so does any copy or restore made as root and chowned afterwards. The directory's own creation event still arrives, but the watcher dropped directory events. The periodic restart re-added the watch without looking for files that had landed meanwhile, and it discarded changes collected but not yet handed over. - A reported new directory is walked and its files indexed in that batch. A directory renamed inside the project moves its rows instead of doubling them. - A maintenance tick restarts the watcher only when the project set changed or a new directory appeared, and each restart is followed by a reconcile. - A periodic reconcile indexes settled files on disk that the index does not know, and logs a warning naming them. Its walk runs in a worker thread. Signed-off-by: sammywachtel --- src/basic_memory/config_models.py | 2 +- src/basic_memory/index/watch_reconcile.py | 196 +++++++ src/basic_memory/index/watch_service.py | 195 ++++++- src/basic_memory/services/initialization.py | 2 + tests/index/test_watch_bulk_copy.py | 534 ++++++++++++++++++++ 5 files changed, 922 insertions(+), 7 deletions(-) create mode 100644 src/basic_memory/index/watch_reconcile.py create mode 100644 tests/index/test_watch_bulk_copy.py diff --git a/src/basic_memory/config_models.py b/src/basic_memory/config_models.py index cd705a7ce..ea5b46481 100644 --- a/src/basic_memory/config_models.py +++ b/src/basic_memory/config_models.py @@ -546,7 +546,7 @@ def __init__(self, **data: Any) -> None: ... watch_project_reload_interval: int = Field( default=300, - description="Seconds between reloading project list in watch service. Higher values reduce CPU usage by minimizing watcher restarts. Default 300s (5 min) balances efficiency with responsiveness to new projects.", + description="Seconds between watch service maintenance ticks. Each tick re-reads the project list; it restarts the watcher only when the projects changed or a new directory appeared, and otherwise reconciles disk against the index, indexing any file the watcher missed. Default 300s (5 min) bounds how long a missed file can stay out of the index.", gt=0, ) diff --git a/src/basic_memory/index/watch_reconcile.py b/src/basic_memory/index/watch_reconcile.py new file mode 100644 index 000000000..de48e57f1 --- /dev/null +++ b/src/basic_memory/index/watch_reconcile.py @@ -0,0 +1,196 @@ +"""Find files on disk that the local index does not know. + +The watcher is the only thing that indexes a file written while the server runs. +Anything that makes it miss an event leaves that file on disk and out of search, +the graph and every tool, with nothing reported. The project scan at startup is +what used to repair that, which is why a restart always seemed to fix it. + +Two things here. `expand_new_directories` turns a reported new directory into the +files inside it, because on Linux those files may never be reported themselves. +`settled_project_files` and `indexed_file_paths` are the two halves of the +reconcile: what is on disk, and what the index holds. The walks and stat calls +are blocking filesystem work, so both walk functions are meant to run in a worker +thread, never on the event loop. +""" + +from __future__ import annotations + +import os +import time +from pathlib import Path + +from watchfiles import Change +from watchfiles.main import FileChange +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from basic_memory import db +from basic_memory.ignore_utils import load_gitignore_patterns, should_ignore_path +from basic_memory.index.filesystem import local_relative_path_is_filtered +from basic_memory.models import Entity + +#: A file's stat signature: (mtime_ns, size). Used to recognize a file that a +#: reconcile already tried to index and that has not changed since. +FileSignature = tuple[int, int] + + +def expand_new_directories( + changes: set[FileChange], + *, + project_root: Path, + ignore_patterns: set[str], +) -> tuple[set[FileChange], tuple[Path, ...]]: + """Add every file inside a newly created directory to the batch as added. + + WHY. On Linux the watcher (watchfiles, over notify's inotify backend) learns + of a new directory from its parent's watch and then adds a watch of its own + for it. If that directory cannot be read at that instant, the add fails and + notify discards the error: nothing is reported, and no file written into the + directory afterwards produces an event. `install -d -o 1000 -g 1000` run as + root does exactly this -- the directory exists, owned by root and closed, for + a moment before it is handed over -- so does a restore or an SSH copy made as + root and chowned after. A bulk copy into a running server then leaves whole + folders unindexed, and which folders depends on scheduling. + + The directory's own creation is still reported, because the parent's watch + is fine. So treat that report as "everything in here is new": walk it and + add each eligible file. Files written into it after this batch are the + reconcile's job, and the restart it prompts gets the directory watched. + + Returns the expanded batch and the new directories found (top-level ones as + reported; their subdirectories are walked, not listed). + """ + root = project_root.expanduser().resolve() + expanded = set(changes) + new_directories: list[Path] = [] + for change, path in changes: + if change != Change.added: + continue + directory = Path(path) + try: + if directory.is_symlink() or not directory.is_dir(): + continue + relative = directory.resolve().relative_to(root).as_posix() + except (OSError, ValueError): + continue + if local_relative_path_is_filtered(relative) or should_ignore_path( + directory, root, ignore_patterns + ): + continue + new_directories.append(directory) + for dirpath, dirnames, filenames in os.walk(directory, followlinks=False): + current = Path(dirpath) + dirnames[:] = [ + name + for name in dirnames + if not name.startswith(".") + and not should_ignore_path(current / name, root, ignore_patterns) + ] + for name in filenames: + file_path = current / name + try: + if file_path.is_symlink() or not file_path.is_file(): + continue + relative_file = file_path.resolve().relative_to(root).as_posix() + except (OSError, ValueError): + continue + if local_relative_path_is_filtered(relative_file) or should_ignore_path( + file_path, root, ignore_patterns + ): + continue + expanded.add((Change.added, str(file_path))) + return expanded, tuple(new_directories) + + +def expand_deleted_directories( + changes: set[FileChange], + *, + project_root: Path, + indexed_paths: set[str], +) -> set[FileChange]: + """Add a delete for every indexed file under a directory the batch reports gone. + + The counterpart of `expand_new_directories`, and what keeps it safe. A + directory renamed or moved inside a project arrives as one deleted directory + and one added directory, with no events for the files in either. Expanding + only the added side would index every file a second time under its new path + while the old rows stayed. With both sides expanded, the move processor pairs + each old path with its new one by checksum and moves the rows, as it does for + a single moved file. A path that still exists is left alone: the delete + planner only acts on confirmed absence anyway. + """ + root = project_root.expanduser().resolve() + expanded = set(changes) + for change, path in changes: + if change != Change.deleted: + continue + candidate = Path(path) + if candidate.exists(): + continue + try: + prefix = candidate.resolve().relative_to(root).as_posix().rstrip("/") + "/" + except ValueError: + continue + if prefix == "./": + continue + for indexed_path in indexed_paths: + if indexed_path.startswith(prefix): + expanded.add((Change.deleted, str(root / indexed_path))) + return expanded + + +def settled_project_files( + project_root: Path, + *, + ignore_patterns: set[str] | None = None, + settle_seconds: float, + now: float | None = None, +) -> dict[str, FileSignature]: + """Return project-relative files that have been still for `settle_seconds`. + + Eligibility is the startup scan's (`scan_local_project_index_files`): the + same ignore rules, hidden paths and symlink handling, so the reconcile never + indexes a file the startup scan would leave alone. + + "Still" is judged on the later of mtime and ctime. A copy that preserves + timestamps (`cp -p`, `tar`, `rsync -a`, a backup restore) gives a file an old + mtime the moment it lands, but its ctime is the time of the copy, so a file + still waiting in the watcher's own debounce window is not taken for a missed + one. + """ + # Imported here: local_project pulls in the services package, which imports + # the watcher lazily; a module-level import would close that loop. + from basic_memory.index.local_project import scan_local_project_index_files + + scan = scan_local_project_index_files( + project_root, + ignore_patterns=( + ignore_patterns + if ignore_patterns is not None + else load_gitignore_patterns(project_root) + ), + ) + cutoff_ns = int(((time.time() if now is None else now) - settle_seconds) * 1_000_000_000) + root = project_root.expanduser().resolve() + settled: dict[str, FileSignature] = {} + for relative_path in scan.file_paths: + try: + stat_result = os.stat(root / relative_path) + except OSError: + # Gone or unreadable between the walk and the stat: not this pass's business. + continue + if max(stat_result.st_mtime_ns, stat_result.st_ctime_ns) > cutoff_ns: + continue + settled[relative_path] = (stat_result.st_mtime_ns, stat_result.st_size) + return settled + + +async def indexed_file_paths( + session_maker: async_sessionmaker[AsyncSession], + project_id: int, +) -> set[str]: + """Return every file path the index holds for one project (one cheap query).""" + query = select(Entity.file_path).where(Entity.project_id == project_id) + async with db.scoped_session(session_maker) as session: + rows = (await session.execute(query)).scalars().all() + return {str(file_path) for file_path in rows} diff --git a/src/basic_memory/index/watch_service.py b/src/basic_memory/index/watch_service.py index b0c10e8de..e15b35ade 100644 --- a/src/basic_memory/index/watch_service.py +++ b/src/basic_memory/index/watch_service.py @@ -5,7 +5,7 @@ import asyncio import os import time -from collections.abc import Sequence +from collections.abc import Callable, Sequence from datetime import datetime from pathlib import Path from typing import Protocol @@ -32,6 +32,13 @@ run_local_watch_event_indexing, ) from basic_memory.index.storage_events import StorageEventIndexRuntime +from basic_memory.index.watch_reconcile import ( + FileSignature, + expand_deleted_directories, + expand_new_directories, + indexed_file_paths, + settled_project_files, +) from basic_memory.models import Project from basic_memory.repository import ProjectRepository from basic_memory.utils import generate_permalink @@ -104,6 +111,7 @@ def __init__( quiet: bool = False, event_index_runtime_factory: WatchEventIndexRuntimeFactory | None = None, constrained_project: str | None = None, + initial_index_pending: Callable[[], bool] | None = None, ) -> None: self.app_config = app_config self.project_repository = project_repository @@ -121,12 +129,160 @@ def __init__( ) self.constrained_project = constrained_project self.console = Console(quiet=quiet) - - async def _schedule_restart(self, stop_event: asyncio.Event) -> None: - """Schedule a watch cycle restart so project config changes are observed.""" - await asyncio.sleep(self.app_config.watch_project_reload_interval) + # State for finding files the watcher missed. + # See _schedule_restart, expand_new_directories and reconcile_unindexed_files. + self.reconcile_enabled = True + self.reconcile_settle_seconds = 60.0 + self._initial_index_pending = initial_index_pending or (lambda: False) + self._batches_in_flight = 0 + self._last_batch_finished = 0.0 + self._reconcile_attempted: dict[int, dict[str, FileSignature]] = {} + self._new_directories_seen = False + self._reconcile_after_restart = False + + async def _schedule_restart( + self, + stop_event: asyncio.Event, + projects: Sequence[Project] = (), + ) -> None: + """Maintain one watch cycle: reconcile, re-check projects, restart when needed. + + This used to stop the cycle on every tick, to pick up project changes, and + that restart was lossy twice over. Stopping discards the changes the watcher + has collected but not yet handed over (watchfiles clears them when it sees + the stop event), and the next cycle's fresh watcher never hears of files + that already exist; while a slow batch was being indexed, a tick dropped + everything written in the meantime. Nothing ever looked for what had been + dropped. + + Now a tick restarts only when there is a reason -- the project set changed, + or a directory appeared that the watcher may not be watching (see + expand_new_directories) -- and every restart is followed, once indexing is + quiet, by a reconcile that indexes whatever the gap cost. A tick with no + reason to restart runs the reconcile on its own: the periodic safety net. + """ + watched = _watched_project_identity(projects) + if self._reconcile_after_restart and self.reconcile_enabled: + await asyncio.sleep(self._quiet_seconds()) + if await self._reconcile_when_idle(projects): + self._reconcile_after_restart = False + while not stop_event.is_set(): + await asyncio.sleep(self.app_config.watch_project_reload_interval) + try: + current = await self._select_projects_to_watch() + except Exception as exc: + # Trigger: the project list could not be read (database unavailable). + # Why: carrying on would leave a stale project set watched forever. + # Outcome: restart the cycle, which retries the read, as upstream did. + logger.warning(f"Watch project reload failed, restarting watch cycle: {exc}") + self._restart(stop_event) + return + if self._project_set_changed(watched, current): + logger.info("Watched project set changed; restarting watch cycle") + self._restart(stop_event) + return + if self._new_directories_seen: + logger.info( + "New directories appeared since the watcher started; restarting the " + "watch cycle so that every one of them is watched" + ) + self._restart(stop_event) + return + # Pick up .gitignore/.bmignore edits, which the restart used to do. + self._ignore_patterns_cache.clear() + if self.reconcile_enabled and await self._reconcile_when_idle(current): + self._reconcile_after_restart = False + + def _restart(self, stop_event: asyncio.Event) -> None: + """End this watch cycle; the next one reconciles what the gap cost.""" + self._reconcile_after_restart = True stop_event.set() + def _quiet_seconds(self) -> float: + """How long indexing must have been idle before a reconcile may run.""" + return 3 * self.app_config.index_delay / 1000 + 5 + + def _project_set_changed( + self, + watched: frozenset[tuple[str, str]], + current: Sequence[Project], + ) -> bool: + """Return whether the projects to watch differ from the ones being watched.""" + return _watched_project_identity(current) != watched + + async def _reconcile_when_idle(self, projects: Sequence[Project]) -> bool: + """Run the reconcile unless indexing is busy right now. Return whether it ran. + + A batch in flight, or one that just finished, can leave changes queued in + the watcher that it will deliver in a moment; reconciling then would index + them twice and report them as missed. The startup scan has the same view + of a cold project. The next tick tries again. + """ + if ( + self._batches_in_flight + or self._initial_index_pending() + or time.monotonic() - self._last_batch_finished < self._quiet_seconds() + ): + logger.debug("Skipping index reconcile: indexing is busy") + return False + try: + await self.reconcile_unindexed_files(projects) + except Exception as exc: + logger.exception(f"Index reconcile failed: {exc}") + self.state.record_error(str(exc)) + await self.write_status() + return True + + async def reconcile_unindexed_files(self, projects: Sequence[Project]) -> int: + """Index every settled file on disk that the index does not know. Return the count. + + The safety net under the watcher. Whatever made + the watcher miss a file, this finds it within one reload interval instead + of at the next restart, and says so in the log, because a reconcile that + has to index something means an event was lost and that is worth knowing. + + A file it already tried, which is still not in the index and has not + changed since, is not retried or reported again on every tick. + """ + reconciled = 0 + for project in projects: + if not self._project_is_configured(project): + continue + project_root = local_project_root(project) + on_disk = await asyncio.to_thread( + settled_project_files, + project_root, + ignore_patterns=self._get_ignore_patterns(project_root), + settle_seconds=self.reconcile_settle_seconds, + ) + indexed = await indexed_file_paths(self.session_maker, project.id) + unindexed = { + path: signature for path, signature in on_disk.items() if path not in indexed + } + already_tried = self._reconcile_attempted.get(project.id, {}) + to_index = sorted( + path + for path, signature in unindexed.items() + if already_tried.get(path) != signature + ) + self._reconcile_attempted[project.id] = unindexed + if not to_index: + continue + + logger.warning( + f"Index reconcile: {len(to_index)} file(s) in project {project.name} were " + "on disk but not in the index, so the watcher missed them; indexing now", + project=project.name, + file_count=len(to_index), + paths=to_index[:20], + ) + await self._handle_changes_isolated( + project, + {(Change.added, str(project_root / path)) for path in to_index}, + ) + reconciled += len(to_index) + return reconciled + def _get_ignore_patterns(self, project_path: Path) -> set[str]: """Return cached ignore patterns for one project root.""" if project_path not in self._ignore_patterns_cache: @@ -141,6 +297,8 @@ async def _watch_projects_cycle( """Run one watchfiles cycle and route batches into project-local indexing.""" project_paths = [project.path for project in projects] previous_filter_roots = self._sorted_watch_filter_roots + # A fresh watcher walks the tree and watches all of it. + self._new_directories_seen = False self._sorted_watch_filter_roots = local_watch_filter_roots(projects) try: @@ -222,7 +380,7 @@ async def run(self) -> None: # pragma: no cover f"{[project.path for project in projects]}" ) stop_event = asyncio.Event() - timer_task = asyncio.create_task(self._schedule_restart(stop_event)) + timer_task = asyncio.create_task(self._schedule_restart(stop_event, projects)) try: await self._watch_projects_cycle(projects, stop_event) @@ -292,6 +450,7 @@ async def _handle_changes_isolated(self, project: Project, changes: set[FileChan this watch cycle. Outcome: log + record the error and continue so other projects still index. """ + self._batches_in_flight += 1 # the reconcile waits for idle try: await self.handle_changes(project, changes) except Exception as exc: @@ -300,6 +459,9 @@ async def _handle_changes_isolated(self, project: Project, changes: set[FileChan ) self.state.record_error(str(exc)) await self.write_status() + finally: + self._batches_in_flight -= 1 + self._last_batch_finished = time.monotonic() async def handle_changes(self, project: Project, changes: set[FileChange]) -> None: """Normalize one project's watchfiles batch and process it through indexing.""" @@ -317,6 +479,21 @@ async def handle_changes(self, project: Project, changes: set[FileChange]) -> No start_time = time.time() project_root = local_project_root(project) + # A new directory's files may never be reported; see expand_new_directories. + changes, new_directories = await asyncio.to_thread( + expand_new_directories, + changes, + project_root=project_root, + ignore_patterns=self._get_ignore_patterns(project_root), + ) + if new_directories: + self._new_directories_seen = True + if any(change == Change.deleted for change, _path in changes): + changes = expand_deleted_directories( + changes, + project_root=project_root, + indexed_paths=await indexed_file_paths(self.session_maker, project.id), + ) request = LocalWatchEventIndexRequest.from_project_changes( project=project, changes=changes, @@ -352,3 +529,9 @@ async def handle_changes(self, project: Project, changes: set[FileChange]) -> No f"duration_ms={duration_ms}" ) await self.write_status() + + +def _watched_project_identity(projects: Sequence[Project]) -> frozenset[tuple[str, str]]: + """Name and path of each watched project: what a watch cycle depends on.""" + # A maintenance tick restarts the cycle only when this changes. + return frozenset((str(project.name), str(project.path)) for project in projects) diff --git a/src/basic_memory/services/initialization.py b/src/basic_memory/services/initialization.py index 01d9ba7fe..b14ee7b00 100644 --- a/src/basic_memory/services/initialization.py +++ b/src/basic_memory/services/initialization.py @@ -279,6 +279,8 @@ async def initialize_file_indexing( quiet=quiet, event_index_runtime_factory=event_index_runtime_factory, constrained_project=constrained_project, + # The watcher's reconcile stays out of the startup scan's way. + initial_index_pending=lambda: bool(_initial_index_tasks), ) # Get active projects diff --git a/tests/index/test_watch_bulk_copy.py b/tests/index/test_watch_bulk_copy.py new file mode 100644 index 000000000..a5ba7cb0f --- /dev/null +++ b/tests/index/test_watch_bulk_copy.py @@ -0,0 +1,534 @@ +"""Files copied in bulk into a running server must all be indexed, without a restart. + +Whole folders of notes copied in this way stayed on disk and out of the index, +with nothing logged. Two causes, both covered here: + +* On Linux a new directory is watched only after it appears. If it cannot be + read at that instant (created closed and then handed to another owner, which + is what `install -d -o ` run as root does), the watch is never added and + no file written into it is ever reported. The end-to-end tests at the bottom + reproduce that with a real watcher and real indexing. +* The watch loop restarted itself on every project-reload tick, and a restart + discards the changes collected but not yet handed over: during a slow batch, + everything written in the meantime. + +The restart tests drive the real watch loop (`WatchService.run`, real +watchfiles, a real directory) with the indexing step replaced by one that +records what it was handed and is slow on its first batch. +""" + +from __future__ import annotations + +import asyncio +import os +import sys +import time +from collections.abc import Sequence +from pathlib import Path +from typing import override + +import pytest +from loguru import logger +from watchfiles import Change + +from basic_memory import db +from basic_memory.config import BasicMemoryConfig +from basic_memory.index.local_runtime import LocalWatchEventIndexRuntimeFactory +from basic_memory.index.watch_reconcile import expand_new_directories +from basic_memory.index.watch_service import WatchService +from basic_memory.models import Project + +#: How long the copy may take to be fully indexed. The scenario needs about four +#: seconds; the rest is room for a slow CI filesystem watcher. +CONVERGENCE_CEILING_SECONDS = 20.0 + +#: How long the first batch takes to "index", in seconds. Longer than the reload +#: interval, so a reload tick always lands while the copy's second half waits. +SLOW_BATCH_SECONDS = 3.0 + + +class RecordingWatchService(WatchService): + """The real watch loop, with indexing replaced by a slow recorder.""" + + def __init__(self, *args, project: Project, **kwargs) -> None: + super().__init__(*args, **kwargs) + self.project = project + self.seen: set[str] = set() + self.first_batch_started = asyncio.Event() + self.batches = 0 + self.cycles = 0 + # This file tests the watcher; the reconcile would hide what it does. + self.reconcile_enabled = False + + @override + async def _select_projects_to_watch(self) -> list[Project]: + return [self.project] + + @override + async def _watch_projects_cycle(self, projects, stop_event) -> None: + self.cycles += 1 + await super()._watch_projects_cycle(projects, stop_event) + + @override + async def handle_changes(self, project, changes) -> None: # type: ignore[override] + self.batches += 1 + first = self.batches == 1 + # As the real handler does: a reported new directory stands for its files. + # On Linux the first note can land before the new folder's watch exists. + changes, _new = expand_new_directories( + changes, project_root=Path(project.path), ignore_patterns=set() + ) + self.seen.update(Path(path).name for _change, path in changes if path.endswith(".md")) + if first: + self.first_batch_started.set() + await asyncio.sleep(SLOW_BATCH_SECONDS) + + @override + async def write_status(self) -> None: + return None + + +class RestartEveryTickWatchService(RecordingWatchService): + """The negative control: the old behavior, a restart on every reload tick. + + Only from the first batch on. A restart before then could swallow the first + write as well, and the control would fail for a reason that is not the defect. + """ + + @override + def _project_set_changed(self, watched, current: Sequence[Project]) -> bool: + return self.first_batch_started.is_set() + + +def _watch_config(app_config: BasicMemoryConfig) -> BasicMemoryConfig: + return app_config.model_copy(update={"watch_project_reload_interval": 1, "index_delay": 200}) + + +async def _bulk_copy_into_busy_watcher( + service: RecordingWatchService, vault: Path +) -> tuple[set[str], float | None]: + """Copy two folders in while the first batch is still indexing. + + Returns the note names the watcher delivered, and how long after the copy + finished the last one arrived (None if they never all did). + """ + expected = {"Charter.md"} | {f"Person {i}.md" for i in range(12)} + run_task = asyncio.create_task(service.run()) + try: + # Let the watcher start. Writing before it does would test nothing. + await asyncio.sleep(1.0) + (vault / "Charter.md").write_text("# Charter\n") + try: + await asyncio.wait_for(service.first_batch_started.wait(), timeout=10) + except TimeoutError: + # Even the first write was lost (a restart between it and its delivery). + return service.seen, None + + # The first batch is now "indexing" for SLOW_BATCH_SECONDS. The rest of the + # copy lands in a folder the copy creates, as a bulk copy does. + people = vault / "people" + people.mkdir() + for i in range(12): + (people / f"Person {i}.md").write_text(f"# Person {i}\n") + copied_at = time.monotonic() + + deadline = copied_at + CONVERGENCE_CEILING_SECONDS + while time.monotonic() < deadline: + if expected <= service.seen: + return service.seen, time.monotonic() - copied_at + await asyncio.sleep(0.1) + return service.seen, None + finally: + service.state.running = False + run_task.cancel() + try: + await run_task + except (asyncio.CancelledError, Exception): + pass + + +def _project(vault: Path) -> Project: + return Project(id=1, name="vault", permalink="vault", path=str(vault)) + + +@pytest.mark.asyncio +async def test_bulk_copy_into_a_busy_watcher_is_fully_delivered( + app_config: BasicMemoryConfig, project_repository, session_maker, tmp_path: Path +) -> None: + vault = tmp_path / "vault" + vault.mkdir() + service = RecordingWatchService( + app_config=_watch_config(app_config), + project_repository=project_repository, + session_maker=session_maker, + project=_project(vault), + ) + + seen, converged_after = await _bulk_copy_into_busy_watcher(service, vault) + + missing = sorted(({"Charter.md"} | {f"Person {i}.md" for i in range(12)}) - seen) + assert converged_after is not None, f"never delivered: {missing}" + assert converged_after < CONVERGENCE_CEILING_SECONDS + # Reload ticks came and went during the slow batch without a restart. + assert service.cycles == 1 + + +@pytest.mark.asyncio +async def test_negative_control_restart_on_every_tick_loses_the_copy( + app_config: BasicMemoryConfig, project_repository, session_maker, tmp_path: Path +) -> None: + """The harness can see the loss: with the old restart, the second folder never arrives. + + If this test starts failing, the scenario above no longer exercises the + defect, and its pass means nothing. + """ + vault = tmp_path / "vault" + vault.mkdir() + service = RestartEveryTickWatchService( + app_config=_watch_config(app_config), + project_repository=project_repository, + session_maker=session_maker, + project=_project(vault), + ) + + seen, converged_after = await _bulk_copy_into_busy_watcher(service, vault) + + assert converged_after is None + assert not ({f"Person {i}.md" for i in range(12)} & seen) + assert service.cycles >= 2 + + +# --- the reconcile ----------------------------------------------------------------- + + +def _reconcile_service(app_config, project_repository, session_maker) -> WatchService: + service = WatchService( + app_config=app_config, + project_repository=project_repository, + session_maker=session_maker, + event_index_runtime_factory=LocalWatchEventIndexRuntimeFactory(), + ) + service.reconcile_settle_seconds = 0 + return service + + +async def _indexed_paths(session_maker, entity_repository) -> set[str]: + async with db.scoped_session(session_maker) as session: + return {entity.file_path for entity in await entity_repository.find_all(session)} + + +@pytest.mark.asyncio +async def test_reconcile_indexes_a_file_the_watcher_never_reported( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, +) -> None: + """A note on disk with no watcher event behind it is found, indexed, and logged.""" + root = Path(project_config.home) + (root / "people").mkdir(parents=True, exist_ok=True) + (root / "people" / "Missed.md").write_text("# Missed\n\nnever reported\n") + service = _reconcile_service(app_config, project_repository, session_maker) + assert "people/Missed.md" not in await _indexed_paths(session_maker, entity_repository) + + warnings: list[str] = [] + sink = logger.add(lambda message: warnings.append(str(message)), level="WARNING") + try: + reconciled = await service.reconcile_unindexed_files([test_project]) + finally: + logger.remove(sink) + + assert reconciled == 1 + assert "people/Missed.md" in await _indexed_paths(session_maker, entity_repository) + assert any("Index reconcile" in line and "watcher missed" in line for line in warnings) + + # A second pass finds nothing to do and says nothing. + warnings.clear() + sink = logger.add(lambda message: warnings.append(str(message)), level="WARNING") + try: + assert await service.reconcile_unindexed_files([test_project]) == 0 + finally: + logger.remove(sink) + assert not any("Index reconcile" in line for line in warnings) + + +@pytest.mark.asyncio +async def test_reconcile_leaves_a_file_still_settling_to_the_watcher( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, +) -> None: + """A file written moments ago may still be in the watcher's debounce: hands off.""" + root = Path(project_config.home) + (root / "Fresh.md").write_text("# Fresh\n") + service = _reconcile_service(app_config, project_repository, session_maker) + service.reconcile_settle_seconds = 3600 + + assert await service.reconcile_unindexed_files([test_project]) == 0 + assert "Fresh.md" not in await _indexed_paths(session_maker, entity_repository) + + +@pytest.mark.asyncio +async def test_reconcile_does_not_retry_a_file_that_will_not_index( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, +) -> None: + """A file indexing cannot take is tried once, not on every tick until it changes.""" + root = Path(project_config.home) + (root / "Stubborn.md").write_text("# Stubborn\n") + handed: list[set[str]] = [] + + class NeverIndexes(WatchService): + @override + async def handle_changes(self, project, changes) -> None: # type: ignore[override] + handed.append({Path(path).name for _change, path in changes}) + + service = NeverIndexes( + app_config=app_config, + project_repository=project_repository, + session_maker=session_maker, + ) + service.reconcile_settle_seconds = 0 + + assert await service.reconcile_unindexed_files([test_project]) == 1 + assert await service.reconcile_unindexed_files([test_project]) == 0 + assert handed == [{"Stubborn.md"}] + + # Once it changes, it is worth another try. + (root / "Stubborn.md").write_text("# Stubborn\n\nedited, longer\n") + assert await service.reconcile_unindexed_files([test_project]) == 1 + + +@pytest.mark.asyncio +async def test_reconcile_waits_while_indexing_is_busy( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, +) -> None: + """Neither a batch in flight nor the startup scan shares the field with a reconcile.""" + calls: list[int] = [] + + class Counting(WatchService): + @override + async def reconcile_unindexed_files(self, projects) -> int: + calls.append(1) + return 0 + + pending = True + service = Counting( + app_config=app_config, + project_repository=project_repository, + session_maker=session_maker, + initial_index_pending=lambda: pending, + ) + + await service._reconcile_when_idle([test_project]) + assert calls == [] + + pending = False + service._batches_in_flight = 1 + await service._reconcile_when_idle([test_project]) + assert calls == [] + + service._batches_in_flight = 0 + await service._reconcile_when_idle([test_project]) + assert calls == [1] + + +# --- a new directory whose files the watcher never reports -------------------------- + + +@pytest.mark.asyncio +async def test_a_reported_new_directory_indexes_the_files_inside_it( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, +) -> None: + """The batch names only the directory, as it does when its own watch failed.""" + root = Path(project_config.home) + people = root / "people" + (people / "nested").mkdir(parents=True) + for name in ("Ann.md", "Bo.md", "nested/Cy.md"): + (people / name).write_text(f"# {Path(name).stem}\n") + (people / ".hidden.md").write_text("# hidden\n") + service = _reconcile_service(app_config, project_repository, session_maker) + + await service.handle_changes(test_project, {(Change.added, str(people))}) + + indexed = await _indexed_paths(session_maker, entity_repository) + assert {"people/Ann.md", "people/Bo.md", "people/nested/Cy.md"} <= indexed + assert "people/.hidden.md" not in indexed + assert service._new_directories_seen is True + + +@pytest.mark.asyncio +async def test_a_directory_moved_inside_the_project_moves_its_rows( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, +) -> None: + """A renamed folder arrives as one deleted and one added directory: rows move, none double.""" + root = Path(project_config.home) + (root / "old").mkdir() + (root / "old" / "Ann.md").write_text("# Ann\n") + (root / "old" / "Bo.md").write_text("# Bo\n") + service = _reconcile_service(app_config, project_repository, session_maker) + await service.handle_changes( + test_project, + {(Change.added, str(root / "old" / "Ann.md")), (Change.added, str(root / "old" / "Bo.md"))}, + ) + async with db.scoped_session(session_maker) as session: + before = {e.file_path: e.id for e in await entity_repository.find_all(session)} + assert {"old/Ann.md", "old/Bo.md"} <= set(before) + + (root / "old").rename(root / "new") + await service.handle_changes( + test_project, {(Change.deleted, str(root / "old")), (Change.added, str(root / "new"))} + ) + + async with db.scoped_session(session_maker) as session: + after = {e.file_path: e.id for e in await entity_repository.find_all(session)} + assert not {"old/Ann.md", "old/Bo.md"} & set(after) + assert after["new/Ann.md"] == before["old/Ann.md"] + assert after["new/Bo.md"] == before["old/Bo.md"] + + +# --- end to end on Linux: the watch on a new directory fails ------------------------ + + +LINUX_NON_ROOT = sys.platform == "linux" and hasattr(os, "geteuid") and os.geteuid() != 0 + + +async def _copy_with_a_closed_directory( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, +) -> tuple[set[str], float | None]: + """Run the real watcher; create `people/` closed, then open it and copy notes in. + + A directory that cannot be read when the watcher learns of it is the shape + `install -d -o 1000 -g 1000` run as root produces for an instant. Here it is + held closed for a beat so the outcome does not depend on scheduling. Returns + what got indexed and how long after the copy the index matched the disk. + """ + root = Path(project_config.home) + config = app_config.model_copy(update={"watch_project_reload_interval": 1, "index_delay": 200}) + service = WatchService( + app_config=config, + project_repository=project_repository, + session_maker=session_maker, + event_index_runtime_factory=LocalWatchEventIndexRuntimeFactory(), + ) + service.reconcile_settle_seconds = 0 + expected = {f"people/Person {i}.md" for i in range(8)} + run_task = asyncio.create_task(service.run()) + try: + await asyncio.sleep(1.5) + people = root / "people" + people.mkdir(mode=0o000) + await asyncio.sleep(1.0) + people.chmod(0o755) + for i in range(8): + (people / f"Person {i}.md").write_text(f"# Person {i}\n") + await asyncio.sleep(0.1) + copied_at = time.monotonic() + deadline = copied_at + E2E_CEILING_SECONDS + indexed: set[str] = set() + while time.monotonic() < deadline: + indexed = await _indexed_paths(session_maker, entity_repository) + if expected <= indexed: + return indexed, time.monotonic() - copied_at + await asyncio.sleep(0.25) + return indexed, None + finally: + service.state.running = False + run_task.cancel() + try: + await run_task + except (asyncio.CancelledError, Exception): + pass + + +#: The new directory is walked in the batch that reports it (about a second), and +#: anything written after is indexed by the reconcile after the restart (reload +#: tick, then the quiet period of 3 x index_delay + 5 s). Twenty seconds is double that. +E2E_CEILING_SECONDS = 20.0 + + +@pytest.mark.skipif(not LINUX_NON_ROOT, reason="needs inotify, and a user a closed directory stops") +@pytest.mark.asyncio +async def test_a_copy_into_a_directory_the_watcher_cannot_watch_is_indexed( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, +) -> None: + indexed, converged_after = await _copy_with_a_closed_directory( + app_config, + project_repository, + session_maker, + test_project, + project_config, + entity_repository, + ) + missing = sorted({f"people/Person {i}.md" for i in range(8)} - indexed) + assert converged_after is not None, f"never indexed: {missing}" + assert converged_after < E2E_CEILING_SECONDS + + +@pytest.mark.skipif(not LINUX_NON_ROOT, reason="needs inotify, and a user a closed directory stops") +@pytest.mark.asyncio +async def test_negative_control_without_the_fix_the_copy_stays_unindexed( + app_config: BasicMemoryConfig, + project_repository, + session_maker, + test_project: Project, + project_config, + entity_repository, + monkeypatch, +) -> None: + """With directory expansion and the reconcile off, the same copy is never indexed. + + This is the defect as shipped: proof that the scenario above reproduces it. + """ + from basic_memory.index import watch_service as watch_service_module + + monkeypatch.setattr( + watch_service_module, "expand_new_directories", lambda changes, **_: (changes, ()) + ) + monkeypatch.setattr(WatchService, "reconcile_unindexed_files", _reconcile_nothing) + + indexed, converged_after = await _copy_with_a_closed_directory( + app_config, + project_repository, + session_maker, + test_project, + project_config, + entity_repository, + ) + assert converged_after is None + assert not {f"people/Person {i}.md" for i in range(8)} & indexed + + +async def _reconcile_nothing(self, projects) -> int: + return 0