From 55f77fc71ef96844e56d90b8a68b6050eda3bd9b Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 24 Sep 2026 15:58:38 +0800 Subject: [PATCH] fix(cascade): commit watcher upserts for a path in delivery order Each watcher event schedules its own upsert coroutine on the loop, and each upsert awaits the database, so two events for the same path could commit in either order and the last committer won. Windows synthesises a 'created' for every file under a freshly created parent directory, which hands the handler the same file four or five times; on the Windows soak box one of those stale duplicates committed after an atomic save's 'added' and put the first write's mtime and lsn back on the row (test_atomic_replace_over_existing_target_keeps _the_row_alive failed 1 run in 4, only there). Serialise the upserts behind one asyncio.Lock per handler; tasks are created in delivery order and the lock is FIFO, so the row ends with the last event. The scanner's sweep still writes on its own path. Verification: the new test fails against the previous watcher with committed == [m2, m1]; passes with the lock. The six existing watcher tests pass; lint-imports 4/4. Co-Authored-By: Claude Fable 5.1 --- src/everos/memory/cascade/watcher.py | 15 +++- .../test_cascade/test_watcher_upsert_order.py | 70 +++++++++++++++++++ 2 files changed, 84 insertions(+), 1 deletion(-) create mode 100644 tests/unit/test_memory/test_cascade/test_watcher_upsert_order.py diff --git a/src/everos/memory/cascade/watcher.py b/src/everos/memory/cascade/watcher.py index 9c30032b9..3810e3b28 100644 --- a/src/everos/memory/cascade/watcher.py +++ b/src/everos/memory/cascade/watcher.py @@ -14,6 +14,7 @@ from __future__ import annotations import asyncio +from collections.abc import Awaitable from pathlib import Path from watchdog.events import FileMovedEvent, FileSystemEvent, FileSystemEventHandler @@ -78,6 +79,14 @@ def __init__( ) -> None: self._memory_root = memory_root self._loop = loop + # Upserts for one path must commit in delivery order. Each one awaits + # the database, so left concurrent the last committer wins and a stale + # duplicate overwrites a newer row. Windows synthesises a ``created`` + # for every file under a freshly created parent directory, handing the + # same file to this handler four or five times; on the soak box one of + # those duplicates landed after an atomic save's ``added`` and put the + # first write's mtime back on the row (1 run in 4). + self._in_order = asyncio.Lock() def on_created(self, event: FileSystemEvent) -> None: self._enqueue(event.src_path, "added") @@ -124,10 +133,14 @@ def _enqueue(self, raw_path: str, change_type: str) -> None: return mtime = _safe_mtime(raw_path) asyncio.run_coroutine_threadsafe( - _enqueue_async(spec, rel, change_type, mtime), + self._serialised(_enqueue_async(spec, rel, change_type, mtime)), self._loop, ) + async def _serialised(self, upsert: Awaitable[None]) -> None: + async with self._in_order: + await upsert + async def _enqueue_async( spec: KindSpec, rel: str, change_type: str, mtime: float diff --git a/tests/unit/test_memory/test_cascade/test_watcher_upsert_order.py b/tests/unit/test_memory/test_cascade/test_watcher_upsert_order.py new file mode 100644 index 000000000..f6a8d62a9 --- /dev/null +++ b/tests/unit/test_memory/test_cascade/test_watcher_upsert_order.py @@ -0,0 +1,70 @@ +"""Upserts for one path commit in delivery order. + +Windows synthesises a ``created`` for every file under a freshly created +parent directory, so one ``mkdir -p`` + write hands the handler the same path +four or five times. Each upsert awaits the database; run concurrently, the +last committer wins, and on the Windows soak box a stale duplicate carrying +the first write's mtime overwrote the row of an atomic save (1 run in 4). +""" + +from __future__ import annotations + +import asyncio +import os +from pathlib import Path + +import pytest +from watchdog.events import FileCreatedEvent + +from everos.core.persistence import MemoryRoot +from everos.memory.cascade import watcher as watcher_mod +from everos.memory.cascade.watcher import _Handler + +_M1 = 1_700_000_000 +_M2 = 1_700_000_060 + + +class _SlowFirstRepo: + """The first upsert commits 50 ms late; the row is whatever committed last.""" + + def __init__(self) -> None: + self.calls = 0 + self.committed: list[float] = [] + + async def upsert( + self, md_path: str, *, kind: str, change_type: str, mtime: float + ) -> int: + self.calls += 1 + if self.calls == 1: + await asyncio.sleep(0.05) + self.committed.append(mtime) + return self.calls + + +async def test_duplicate_events_commit_in_delivery_order( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + repo = _SlowFirstRepo() + monkeypatch.setattr(watcher_mod, "md_change_state_repo", repo) + md = tmp_path.joinpath( + "default_app", + "default_project", + "users", + "u1", + "episodes", + "episode-2026-01-01.md", + ) + md.parent.mkdir(parents=True) + md.write_text("v1", encoding="utf-8") + handler = _Handler(MemoryRoot(tmp_path), asyncio.get_running_loop()) + + os.utime(md, (_M1, _M1)) + handler.on_created(FileCreatedEvent(str(md))) # the synthetic duplicate + os.utime(md, (_M2, _M2)) + handler.on_created(FileCreatedEvent(str(md))) # the save that must win + await asyncio.sleep(0.2) + + assert repo.committed == [_M1, _M2], ( + "the stale duplicate committed after the newer event; the row now " + "carries the old mtime" + )