Compare commits

..
Author SHA1 Message Date
stumpylog cf252b144c 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?
2026-09-08 19:33:59 -07:00
3 changed files with 440 additions and 55 deletions
@@ -72,6 +72,24 @@ class TrackedFile:
return False 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: class FileStabilityTracker:
""" """
Tracks file events and determines when files are stable for consumption. Tracks file events and determines when files are stable for consumption.
@@ -314,7 +332,7 @@ def _consume_file(
consumption_dir: Path, consumption_dir: Path,
*, *,
subdirs_as_tags: bool, subdirs_as_tags: bool,
) -> bool: ) -> str | None:
""" """
Queue a file for consumption. Queue a file for consumption.
@@ -324,18 +342,18 @@ def _consume_file(
subdirs_as_tags: Whether to create tags from subdirectory names. subdirs_as_tags: Whether to create tags from subdirectory names.
Returns: Returns:
True if the file was successfully handed to Celery, False otherwise. The Celery task id if the file was successfully handed to Celery,
Callers must not record the file as queued on failure, or the rescan None otherwise. Callers must not record the file as queued on
will never retry it. failure, or the rescan will never retry it.
""" """
# Verify file still exists and is accessible # Verify file still exists and is accessible
try: try:
if not filepath.is_file(): if not filepath.is_file():
logger.debug(f"Not consuming {filepath}: not a file or doesn't exist") logger.debug(f"Not consuming {filepath}: not a file or doesn't exist")
return False return None
except OSError as e: except OSError as e:
logger.warning(f"Not consuming {filepath}: {e}") logger.warning(f"Not consuming {filepath}: {e}")
return False return None
# Get tags from path if configured # Get tags from path if configured
tag_ids: list[int] | None = None tag_ids: list[int] | None = None
@@ -348,7 +366,7 @@ def _consume_file(
# Queue for consumption # Queue for consumption
try: try:
logger.info(f"Adding {filepath} to the task queue") logger.info(f"Adding {filepath} to the task queue")
consume_file.apply_async( result = consume_file.apply_async(
kwargs={ kwargs={
"input_doc": ConsumableDocument( "input_doc": ConsumableDocument(
source=DocumentSource.ConsumeFolder, source=DocumentSource.ConsumeFolder,
@@ -360,9 +378,9 @@ def _consume_file(
) )
except Exception: except Exception:
logger.exception(f"Error while queuing document {filepath}") logger.exception(f"Error while queuing document {filepath}")
return False return None
return True return result.id
class Command(BaseCommand): class Command(BaseCommand):
@@ -479,18 +497,19 @@ class Command(BaseCommand):
recursive: bool, recursive: bool,
subdirs_as_tags: bool, subdirs_as_tags: bool,
consumer_filter: ConsumerFilter, consumer_filter: ConsumerFilter,
) -> set[Path]: ) -> dict[Path, QueuedFile]:
""" """
Process any existing files in the consumption directory. Process any existing files in the consumption directory.
Returns the set of resolved paths that were queued, so the watch loop Returns a dict mapping each resolved path that was queued to its
can seed its in-flight set and avoid re-queuing them on the first QueuedFile state, so the watch loop can seed its in-flight dict and
rescan before the consume tasks have removed them from disk. 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}") logger.info(f"Processing existing files in {directory}")
glob_pattern = "**/*" if recursive else "*" glob_pattern = "**/*" if recursive else "*"
queued: set[Path] = set() queued: dict[Path, QueuedFile] = {}
for filepath in directory.glob(glob_pattern): for filepath in directory.glob(glob_pattern):
# Use filter to check if file should be processed # Use filter to check if file should be processed
@@ -500,12 +519,17 @@ class Command(BaseCommand):
if not consumer_filter(Change.added, str(filepath)): if not consumer_filter(Change.added, str(filepath)):
continue continue
if _consume_file( task_id = _consume_file(
filepath=filepath, filepath=filepath,
consumption_dir=directory, consumption_dir=directory,
subdirs_as_tags=subdirs_as_tags, 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 return queued
@@ -516,21 +540,43 @@ class Command(BaseCommand):
recursive: bool, recursive: bool,
consumer_filter: ConsumerFilter, consumer_filter: ConsumerFilter,
tracker: FileStabilityTracker, tracker: FileStabilityTracker,
queued: set[Path], queued: dict[Path, QueuedFile],
) -> None: ) -> None:
""" """
Re-inject on-disk files the watcher never reported into the tracker. 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 Acts as a safety net for files stranded by the watcher-recreation gap
(see ``rescan_interval_s``). Files already being tracked or already (see ``rescan_interval_s``). Files already being tracked, or already
queued and awaiting consumption are skipped, so a file is never queued queued and still in flight (or completed but with unchanged content),
twice. Queued paths that have since left the directory are pruned so a are skipped, so a file is never queued twice and a permanently broken
later file reusing the same name is not skipped forever. 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 # Long-running process: drop stale DB connections before querying (#4265)
for path in list(queued): db.close_old_connections()
if not path.exists():
queued.discard(path) # 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 "*" glob_pattern = "**/*" if recursive else "*"
@@ -558,7 +604,7 @@ class Command(BaseCommand):
polling_interval: float, polling_interval: float,
stability_delay: float, stability_delay: float,
is_testing: bool, is_testing: bool,
queued: set[Path] | None = None, queued: dict[Path, QueuedFile] | None = None,
) -> None: ) -> None:
"""Watch directory for changes and process stable files.""" """Watch directory for changes and process stable files."""
use_polling = polling_interval > 0 use_polling = polling_interval > 0
@@ -567,7 +613,7 @@ class Command(BaseCommand):
# Resolved paths that have been queued and are awaiting consumption. # Resolved paths that have been queued and are awaiting consumption.
# Seeded from the startup scan so the first rescan does not re-queue # Seeded from the startup scan so the first rescan does not re-queue
# files whose consume tasks have not yet removed them from disk. # 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) # Full-glob safety net cadence (0 disables)
rescan_interval_s = self.rescan_interval_s rescan_interval_s = self.rescan_interval_s
@@ -644,7 +690,7 @@ class Command(BaseCommand):
# Consumed (or otherwise removed); a later file # Consumed (or otherwise removed); a later file
# reusing this name must not be skipped as # reusing this name must not be skipped as
# already-queued. # already-queued.
queued.discard(path) queued.pop(path, None)
if not path.is_file(): if not path.is_file():
continue continue
if path in queued: if path in queued:
@@ -663,12 +709,17 @@ class Command(BaseCommand):
# rescan does not re-queue them while the consume task # rescan does not re-queue them while the consume task
# has yet to remove them from disk, but does retry a # has yet to remove them from disk, but does retry a
# failed publish instead of stranding it # failed publish instead of stranding it
if _consume_file( task_id = _consume_file(
filepath=stable_path, filepath=stable_path,
consumption_dir=directory, consumption_dir=directory,
subdirs_as_tags=subdirs_as_tags, 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 # Exit watch loop to reconfigure timeout
break break
@@ -677,13 +728,16 @@ class Command(BaseCommand):
if rescan_timeout_ms > 0 and ( if rescan_timeout_ms > 0 and (
monotonic() - last_rescan >= rescan_interval_s monotonic() - last_rescan >= rescan_interval_s
): ):
self._rescan_existing_files( try:
directory=directory, self._rescan_existing_files(
recursive=recursive, directory=directory,
consumer_filter=consumer_filter, recursive=recursive,
tracker=tracker, consumer_filter=consumer_filter,
queued=queued, tracker=tracker,
) queued=queued,
)
except Exception:
logger.exception("Error during consume folder rescan")
last_rescan = monotonic() last_rescan = monotonic()
# Determine next timeout # Determine next timeout
@@ -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"
)
+182 -16
View File
@@ -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 Command
from documents.management.commands.document_consumer import ConsumerFilter from documents.management.commands.document_consumer import ConsumerFilter
from documents.management.commands.document_consumer import FileStabilityTracker 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 TrackedFile
from documents.management.commands.document_consumer import _consume_file from documents.management.commands.document_consumer import _consume_file
from documents.management.commands.document_consumer import _tags_from_path from documents.management.commands.document_consumer import _tags_from_path
from documents.models import PaperlessTask
from documents.models import Tag from documents.models import Tag
if TYPE_CHECKING: if TYPE_CHECKING:
@@ -445,13 +447,14 @@ class TestConsumeFile:
target = consumption_dir / "document.pdf" target = consumption_dir / "document.pdf"
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
mock_consume_file_delay.apply_async.return_value.id = "abc123"
result = _consume_file( result = _consume_file(
filepath=target, filepath=target,
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, 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() mock_consume_file_delay.apply_async.assert_called_once()
call_args = mock_consume_file_delay.apply_async.call_args call_args = mock_consume_file_delay.apply_async.call_args
consumable_doc = call_args.kwargs["kwargs"]["input_doc"] consumable_doc = call_args.kwargs["kwargs"]["input_doc"]
@@ -470,7 +473,7 @@ class TestConsumeFile:
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, subdirs_as_tags=False,
) )
assert result is False assert result is None
mock_consume_file_delay.apply_async.assert_not_called() mock_consume_file_delay.apply_async.assert_not_called()
def test_consume_directory( def test_consume_directory(
@@ -487,7 +490,7 @@ class TestConsumeFile:
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, subdirs_as_tags=False,
) )
assert result is False assert result is None
mock_consume_file_delay.apply_async.assert_not_called() mock_consume_file_delay.apply_async.assert_not_called()
def test_consume_with_permission_error( def test_consume_with_permission_error(
@@ -507,7 +510,7 @@ class TestConsumeFile:
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, subdirs_as_tags=False,
) )
assert result is False assert result is None
mock_consume_file_delay.apply_async.assert_not_called() mock_consume_file_delay.apply_async.assert_not_called()
def test_consume_with_apply_async_failure( def test_consume_with_apply_async_failure(
@@ -527,7 +530,7 @@ class TestConsumeFile:
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, subdirs_as_tags=False,
) )
assert result is False assert result is None
def test_consume_with_tags_error( def test_consume_with_tags_error(
self, self,
@@ -545,12 +548,13 @@ class TestConsumeFile:
side_effect=DatabaseError("Something happened"), side_effect=DatabaseError("Something happened"),
) )
mock_consume_file_delay.apply_async.return_value.id = "abc123"
result = _consume_file( result = _consume_file(
filepath=target, filepath=target,
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=True, subdirs_as_tags=True,
) )
assert result is True assert result == "abc123"
mock_consume_file_delay.apply_async.assert_called_once() mock_consume_file_delay.apply_async.assert_called_once()
call_args = mock_consume_file_delay.apply_async.call_args call_args = mock_consume_file_delay.apply_async.call_args
overrides = call_args.kwargs["kwargs"]["overrides"] overrides = call_args.kwargs["kwargs"]["overrides"]
@@ -1116,6 +1120,7 @@ class TestCommandWatchEdgeCases:
Tag.objects.all().delete() Tag.objects.all().delete()
@pytest.mark.django_db
class TestRescanExistingFiles: class TestRescanExistingFiles:
""" """
Unit tests for the rescan safety net. Unit tests for the rescan safety net.
@@ -1134,12 +1139,23 @@ class TestRescanExistingFiles:
ignore_patterns=[], 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( def _rescan(
self, self,
directory: Path, directory: Path,
consumer_filter: ConsumerFilter, consumer_filter: ConsumerFilter,
tracker: FileStabilityTracker, tracker: FileStabilityTracker,
queued: set[Path], queued: dict[Path, QueuedFile],
*, *,
recursive: bool = False, recursive: bool = False,
) -> None: ) -> None:
@@ -1162,7 +1178,7 @@ class TestRescanExistingFiles:
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
tracker = FileStabilityTracker(stability_delay=0.1) 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.is_tracking(target) is True
assert tracker.pending_count == 1 assert tracker.pending_count == 1
@@ -1179,7 +1195,7 @@ class TestRescanExistingFiles:
tracker = FileStabilityTracker(stability_delay=0.1) tracker = FileStabilityTracker(stability_delay=0.1)
tracker.track(target, Change.added) 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 assert tracker.pending_count == 1
@@ -1193,11 +1209,17 @@ class TestRescanExistingFiles:
target = consumption_dir / "inflight.pdf" target = consumption_dir / "inflight.pdf"
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
tracker = FileStabilityTracker(stability_delay=0.1) 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) self._rescan(consumption_dir, pdf_only_filter, tracker, queued)
assert tracker.pending_count == 0 assert tracker.pending_count == 0
assert target.resolve() in queued
def test_prunes_vanished_queued_paths( def test_prunes_vanished_queued_paths(
self, self,
@@ -1207,7 +1229,7 @@ class TestRescanExistingFiles:
"""Queued paths no longer on disk are dropped so the name can recur.""" """Queued paths no longer on disk are dropped so the name can recur."""
gone = (consumption_dir / "gone.pdf").resolve() gone = (consumption_dir / "gone.pdf").resolve()
tracker = FileStabilityTracker(stability_delay=0.1) 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) self._rescan(consumption_dir, pdf_only_filter, tracker, queued)
@@ -1222,7 +1244,7 @@ class TestRescanExistingFiles:
(consumption_dir / "notes.xyz").write_bytes(b"content") (consumption_dir / "notes.xyz").write_bytes(b"content")
tracker = FileStabilityTracker(stability_delay=0.1) 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 assert tracker.pending_count == 0
@@ -1239,13 +1261,151 @@ class TestRescanExistingFiles:
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
shallow = FileStabilityTracker(stability_delay=0.1) 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 assert shallow.pending_count == 0
deep = FileStabilityTracker(stability_delay=0.1) 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 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: class TestProcessExistingFilesQueued:
"""Tests that startup processing reports which paths it queued.""" """Tests that startup processing reports which paths it queued."""
@@ -1258,7 +1418,8 @@ class TestProcessExistingFilesQueued:
mock_consume_file_delay: MagicMock, mock_consume_file_delay: MagicMock,
settings: Settings, settings: Settings,
) -> None: ) -> 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" target = consumption_dir / "document.pdf"
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
settings.CONSUMER_IGNORE_PATTERNS = [] settings.CONSUMER_IGNORE_PATTERNS = []
@@ -1271,6 +1432,9 @@ class TestProcessExistingFilesQueued:
) )
assert target.resolve() in queued 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 @pytest.mark.management
@@ -1295,11 +1459,13 @@ class TestCommandRetryAfterQueueFailure:
"""A publish failure from the watch loop is retried by the rescan.""" """A publish failure from the watch loop is retried by the rescan."""
apply_async = mock_consume_file_delay.apply_async 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: if apply_async.call_count == 1:
raise Exception("broker down") raise Exception("broker down")
return apply_async.return_value
apply_async.side_effect = fail_first_call 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) thread = start_consumer(stability_delay=0.1, rescan_interval=0.3)