From 76b291e21463900c8847192c4272ec20da8f3892 Mon Sep 17 00:00:00 2001 From: zmylol <1014124136@qq.com> Date: Sun, 27 Sep 2026 14:56:04 +0800 Subject: [PATCH] fix(reload): filter filesystem events before queueing --- pyproject.toml | 2 +- taskiq/cli/worker/process_manager.py | 13 +++ tests/cli/worker/test_process_manager.py | 101 +++++++++++++++++++++++ uv.lock | 2 +- 4 files changed, 116 insertions(+), 2 deletions(-) create mode 100644 tests/cli/worker/test_process_manager.py diff --git a/pyproject.toml b/pyproject.toml index 00748247..1d0a00e7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -62,7 +62,7 @@ opentelemetry = [ "psutil>=7", ] orjson = ["orjson>=3"] -reload = ["watchdog>=4", "gitignore-parser>=0"] +reload = ["watchdog>=4.0.2", "gitignore-parser>=0"] uv = ["uvloop>=0.16.0,<1; sys_platform != 'win32'"] zmq = ["pyzmq>=26"] diff --git a/taskiq/cli/worker/process_manager.py b/taskiq/cli/worker/process_manager.py index b5d27427..0abd761e 100644 --- a/taskiq/cli/worker/process_manager.py +++ b/taskiq/cli/worker/process_manager.py @@ -12,6 +12,12 @@ from typing import Any try: + from watchdog.events import ( + FileCreatedEvent, + FileDeletedEvent, + FileModifiedEvent, + FileMovedEvent, + ) from watchdog.observers import Observer from taskiq.cli.watcher import FileWatcher @@ -181,6 +187,13 @@ def __init__( ), path=path_to_watch, recursive=True, + # Filter open/close events before watchdog queues them. + event_filter=[ + FileModifiedEvent, + FileCreatedEvent, + FileDeletedEvent, + FileMovedEvent, + ], ) shutdown_handler = get_signal_handler(self.action_queue, ShutdownAction()) diff --git a/tests/cli/worker/test_process_manager.py b/tests/cli/worker/test_process_manager.py new file mode 100644 index 00000000..3a5c145f --- /dev/null +++ b/tests/cli/worker/test_process_manager.py @@ -0,0 +1,101 @@ +from pathlib import Path +from queue import Queue +from unittest.mock import Mock + +import pytest +from watchdog import events +from watchdog.events import ( + FileClosedEvent, + FileCreatedEvent, + FileDeletedEvent, + FileModifiedEvent, + FileMovedEvent, + FileOpenedEvent, + FileSystemEvent, +) +from watchdog.observers.api import BaseObserver, EventEmitter + +from taskiq.cli.worker import process_manager +from taskiq.cli.worker.args import WorkerArgs +from taskiq.cli.worker.process_manager import ProcessManager, ReloadAllAction + + +@pytest.fixture(params=[[], ["first", "second"]], ids=["default", "multiple-dirs"]) +def reloading_manager( + monkeypatch: pytest.MonkeyPatch, + tmp_path: Path, + request: pytest.FixtureRequest, +) -> tuple[ProcessManager, BaseObserver]: + monkeypatch.chdir(tmp_path) + monkeypatch.setattr(process_manager.signal, "signal", Mock()) + monkeypatch.setattr(process_manager, "Queue", Queue) + reload_dirs = [str(tmp_path / name) for name in request.param] + for directory in reload_dirs: + Path(directory).mkdir() + + # Leave the observer stopped so dispatch cannot hide unwanted queued events. + observer = BaseObserver(EventEmitter) + manager = ProcessManager( + WorkerArgs( + broker="example:broker", + modules=[], + reload=True, + reload_dirs=reload_dirs, + no_gitignore=True, + ), + worker_function=Mock(), + observer=observer, + ) + assert {emitter.watch.path for emitter in observer.emitters} == set( + reload_dirs or ["."], + ) + assert all(emitter.watch.is_recursive for emitter in observer.emitters) + return manager, observer + + +@pytest.mark.parametrize( + "event_class", + [ + FileOpenedEvent, + FileClosedEvent, + getattr(events, "FileClosedNoWriteEvent", None), + ], + ids=["opened", "closed", "closed-no-write"], +) +def test_open_and_close_events_never_enter_observer_queue( + reloading_manager: tuple[ProcessManager, BaseObserver], + event_class: type[FileSystemEvent] | None, +) -> None: + if event_class is None: + pytest.skip("This watchdog version does not emit read-only close events") + + manager, observer = reloading_manager + for emitter in observer.emitters: + for index in range(10): + emitter.queue_event(event_class(f"task_{index}.py")) + assert observer.event_queue.empty() + assert manager.action_queue.empty() + + +@pytest.mark.parametrize( + "event", + [ + FileCreatedEvent("task.py"), + FileModifiedEvent("task.py"), + FileDeletedEvent("task.py"), + FileMovedEvent("task.py", "renamed_task.py"), + ], + ids=["created", "modified", "deleted", "moved"], +) +def test_file_changes_queue_worker_reload( + reloading_manager: tuple[ProcessManager, BaseObserver], + event: FileSystemEvent, +) -> None: + manager, observer = reloading_manager + for emitter in observer.emitters: + emitter.queue_event(event) + assert not observer.event_queue.empty() + observer.dispatch_events(observer.event_queue) + assert isinstance(manager.action_queue.get_nowait(), ReloadAllAction) + assert manager.action_queue.empty() + assert observer.event_queue.empty() diff --git a/uv.lock b/uv.lock index 06c411f0..abfa96de 100644 --- a/uv.lock +++ b/uv.lock @@ -1827,7 +1827,7 @@ requires-dist = [ { name = "taskiq-dependencies", specifier = ">=1.3.1,<2" }, { name = "typing-extensions", marker = "python_full_version < '3.11'", specifier = ">=3.10.0.0" }, { name = "uvloop", marker = "sys_platform != 'win32' and extra == 'uv'", specifier = ">=0.16.0,<1" }, - { name = "watchdog", marker = "extra == 'reload'", specifier = ">=4" }, + { name = "watchdog", marker = "extra == 'reload'", specifier = ">=4.0.2" }, ] provides-extras = ["cbor", "metrics", "msgpack", "opentelemetry", "orjson", "reload", "uv", "zmq"]