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