mirror of
https://github.com/paperless-ngx/paperless-ngx.git
synced 2026-09-09 03:07:59 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cf252b144c |
@@ -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"
|
||||||
|
)
|
||||||
@@ -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)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user