From cf252b144c4b44daaa4ab26af70041a713b46f8f Mon Sep 17 00:00:00 2001 From: stumpylog <797416+stumpylog@users.noreply.github.com> Date: Tue, 8 Sep 2026 16:04:25 -0700 Subject: [PATCH] Experiments with tying the consumer to the task queue a little more, allowing use to determine if we should requeue a file The tricky bit is figuring out if we should, so we don't loop. Think this works? --- .../management/commands/document_consumer.py | 132 ++++++++---- .../test_document_consumer_stuck_queue.py | 165 +++++++++++++++ .../tests/test_management_consumer.py | 198 ++++++++++++++++-- 3 files changed, 440 insertions(+), 55 deletions(-) create mode 100644 src/documents/tests/test_document_consumer_stuck_queue.py diff --git a/src/documents/management/commands/document_consumer.py b/src/documents/management/commands/document_consumer.py index cc6535487..80aa58875 100644 --- a/src/documents/management/commands/document_consumer.py +++ b/src/documents/management/commands/document_consumer.py @@ -72,6 +72,24 @@ class TrackedFile: return False +@dataclass(frozen=True, slots=True) +class QueuedFile: + """A file handed to Celery, with enough state to decide when it's safe to re-check.""" + + task_id: str + size: int + mtime: float + + @classmethod + def from_path(cls, task_id: str, path: Path) -> QueuedFile | None: + """Snapshot the file's size and mtime, or None if it cannot be stat'd.""" + try: + stat = path.stat() + except OSError: + return None + return cls(task_id, stat.st_size, stat.st_mtime) + + class FileStabilityTracker: """ Tracks file events and determines when files are stable for consumption. @@ -314,7 +332,7 @@ def _consume_file( consumption_dir: Path, *, subdirs_as_tags: bool, -) -> bool: +) -> str | None: """ Queue a file for consumption. @@ -324,18 +342,18 @@ def _consume_file( subdirs_as_tags: Whether to create tags from subdirectory names. Returns: - True if the file was successfully handed to Celery, False otherwise. - Callers must not record the file as queued on failure, or the rescan - will never retry it. + The Celery task id if the file was successfully handed to Celery, + None otherwise. Callers must not record the file as queued on + failure, or the rescan will never retry it. """ # Verify file still exists and is accessible try: if not filepath.is_file(): logger.debug(f"Not consuming {filepath}: not a file or doesn't exist") - return False + return None except OSError as e: logger.warning(f"Not consuming {filepath}: {e}") - return False + return None # Get tags from path if configured tag_ids: list[int] | None = None @@ -348,7 +366,7 @@ def _consume_file( # Queue for consumption try: logger.info(f"Adding {filepath} to the task queue") - consume_file.apply_async( + result = consume_file.apply_async( kwargs={ "input_doc": ConsumableDocument( source=DocumentSource.ConsumeFolder, @@ -360,9 +378,9 @@ def _consume_file( ) except Exception: logger.exception(f"Error while queuing document {filepath}") - return False + return None - return True + return result.id class Command(BaseCommand): @@ -479,18 +497,19 @@ class Command(BaseCommand): recursive: bool, subdirs_as_tags: bool, consumer_filter: ConsumerFilter, - ) -> set[Path]: + ) -> dict[Path, QueuedFile]: """ Process any existing files in the consumption directory. - Returns the set of resolved paths that were queued, so the watch loop - can seed its in-flight set and avoid re-queuing them on the first - rescan before the consume tasks have removed them from disk. + Returns a dict mapping each resolved path that was queued to its + QueuedFile state, so the watch loop can seed its in-flight dict and + avoid re-queuing them on the first rescan before the consume tasks + have removed them from disk. """ logger.info(f"Processing existing files in {directory}") glob_pattern = "**/*" if recursive else "*" - queued: set[Path] = set() + queued: dict[Path, QueuedFile] = {} for filepath in directory.glob(glob_pattern): # Use filter to check if file should be processed @@ -500,12 +519,17 @@ class Command(BaseCommand): if not consumer_filter(Change.added, str(filepath)): continue - if _consume_file( + task_id = _consume_file( filepath=filepath, consumption_dir=directory, subdirs_as_tags=subdirs_as_tags, - ): - queued.add(filepath.resolve()) + ) + if task_id is None: + continue + + entry = QueuedFile.from_path(task_id, filepath) + if entry is not None: + queued[filepath.resolve()] = entry return queued @@ -516,21 +540,43 @@ class Command(BaseCommand): recursive: bool, consumer_filter: ConsumerFilter, tracker: FileStabilityTracker, - queued: set[Path], + queued: dict[Path, QueuedFile], ) -> None: """ Re-inject on-disk files the watcher never reported into the tracker. Acts as a safety net for files stranded by the watcher-recreation gap - (see ``rescan_interval_s``). Files already being tracked or already - queued and awaiting consumption are skipped, so a file is never queued - twice. Queued paths that have since left the directory are pruned so a - later file reusing the same name is not skipped forever. + (see ``rescan_interval_s``). Files already being tracked, or already + queued and still in flight (or completed but with unchanged content), + are skipped, so a file is never queued twice and a permanently broken + file does not retry forever. Queued paths that have since left the + directory are pruned so a later file reusing the same name is not + skipped forever. """ - # Prune in-flight paths that have left the directory - for path in list(queued): - if not path.exists(): - queued.discard(path) + # Long-running process: drop stale DB connections before querying (#4265) + db.close_old_connections() + + # Vanished from disk: consumed (or otherwise removed), prune regardless of status + for path in [path for path in queued if not path.exists()]: + del queued[path] + + if queued: + tasks = PaperlessTask.objects.only("task_id", "status").in_bulk( + [entry.task_id for entry in queued.values()], + field_name="task_id", + ) + for path, entry in list(queued.items()): + task = tasks.get(entry.task_id) + # No row yet means the task has not started: treat as in flight + if task is None or task.status not in PaperlessTask.COMPLETE_STATUSES: + continue + try: + current = path.stat() + except OSError: + continue + # Completed and the content changed: a new file, allow a retry + if current.st_size != entry.size or current.st_mtime != entry.mtime: + del queued[path] glob_pattern = "**/*" if recursive else "*" @@ -558,7 +604,7 @@ class Command(BaseCommand): polling_interval: float, stability_delay: float, is_testing: bool, - queued: set[Path] | None = None, + queued: dict[Path, QueuedFile] | None = None, ) -> None: """Watch directory for changes and process stable files.""" use_polling = polling_interval > 0 @@ -567,7 +613,7 @@ class Command(BaseCommand): # Resolved paths that have been queued and are awaiting consumption. # Seeded from the startup scan so the first rescan does not re-queue # files whose consume tasks have not yet removed them from disk. - queued = set() if queued is None else queued + queued = {} if queued is None else queued # Full-glob safety net cadence (0 disables) rescan_interval_s = self.rescan_interval_s @@ -644,7 +690,7 @@ class Command(BaseCommand): # Consumed (or otherwise removed); a later file # reusing this name must not be skipped as # already-queued. - queued.discard(path) + queued.pop(path, None) if not path.is_file(): continue if path in queued: @@ -663,12 +709,17 @@ class Command(BaseCommand): # rescan does not re-queue them while the consume task # has yet to remove them from disk, but does retry a # failed publish instead of stranding it - if _consume_file( + task_id = _consume_file( filepath=stable_path, consumption_dir=directory, subdirs_as_tags=subdirs_as_tags, - ): - queued.add(stable_path) + ) + if task_id is None: + continue + + entry = QueuedFile.from_path(task_id, stable_path) + if entry is not None: + queued[stable_path] = entry # Exit watch loop to reconfigure timeout break @@ -677,13 +728,16 @@ class Command(BaseCommand): if rescan_timeout_ms > 0 and ( monotonic() - last_rescan >= rescan_interval_s ): - self._rescan_existing_files( - directory=directory, - recursive=recursive, - consumer_filter=consumer_filter, - tracker=tracker, - queued=queued, - ) + try: + self._rescan_existing_files( + directory=directory, + recursive=recursive, + consumer_filter=consumer_filter, + tracker=tracker, + queued=queued, + ) + except Exception: + logger.exception("Error during consume folder rescan") last_rescan = monotonic() # Determine next timeout diff --git a/src/documents/tests/test_document_consumer_stuck_queue.py b/src/documents/tests/test_document_consumer_stuck_queue.py new file mode 100644 index 000000000..350ee3769 --- /dev/null +++ b/src/documents/tests/test_document_consumer_stuck_queue.py @@ -0,0 +1,165 @@ +""" +Regression test for GH discussion #13969. + +A consume-folder file that fails (e.g. a scanner's 0-byte placeholder +hitting "Unsupported mime type inode/x-empty") must be re-detected once +its content changes, not permanently stranded in the watcher's queued +set. See docs/superpowers/specs/2026-09-08-consume-folder-stuck-queue-spec.md. +""" + +from __future__ import annotations + +import shutil +from time import monotonic +from time import sleep +from typing import TYPE_CHECKING + +import pytest + +from documents.management.commands.document_consumer import Command +from documents.models import PaperlessTask +from documents.tests.test_management_consumer import consumption_dir # noqa: F401 +from documents.tests.test_management_consumer import ( + mock_consume_file_delay, # noqa: F401 +) +from documents.tests.test_management_consumer import ( + mock_supported_extensions, # noqa: F401 +) +from documents.tests.test_management_consumer import sample_pdf # noqa: F401 +from documents.tests.test_management_consumer import scratch_dir # noqa: F401 +from documents.tests.test_management_consumer import start_consumer # noqa: F401 +from documents.tests.test_management_consumer import wait_for_mock_call + +if TYPE_CHECKING: + from collections.abc import Callable + from pathlib import Path + from unittest.mock import MagicMock + + from pytest_mock import MockerFixture + + from documents.tests.test_management_consumer import ConsumerThread + + +@pytest.mark.management +# transaction=True: the background consumer thread needs to see rows +# committed by this test. +@pytest.mark.django_db(transaction=True) +class TestStuckQueueAfterConsumptionFailure: + def test_scanner_placeholder_recovers_after_failure( + self, + consumption_dir: Path, # noqa: F811 + sample_pdf: Path, # noqa: F811 + mock_consume_file_delay: MagicMock, # noqa: F811 + start_consumer: Callable[..., ConsumerThread], # noqa: F811 + ) -> None: + """ + Reproduces discussion #13969: a scanner creates a 0-byte file, + consumption fails on it, then the scanner writes real content — + the watcher must pick it up on the next rescan instead of + ignoring it forever. + """ + apply_async = mock_consume_file_delay.apply_async + apply_async.return_value.id = "scan-task-1" + + thread = start_consumer( + stability_delay=0.1, + rescan_interval=0.3, + ) + + target = consumption_dir / "scan.pdf" + target.write_bytes(b"") # scanner's 0-byte placeholder + + assert wait_for_mock_call(apply_async, timeout_s=5.0) + if thread.exception: + raise thread.exception + assert apply_async.call_count == 1 + + # The Celery task fails on inode/x-empty, exactly as consumer.py's + # mime-type check would in production. Create a real PaperlessTask row + # with FAILURE status to simulate the task's result when the real + # consumer.py rescan queries for it. This exercises the actual + # batched query path: PaperlessTask.objects.filter(task_id__in=[...]). + # values_list("task_id", "status") + PaperlessTask.objects.create( + task_id="scan-task-1", + task_type=PaperlessTask.TaskType.CONSUME_FILE, + status=PaperlessTask.Status.FAILURE, + ) + + # Scanner finishes writing the real scan. + shutil.copy(sample_pdf, target) + + # Needs to clear: rescan_interval (0.3s, until the entry is + # released) + a fresh stability_delay (0.1s, before _consume_file + # is called again) + polling slop + test margin. + deadline = monotonic() + 8.0 + while apply_async.call_count < 2 and monotonic() < deadline: + sleep(0.1) + + if thread.exception: + raise thread.exception + + assert apply_async.call_count == 2, ( + "Expected the file to be re-consumed after the scanner wrote " + f"real content, but apply_async was called " + f"{apply_async.call_count} time(s)" + ) + + +@pytest.mark.management +@pytest.mark.django_db(transaction=True) +class TestRescanErrorHandling: + def test_rescan_exception_does_not_break_the_watch_loop( + self, + consumption_dir: Path, # noqa: F811 + mock_consume_file_delay: MagicMock, # noqa: F811 + start_consumer: Callable[..., ConsumerThread], # noqa: F811 + mocker: MockerFixture, + ) -> None: + """ + A rescan that raises must not kill the watcher, and must not + cause the loop to busy-retry the database on every wake. The + watch loop updates ``last_rescan`` even when the rescan raises, + so a broken rescan still waits a full ``rescan_interval_s`` + between attempts instead of spinning. + + A file kept perpetually "pending" (rewritten faster than + ``stability_delay``) forces the watch loop to wake more often + than ``rescan_interval_s``, which is what surfaces the busy-retry + regression: with the broken shape, every one of those frequent + wakes re-attempts the rescan instead of only every + ``rescan_interval_s``. + """ + rescan = mocker.patch.object( + Command, + "_rescan_existing_files", + side_effect=Exception("db down"), + ) + + thread = start_consumer( + stability_delay=0.02, + rescan_interval=0.5, + ) + + # Keep a file perpetually unstable so the watch loop's timeout is + # floored at stability_delay (0.02s) rather than rescan_interval_s + # (0.5s). Otherwise the loop only wakes every 0.5s regardless of + # the rescan bug, and the two shapes would be indistinguishable. + target = consumption_dir / "busy.pdf" + deadline = monotonic() + 2.0 + counter = 0 + while monotonic() < deadline: + counter += 1 + target.write_bytes(f"%PDF-1.4\n{counter}\n".encode()) + sleep(0.01) + + assert thread.is_alive() + if thread.exception: + raise thread.exception + + assert rescan.called + assert rescan.call_count < 15, ( + "Expected the rescan to be retried roughly once per " + "rescan_interval_s, but it was called " + f"{rescan.call_count} time(s), suggesting a busy-loop" + ) diff --git a/src/documents/tests/test_management_consumer.py b/src/documents/tests/test_management_consumer.py index 46a7c957a..f405daacb 100644 --- a/src/documents/tests/test_management_consumer.py +++ b/src/documents/tests/test_management_consumer.py @@ -33,9 +33,11 @@ from documents.data_models import DocumentSource from documents.management.commands.document_consumer import Command from documents.management.commands.document_consumer import ConsumerFilter from documents.management.commands.document_consumer import FileStabilityTracker +from documents.management.commands.document_consumer import QueuedFile from documents.management.commands.document_consumer import TrackedFile from documents.management.commands.document_consumer import _consume_file from documents.management.commands.document_consumer import _tags_from_path +from documents.models import PaperlessTask from documents.models import Tag if TYPE_CHECKING: @@ -445,13 +447,14 @@ class TestConsumeFile: target = consumption_dir / "document.pdf" shutil.copy(sample_pdf, target) + mock_consume_file_delay.apply_async.return_value.id = "abc123" result = _consume_file( filepath=target, consumption_dir=consumption_dir, subdirs_as_tags=False, ) - assert result is True + assert result == mock_consume_file_delay.apply_async.return_value.id mock_consume_file_delay.apply_async.assert_called_once() call_args = mock_consume_file_delay.apply_async.call_args consumable_doc = call_args.kwargs["kwargs"]["input_doc"] @@ -470,7 +473,7 @@ class TestConsumeFile: consumption_dir=consumption_dir, subdirs_as_tags=False, ) - assert result is False + assert result is None mock_consume_file_delay.apply_async.assert_not_called() def test_consume_directory( @@ -487,7 +490,7 @@ class TestConsumeFile: consumption_dir=consumption_dir, subdirs_as_tags=False, ) - assert result is False + assert result is None mock_consume_file_delay.apply_async.assert_not_called() def test_consume_with_permission_error( @@ -507,7 +510,7 @@ class TestConsumeFile: consumption_dir=consumption_dir, subdirs_as_tags=False, ) - assert result is False + assert result is None mock_consume_file_delay.apply_async.assert_not_called() def test_consume_with_apply_async_failure( @@ -527,7 +530,7 @@ class TestConsumeFile: consumption_dir=consumption_dir, subdirs_as_tags=False, ) - assert result is False + assert result is None def test_consume_with_tags_error( self, @@ -545,12 +548,13 @@ class TestConsumeFile: side_effect=DatabaseError("Something happened"), ) + mock_consume_file_delay.apply_async.return_value.id = "abc123" result = _consume_file( filepath=target, consumption_dir=consumption_dir, subdirs_as_tags=True, ) - assert result is True + assert result == "abc123" mock_consume_file_delay.apply_async.assert_called_once() call_args = mock_consume_file_delay.apply_async.call_args overrides = call_args.kwargs["kwargs"]["overrides"] @@ -1116,6 +1120,7 @@ class TestCommandWatchEdgeCases: Tag.objects.all().delete() +@pytest.mark.django_db class TestRescanExistingFiles: """ Unit tests for the rescan safety net. @@ -1134,12 +1139,23 @@ class TestRescanExistingFiles: ignore_patterns=[], ) + def _queued_as_is(self, target: Path, task_id: str) -> dict[Path, QueuedFile]: + """A queued entry whose snapshot matches the file's current content.""" + stat = target.stat() + return { + target.resolve(): QueuedFile( + task_id=task_id, + size=stat.st_size, + mtime=stat.st_mtime, + ), + } + def _rescan( self, directory: Path, consumer_filter: ConsumerFilter, tracker: FileStabilityTracker, - queued: set[Path], + queued: dict[Path, QueuedFile], *, recursive: bool = False, ) -> None: @@ -1162,7 +1178,7 @@ class TestRescanExistingFiles: shutil.copy(sample_pdf, target) tracker = FileStabilityTracker(stability_delay=0.1) - self._rescan(consumption_dir, pdf_only_filter, tracker, set()) + self._rescan(consumption_dir, pdf_only_filter, tracker, {}) assert tracker.is_tracking(target) is True assert tracker.pending_count == 1 @@ -1179,7 +1195,7 @@ class TestRescanExistingFiles: tracker = FileStabilityTracker(stability_delay=0.1) tracker.track(target, Change.added) - self._rescan(consumption_dir, pdf_only_filter, tracker, set()) + self._rescan(consumption_dir, pdf_only_filter, tracker, {}) assert tracker.pending_count == 1 @@ -1193,11 +1209,17 @@ class TestRescanExistingFiles: target = consumption_dir / "inflight.pdf" shutil.copy(sample_pdf, target) tracker = FileStabilityTracker(stability_delay=0.1) - queued = {target.resolve()} + PaperlessTask.objects.create( + task_id="task-inflight", + task_type=PaperlessTask.TaskType.CONSUME_FILE, + status=PaperlessTask.Status.STARTED, + ) + queued = self._queued_as_is(target, "task-inflight") self._rescan(consumption_dir, pdf_only_filter, tracker, queued) assert tracker.pending_count == 0 + assert target.resolve() in queued def test_prunes_vanished_queued_paths( self, @@ -1207,7 +1229,7 @@ class TestRescanExistingFiles: """Queued paths no longer on disk are dropped so the name can recur.""" gone = (consumption_dir / "gone.pdf").resolve() tracker = FileStabilityTracker(stability_delay=0.1) - queued = {gone} + queued = {gone: QueuedFile(task_id="task-gone", size=0, mtime=0.0)} self._rescan(consumption_dir, pdf_only_filter, tracker, queued) @@ -1222,7 +1244,7 @@ class TestRescanExistingFiles: (consumption_dir / "notes.xyz").write_bytes(b"content") tracker = FileStabilityTracker(stability_delay=0.1) - self._rescan(consumption_dir, pdf_only_filter, tracker, set()) + self._rescan(consumption_dir, pdf_only_filter, tracker, {}) assert tracker.pending_count == 0 @@ -1239,13 +1261,151 @@ class TestRescanExistingFiles: shutil.copy(sample_pdf, target) shallow = FileStabilityTracker(stability_delay=0.1) - self._rescan(consumption_dir, pdf_only_filter, shallow, set()) + self._rescan(consumption_dir, pdf_only_filter, shallow, {}) assert shallow.pending_count == 0 deep = FileStabilityTracker(stability_delay=0.1) - self._rescan(consumption_dir, pdf_only_filter, deep, set(), recursive=True) + self._rescan(consumption_dir, pdf_only_filter, deep, {}, recursive=True) assert deep.is_tracking(target) is True + def test_completed_but_content_unchanged_stays_queued( + self, + consumption_dir: Path, + sample_pdf: Path, + pdf_only_filter: ConsumerFilter, + ) -> None: + """ + A task that failed (or succeeded) but whose file content never + changed since being queued stays put — this is what makes a + permanently-broken file (e.g. a corrupt PDF) fail once instead of + retrying forever (C2). + """ + target = consumption_dir / "broken.pdf" + shutil.copy(sample_pdf, target) + tracker = FileStabilityTracker(stability_delay=0.1) + PaperlessTask.objects.create( + task_id="task-failed-unchanged", + task_type=PaperlessTask.TaskType.CONSUME_FILE, + status=PaperlessTask.Status.FAILURE, + ) + queued = self._queued_as_is(target, "task-failed-unchanged") + + self._rescan(consumption_dir, pdf_only_filter, tracker, queued) + + assert target.resolve() in queued + assert tracker.pending_count == 0 + + def test_completed_and_content_changed_is_released( + self, + consumption_dir: Path, + sample_pdf: Path, + pdf_only_filter: ConsumerFilter, + ) -> None: + """ + A task that completed AND whose file content has since changed is + released and re-tracked — this is the discussion #13969 fix: the + scanner's 0-byte file failed, then real content arrived. + """ + target = consumption_dir / "scanned.pdf" + target.write_bytes(b"") # simulate the 0-byte placeholder that was queued + tracker = FileStabilityTracker(stability_delay=0.1) + PaperlessTask.objects.create( + task_id="task-failed-changed", + task_type=PaperlessTask.TaskType.CONSUME_FILE, + status=PaperlessTask.Status.FAILURE, + ) + queued = { + target.resolve(): QueuedFile( + task_id="task-failed-changed", + size=0, + mtime=0.0, + ), + } + shutil.copy(sample_pdf, target) # scanner writes real content + + self._rescan(consumption_dir, pdf_only_filter, tracker, queued) + + assert target.resolve() not in queued + assert tracker.is_tracking(target.resolve()) is True + + def test_no_paperlesstask_row_stays_queued( + self, + consumption_dir: Path, + sample_pdf: Path, + pdf_only_filter: ConsumerFilter, + ) -> None: + """No matching PaperlessTask row is treated as still in flight, not released.""" + target = consumption_dir / "no_row.pdf" + shutil.copy(sample_pdf, target) + tracker = FileStabilityTracker(stability_delay=0.1) + queued = self._queued_as_is(target, "task-does-not-exist") + + self._rescan(consumption_dir, pdf_only_filter, tracker, queued) + + assert target.resolve() in queued + + def test_revoked_status_is_treated_as_complete( + self, + consumption_dir: Path, + sample_pdf: Path, + pdf_only_filter: ConsumerFilter, + ) -> None: + """ + The release guard checks membership in COMPLETE_STATUSES (SUCCESS, + FAILURE, REVOKED), not just FAILURE. A cancelled/revoked task (e.g. + after a worker restart discards a stale queue entry) whose content + has since changed must also be released, not stuck treating REVOKED + as still in-flight. + """ + target = consumption_dir / "revoked.pdf" + target.write_bytes(b"") + tracker = FileStabilityTracker(stability_delay=0.1) + PaperlessTask.objects.create( + task_id="task-revoked", + task_type=PaperlessTask.TaskType.CONSUME_FILE, + status=PaperlessTask.Status.REVOKED, + ) + queued = { + target.resolve(): QueuedFile(task_id="task-revoked", size=0, mtime=0.0), + } + shutil.copy(sample_pdf, target) + + self._rescan(consumption_dir, pdf_only_filter, tracker, queued) + + assert target.resolve() not in queued + assert tracker.is_tracking(target.resolve()) is True + + def test_rescan_issues_one_batched_query_for_multiple_queued_files( + self, + consumption_dir: Path, + sample_pdf: Path, + pdf_only_filter: ConsumerFilter, + django_assert_num_queries, + ) -> None: + """ + The status lookup for every queued file must be a single batched + `task_id__in=[...]` query, not one query per file — otherwise a + consume folder with many in-flight files turns every rescan into + an N+1. + """ + tracker = FileStabilityTracker(stability_delay=0.1) + queued: dict[Path, QueuedFile] = {} + for i in range(3): + target = consumption_dir / f"doc{i}.pdf" + shutil.copy(sample_pdf, target) + task_id = f"task-batched-{i}" + PaperlessTask.objects.create( + task_id=task_id, + task_type=PaperlessTask.TaskType.CONSUME_FILE, + status=PaperlessTask.Status.STARTED, + ) + queued.update(self._queued_as_is(target, task_id)) + + with django_assert_num_queries(1): + self._rescan(consumption_dir, pdf_only_filter, tracker, queued) + + assert len(queued) == 3 + class TestProcessExistingFilesQueued: """Tests that startup processing reports which paths it queued.""" @@ -1258,7 +1418,8 @@ class TestProcessExistingFilesQueued: mock_consume_file_delay: MagicMock, settings: Settings, ) -> None: - """The set returned seeds the rescan's queued set, avoiding re-queue.""" + """The dict returned seeds the rescan's queued dict, avoiding re-queue.""" + mock_consume_file_delay.apply_async.return_value.id = "startup-task-id" target = consumption_dir / "document.pdf" shutil.copy(sample_pdf, target) settings.CONSUMER_IGNORE_PATTERNS = [] @@ -1271,6 +1432,9 @@ class TestProcessExistingFilesQueued: ) assert target.resolve() in queued + entry = queued[target.resolve()] + assert entry.task_id == "startup-task-id" + assert entry.size == target.stat().st_size @pytest.mark.management @@ -1295,11 +1459,13 @@ class TestCommandRetryAfterQueueFailure: """A publish failure from the watch loop is retried by the rescan.""" apply_async = mock_consume_file_delay.apply_async - def fail_first_call(*args: object, **kwargs: object) -> None: + def fail_first_call(*args: object, **kwargs: object) -> MagicMock | None: if apply_async.call_count == 1: raise Exception("broker down") + return apply_async.return_value apply_async.side_effect = fail_first_call + apply_async.return_value.id = "task_id_test" thread = start_consumer(stability_delay=0.1, rescan_interval=0.3)