mirror of
https://github.com/paperless-ngx/paperless-ngx.git
synced 2026-09-09 03:07:59 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cf252b144c |
@@ -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
|
||||
|
||||
@@ -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"
|
||||
)
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -1,9 +1,7 @@
|
||||
import logging
|
||||
|
||||
from celery import Task
|
||||
from celery import shared_task
|
||||
|
||||
from documents.models import PaperlessTask
|
||||
from paperless_mail.mail import MailAccountHandler
|
||||
from paperless_mail.mail import MailError
|
||||
from paperless_mail.models import MailAccount
|
||||
@@ -12,27 +10,8 @@ from paperless_mail.models import MailRule
|
||||
logger = logging.getLogger("paperless.mail.tasks")
|
||||
|
||||
|
||||
@shared_task(bind=True)
|
||||
def process_mail_accounts(self: Task, account_ids: list[int] | None = None) -> str:
|
||||
# A scheduled check can still be running (or queued) when the next one
|
||||
# fires, e.g. a large attachment batch that takes longer to process than
|
||||
# the check interval. ProcessedMail dedup only records a message once its
|
||||
# handling has finished, so an overlapping run can still pick up the same
|
||||
# not-yet-recorded message. Skip outright rather than race it.
|
||||
other_mail_fetch_running = (
|
||||
PaperlessTask.objects.filter(
|
||||
task_type=PaperlessTask.TaskType.MAIL_FETCH,
|
||||
status__in=[PaperlessTask.Status.PENDING, PaperlessTask.Status.STARTED],
|
||||
)
|
||||
.exclude(task_id=self.request.id)
|
||||
.exists()
|
||||
)
|
||||
if other_mail_fetch_running:
|
||||
logger.info(
|
||||
"Mail account processing is already running; skipping this run.",
|
||||
)
|
||||
return "Skipped: mail account processing already in progress."
|
||||
|
||||
@shared_task
|
||||
def process_mail_accounts(account_ids: list[int] | None = None) -> str:
|
||||
total_new_documents = 0
|
||||
accounts = (
|
||||
MailAccount.objects.filter(pk__in=account_ids)
|
||||
|
||||
@@ -1,89 +0,0 @@
|
||||
from unittest import mock
|
||||
|
||||
import pytest
|
||||
|
||||
from documents.models import PaperlessTask
|
||||
from paperless_mail import tasks
|
||||
from paperless_mail.tests.factories import MailAccountFactory
|
||||
from paperless_mail.tests.factories import MailRuleFactory
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
class TestProcessMailAccountsOverlap:
|
||||
def test_skips_when_another_mail_fetch_task_is_running(self) -> None:
|
||||
account = MailAccountFactory.create()
|
||||
MailRuleFactory.create(account=account, enabled=True)
|
||||
|
||||
PaperlessTask.objects.create(
|
||||
task_id="other-running-task",
|
||||
task_type=PaperlessTask.TaskType.MAIL_FETCH,
|
||||
trigger_source=PaperlessTask.TriggerSource.SCHEDULED,
|
||||
status=PaperlessTask.Status.STARTED,
|
||||
)
|
||||
|
||||
with mock.patch.object(
|
||||
tasks.MailAccountHandler,
|
||||
"handle_mail_account",
|
||||
) as mocked_handle:
|
||||
result = tasks.process_mail_accounts()
|
||||
|
||||
mocked_handle.assert_not_called()
|
||||
assert result == "Skipped: mail account processing already in progress."
|
||||
|
||||
def test_runs_when_no_other_mail_fetch_task_is_running(self) -> None:
|
||||
account = MailAccountFactory.create()
|
||||
MailRuleFactory.create(account=account, enabled=True)
|
||||
|
||||
with mock.patch.object(
|
||||
tasks.MailAccountHandler,
|
||||
"handle_mail_account",
|
||||
return_value=0,
|
||||
) as mocked_handle:
|
||||
result = tasks.process_mail_accounts()
|
||||
|
||||
mocked_handle.assert_called_once()
|
||||
assert result == "No new documents were added."
|
||||
|
||||
def test_ignores_completed_mail_fetch_tasks(self) -> None:
|
||||
account = MailAccountFactory.create()
|
||||
MailRuleFactory.create(account=account, enabled=True)
|
||||
|
||||
PaperlessTask.objects.create(
|
||||
task_id="finished-task",
|
||||
task_type=PaperlessTask.TaskType.MAIL_FETCH,
|
||||
trigger_source=PaperlessTask.TriggerSource.SCHEDULED,
|
||||
status=PaperlessTask.Status.SUCCESS,
|
||||
)
|
||||
|
||||
with mock.patch.object(
|
||||
tasks.MailAccountHandler,
|
||||
"handle_mail_account",
|
||||
return_value=0,
|
||||
) as mocked_handle:
|
||||
result = tasks.process_mail_accounts()
|
||||
|
||||
mocked_handle.assert_called_once()
|
||||
assert result == "No new documents were added."
|
||||
|
||||
def test_does_not_skip_due_to_its_own_task_row(self) -> None:
|
||||
account = MailAccountFactory.create()
|
||||
MailRuleFactory.create(account=account, enabled=True)
|
||||
|
||||
PaperlessTask.objects.create(
|
||||
task_id="self-task-id",
|
||||
task_type=PaperlessTask.TaskType.MAIL_FETCH,
|
||||
trigger_source=PaperlessTask.TriggerSource.SCHEDULED,
|
||||
status=PaperlessTask.Status.STARTED,
|
||||
)
|
||||
|
||||
with mock.patch.object(
|
||||
tasks.MailAccountHandler,
|
||||
"handle_mail_account",
|
||||
return_value=0,
|
||||
) as mocked_handle:
|
||||
result = tasks.process_mail_accounts.apply(
|
||||
task_id="self-task-id",
|
||||
).result
|
||||
|
||||
mocked_handle.assert_called_once()
|
||||
assert result == "No new documents were added."
|
||||
Reference in New Issue
Block a user