mirror of
https://github.com/paperless-ngx/paperless-ngx.git
synced 2026-09-09 11:17:58 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
82f76cdcc9 |
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -72,24 +72,6 @@ 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.
|
||||
@@ -332,7 +314,7 @@ def _consume_file(
|
||||
consumption_dir: Path,
|
||||
*,
|
||||
subdirs_as_tags: bool,
|
||||
) -> str | None:
|
||||
) -> bool:
|
||||
"""
|
||||
Queue a file for consumption.
|
||||
|
||||
@@ -342,18 +324,18 @@ def _consume_file(
|
||||
subdirs_as_tags: Whether to create tags from subdirectory names.
|
||||
|
||||
Returns:
|
||||
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.
|
||||
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.
|
||||
"""
|
||||
# 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 None
|
||||
return False
|
||||
except OSError as e:
|
||||
logger.warning(f"Not consuming {filepath}: {e}")
|
||||
return None
|
||||
return False
|
||||
|
||||
# Get tags from path if configured
|
||||
tag_ids: list[int] | None = None
|
||||
@@ -366,7 +348,7 @@ def _consume_file(
|
||||
# Queue for consumption
|
||||
try:
|
||||
logger.info(f"Adding {filepath} to the task queue")
|
||||
result = consume_file.apply_async(
|
||||
consume_file.apply_async(
|
||||
kwargs={
|
||||
"input_doc": ConsumableDocument(
|
||||
source=DocumentSource.ConsumeFolder,
|
||||
@@ -378,9 +360,9 @@ def _consume_file(
|
||||
)
|
||||
except Exception:
|
||||
logger.exception(f"Error while queuing document {filepath}")
|
||||
return None
|
||||
return False
|
||||
|
||||
return result.id
|
||||
return True
|
||||
|
||||
|
||||
class Command(BaseCommand):
|
||||
@@ -497,19 +479,18 @@ class Command(BaseCommand):
|
||||
recursive: bool,
|
||||
subdirs_as_tags: bool,
|
||||
consumer_filter: ConsumerFilter,
|
||||
) -> dict[Path, QueuedFile]:
|
||||
) -> set[Path]:
|
||||
"""
|
||||
Process any existing files in the consumption directory.
|
||||
|
||||
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.
|
||||
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.
|
||||
"""
|
||||
logger.info(f"Processing existing files in {directory}")
|
||||
|
||||
glob_pattern = "**/*" if recursive else "*"
|
||||
queued: dict[Path, QueuedFile] = {}
|
||||
queued: set[Path] = set()
|
||||
|
||||
for filepath in directory.glob(glob_pattern):
|
||||
# Use filter to check if file should be processed
|
||||
@@ -519,17 +500,12 @@ class Command(BaseCommand):
|
||||
if not consumer_filter(Change.added, str(filepath)):
|
||||
continue
|
||||
|
||||
task_id = _consume_file(
|
||||
if _consume_file(
|
||||
filepath=filepath,
|
||||
consumption_dir=directory,
|
||||
subdirs_as_tags=subdirs_as_tags,
|
||||
)
|
||||
if task_id is None:
|
||||
continue
|
||||
|
||||
entry = QueuedFile.from_path(task_id, filepath)
|
||||
if entry is not None:
|
||||
queued[filepath.resolve()] = entry
|
||||
):
|
||||
queued.add(filepath.resolve())
|
||||
|
||||
return queued
|
||||
|
||||
@@ -540,43 +516,21 @@ class Command(BaseCommand):
|
||||
recursive: bool,
|
||||
consumer_filter: ConsumerFilter,
|
||||
tracker: FileStabilityTracker,
|
||||
queued: dict[Path, QueuedFile],
|
||||
queued: set[Path],
|
||||
) -> 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 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.
|
||||
(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.
|
||||
"""
|
||||
# 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]
|
||||
# Prune in-flight paths that have left the directory
|
||||
for path in list(queued):
|
||||
if not path.exists():
|
||||
queued.discard(path)
|
||||
|
||||
glob_pattern = "**/*" if recursive else "*"
|
||||
|
||||
@@ -604,7 +558,7 @@ class Command(BaseCommand):
|
||||
polling_interval: float,
|
||||
stability_delay: float,
|
||||
is_testing: bool,
|
||||
queued: dict[Path, QueuedFile] | None = None,
|
||||
queued: set[Path] | None = None,
|
||||
) -> None:
|
||||
"""Watch directory for changes and process stable files."""
|
||||
use_polling = polling_interval > 0
|
||||
@@ -613,7 +567,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 = {} if queued is None else queued
|
||||
queued = set() if queued is None else queued
|
||||
|
||||
# Full-glob safety net cadence (0 disables)
|
||||
rescan_interval_s = self.rescan_interval_s
|
||||
@@ -690,7 +644,7 @@ class Command(BaseCommand):
|
||||
# Consumed (or otherwise removed); a later file
|
||||
# reusing this name must not be skipped as
|
||||
# already-queued.
|
||||
queued.pop(path, None)
|
||||
queued.discard(path)
|
||||
if not path.is_file():
|
||||
continue
|
||||
if path in queued:
|
||||
@@ -709,17 +663,12 @@ 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
|
||||
task_id = _consume_file(
|
||||
if _consume_file(
|
||||
filepath=stable_path,
|
||||
consumption_dir=directory,
|
||||
subdirs_as_tags=subdirs_as_tags,
|
||||
)
|
||||
if task_id is None:
|
||||
continue
|
||||
|
||||
entry = QueuedFile.from_path(task_id, stable_path)
|
||||
if entry is not None:
|
||||
queued[stable_path] = entry
|
||||
):
|
||||
queued.add(stable_path)
|
||||
|
||||
# Exit watch loop to reconfigure timeout
|
||||
break
|
||||
@@ -728,16 +677,13 @@ class Command(BaseCommand):
|
||||
if rescan_timeout_ms > 0 and (
|
||||
monotonic() - last_rescan >= rescan_interval_s
|
||||
):
|
||||
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")
|
||||
self._rescan_existing_files(
|
||||
directory=directory,
|
||||
recursive=recursive,
|
||||
consumer_filter=consumer_filter,
|
||||
tracker=tracker,
|
||||
queued=queued,
|
||||
)
|
||||
last_rescan = monotonic()
|
||||
|
||||
# Determine next timeout
|
||||
|
||||
@@ -1,165 +0,0 @@
|
||||
"""
|
||||
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,11 +33,9 @@ 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:
|
||||
@@ -447,14 +445,13 @@ 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 == mock_consume_file_delay.apply_async.return_value.id
|
||||
assert result is True
|
||||
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"]
|
||||
@@ -473,7 +470,7 @@ class TestConsumeFile:
|
||||
consumption_dir=consumption_dir,
|
||||
subdirs_as_tags=False,
|
||||
)
|
||||
assert result is None
|
||||
assert result is False
|
||||
mock_consume_file_delay.apply_async.assert_not_called()
|
||||
|
||||
def test_consume_directory(
|
||||
@@ -490,7 +487,7 @@ class TestConsumeFile:
|
||||
consumption_dir=consumption_dir,
|
||||
subdirs_as_tags=False,
|
||||
)
|
||||
assert result is None
|
||||
assert result is False
|
||||
mock_consume_file_delay.apply_async.assert_not_called()
|
||||
|
||||
def test_consume_with_permission_error(
|
||||
@@ -510,7 +507,7 @@ class TestConsumeFile:
|
||||
consumption_dir=consumption_dir,
|
||||
subdirs_as_tags=False,
|
||||
)
|
||||
assert result is None
|
||||
assert result is False
|
||||
mock_consume_file_delay.apply_async.assert_not_called()
|
||||
|
||||
def test_consume_with_apply_async_failure(
|
||||
@@ -530,7 +527,7 @@ class TestConsumeFile:
|
||||
consumption_dir=consumption_dir,
|
||||
subdirs_as_tags=False,
|
||||
)
|
||||
assert result is None
|
||||
assert result is False
|
||||
|
||||
def test_consume_with_tags_error(
|
||||
self,
|
||||
@@ -548,13 +545,12 @@ 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 == "abc123"
|
||||
assert result is True
|
||||
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"]
|
||||
@@ -1120,7 +1116,6 @@ class TestCommandWatchEdgeCases:
|
||||
Tag.objects.all().delete()
|
||||
|
||||
|
||||
@pytest.mark.django_db
|
||||
class TestRescanExistingFiles:
|
||||
"""
|
||||
Unit tests for the rescan safety net.
|
||||
@@ -1139,23 +1134,12 @@ 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: dict[Path, QueuedFile],
|
||||
queued: set[Path],
|
||||
*,
|
||||
recursive: bool = False,
|
||||
) -> None:
|
||||
@@ -1178,7 +1162,7 @@ class TestRescanExistingFiles:
|
||||
shutil.copy(sample_pdf, target)
|
||||
tracker = FileStabilityTracker(stability_delay=0.1)
|
||||
|
||||
self._rescan(consumption_dir, pdf_only_filter, tracker, {})
|
||||
self._rescan(consumption_dir, pdf_only_filter, tracker, set())
|
||||
|
||||
assert tracker.is_tracking(target) is True
|
||||
assert tracker.pending_count == 1
|
||||
@@ -1195,7 +1179,7 @@ class TestRescanExistingFiles:
|
||||
tracker = FileStabilityTracker(stability_delay=0.1)
|
||||
tracker.track(target, Change.added)
|
||||
|
||||
self._rescan(consumption_dir, pdf_only_filter, tracker, {})
|
||||
self._rescan(consumption_dir, pdf_only_filter, tracker, set())
|
||||
|
||||
assert tracker.pending_count == 1
|
||||
|
||||
@@ -1209,17 +1193,11 @@ class TestRescanExistingFiles:
|
||||
target = consumption_dir / "inflight.pdf"
|
||||
shutil.copy(sample_pdf, target)
|
||||
tracker = FileStabilityTracker(stability_delay=0.1)
|
||||
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")
|
||||
queued = {target.resolve()}
|
||||
|
||||
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,
|
||||
@@ -1229,7 +1207,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: QueuedFile(task_id="task-gone", size=0, mtime=0.0)}
|
||||
queued = {gone}
|
||||
|
||||
self._rescan(consumption_dir, pdf_only_filter, tracker, queued)
|
||||
|
||||
@@ -1244,7 +1222,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, {})
|
||||
self._rescan(consumption_dir, pdf_only_filter, tracker, set())
|
||||
|
||||
assert tracker.pending_count == 0
|
||||
|
||||
@@ -1261,151 +1239,13 @@ class TestRescanExistingFiles:
|
||||
shutil.copy(sample_pdf, target)
|
||||
|
||||
shallow = FileStabilityTracker(stability_delay=0.1)
|
||||
self._rescan(consumption_dir, pdf_only_filter, shallow, {})
|
||||
self._rescan(consumption_dir, pdf_only_filter, shallow, set())
|
||||
assert shallow.pending_count == 0
|
||||
|
||||
deep = FileStabilityTracker(stability_delay=0.1)
|
||||
self._rescan(consumption_dir, pdf_only_filter, deep, {}, recursive=True)
|
||||
self._rescan(consumption_dir, pdf_only_filter, deep, set(), 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."""
|
||||
@@ -1418,8 +1258,7 @@ class TestProcessExistingFilesQueued:
|
||||
mock_consume_file_delay: MagicMock,
|
||||
settings: Settings,
|
||||
) -> None:
|
||||
"""The dict returned seeds the rescan's queued dict, avoiding re-queue."""
|
||||
mock_consume_file_delay.apply_async.return_value.id = "startup-task-id"
|
||||
"""The set returned seeds the rescan's queued set, avoiding re-queue."""
|
||||
target = consumption_dir / "document.pdf"
|
||||
shutil.copy(sample_pdf, target)
|
||||
settings.CONSUMER_IGNORE_PATTERNS = []
|
||||
@@ -1432,9 +1271,6 @@ 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
|
||||
@@ -1459,13 +1295,11 @@ 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) -> MagicMock | None:
|
||||
def fail_first_call(*args: object, **kwargs: object) -> 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)
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user