Compare commits

..
3 Commits
Author SHA1 Message Date
shamoon 47562449f4 Fix: ensure apply ai suggestions runs after document created 2026-09-02 08:58:45 -07:00
Trenton H c2a9532b8f Fix: Handle Celery enqueue failures when enqueuing files for consumption (#13935) 2026-09-02 15:58:27 +00:00
Trenton H 713c857a08 Fix: Handle Celery mail task chord errors (#13936)
* Fix: mail rule loops forever when all attachments are duplicates

When every attachment in a mail is rejected as a duplicate, the chord's
header tasks all fail. Celery's default task_allow_error_cb_on_chord_header
skips the error callback in that case, so no ProcessedMail row is ever
created, and the same mail is refetched and reprocessed on every poll for
as long as it stays in the rule's maximum_age window.

* Minor simplifications and cleanup
2026-09-02 08:41:30 -07:00
101 changed files with 23785 additions and 23650 deletions
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
@@ -314,7 +314,7 @@ def _consume_file(
consumption_dir: Path, consumption_dir: Path,
*, *,
subdirs_as_tags: bool, subdirs_as_tags: bool,
) -> None: ) -> bool:
""" """
Queue a file for consumption. Queue a file for consumption.
@@ -322,15 +322,20 @@ def _consume_file(
filepath: Path to the file to consume. filepath: Path to the file to consume.
consumption_dir: Base consumption directory. consumption_dir: Base consumption directory.
subdirs_as_tags: Whether to create tags from subdirectory names. 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.
""" """
# 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 return False
except OSError as e: except OSError as e:
logger.warning(f"Not consuming {filepath}: {e}") logger.warning(f"Not consuming {filepath}: {e}")
return return False
# Get tags from path if configured # Get tags from path if configured
tag_ids: list[int] | None = None tag_ids: list[int] | None = None
@@ -355,6 +360,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 True
class Command(BaseCommand): class Command(BaseCommand):
@@ -492,12 +500,12 @@ class Command(BaseCommand):
if not consumer_filter(Change.added, str(filepath)): if not consumer_filter(Change.added, str(filepath)):
continue continue
_consume_file( if _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()) queued.add(filepath.resolve())
return queued return queued
@@ -651,14 +659,16 @@ class Command(BaseCommand):
# Check for stable files # Check for stable files
for stable_path in tracker.get_stable_files(): for stable_path in tracker.get_stable_files():
_consume_file( # Only remember files that were actually queued, so the
# 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(
filepath=stable_path, filepath=stable_path,
consumption_dir=directory, consumption_dir=directory,
subdirs_as_tags=subdirs_as_tags, subdirs_as_tags=subdirs_as_tags,
) ):
# Remember it so the rescan does not re-queue it while queued.add(stable_path)
# the consume task has yet to remove it from disk
queued.add(stable_path)
# Exit watch loop to reconfigure timeout # Exit watch loop to reconfigure timeout
break break
+1 -1
View File
@@ -1003,7 +1003,7 @@ def run_workflows(
# kwargs so the PaperlessTask record can note the # kwargs so the PaperlessTask record can note the
# document, see _extract_input_data # document, see _extract_input_data
apply_ai_suggestions.delay( apply_ai_suggestions.delay_on_commit(
action_id=action.pk, action_id=action.pk,
document_id=document.pk, document_id=document.pk,
) )
@@ -445,12 +445,13 @@ class TestConsumeFile:
target = consumption_dir / "document.pdf" target = consumption_dir / "document.pdf"
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
_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
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"]
@@ -464,11 +465,12 @@ class TestConsumeFile:
mock_consume_file_delay: MagicMock, mock_consume_file_delay: MagicMock,
) -> None: ) -> None:
"""Test _consume_file handles nonexistent files gracefully.""" """Test _consume_file handles nonexistent files gracefully."""
_consume_file( result = _consume_file(
filepath=consumption_dir / "nonexistent.pdf", filepath=consumption_dir / "nonexistent.pdf",
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, subdirs_as_tags=False,
) )
assert result is False
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(
@@ -480,11 +482,12 @@ class TestConsumeFile:
subdir = consumption_dir / "subdir" subdir = consumption_dir / "subdir"
subdir.mkdir() subdir.mkdir()
_consume_file( result = _consume_file(
filepath=subdir, filepath=subdir,
consumption_dir=consumption_dir, consumption_dir=consumption_dir,
subdirs_as_tags=False, subdirs_as_tags=False,
) )
assert result is False
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(
@@ -499,13 +502,33 @@ class TestConsumeFile:
shutil.copy(sample_pdf, target) shutil.copy(sample_pdf, target)
mocker.patch.object(Path, "is_file", side_effect=PermissionError("denied")) mocker.patch.object(Path, "is_file", side_effect=PermissionError("denied"))
_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 False
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(
self,
consumption_dir: Path,
sample_pdf: Path,
mock_consume_file_delay: MagicMock,
) -> None:
"""Test _consume_file reports failure when apply_async raises."""
target = consumption_dir / "document.pdf"
shutil.copy(sample_pdf, target)
mock_consume_file_delay.apply_async.side_effect = Exception("broker down")
result = _consume_file(
filepath=target,
consumption_dir=consumption_dir,
subdirs_as_tags=False,
)
assert result is False
def test_consume_with_tags_error( def test_consume_with_tags_error(
self, self,
consumption_dir: Path, consumption_dir: Path,
@@ -522,11 +545,12 @@ class TestConsumeFile:
side_effect=DatabaseError("Something happened"), side_effect=DatabaseError("Something happened"),
) )
_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
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"]
@@ -1249,6 +1273,52 @@ class TestProcessExistingFilesQueued:
assert target.resolve() in queued assert target.resolve() in queued
@pytest.mark.management
@pytest.mark.django_db
class TestCommandRetryAfterQueueFailure:
"""
Regression test for GH #13923.
A file whose ``apply_async`` publish fails (e.g. broker briefly down)
must not be marked as queued, so the periodic rescan retries it once
the broker recovers, instead of stranding it until the consumer
process is restarted.
"""
def test_watch_loop_retries_failed_publish_on_rescan(
self,
consumption_dir: Path,
sample_pdf: Path,
mock_consume_file_delay: MagicMock,
start_consumer: Callable[..., ConsumerThread],
) -> None:
"""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:
if apply_async.call_count == 1:
raise Exception("broker down")
apply_async.side_effect = fail_first_call
thread = start_consumer(stability_delay=0.1, rescan_interval=0.3)
target = consumption_dir / "document.pdf"
shutil.copy(sample_pdf, target)
deadline = monotonic() + 5.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 failed publish to be retried by the rescan, "
f"but apply_async was only called {apply_async.call_count} time(s)"
)
@pytest.mark.management @pytest.mark.management
@pytest.mark.django_db @pytest.mark.django_db
class TestCommandRescanRecovery: class TestCommandRescanRecovery:
+5 -1
View File
@@ -5621,11 +5621,15 @@ class TestApplyAISuggestionsWorkflowAction(
action = self.make_action() action = self.make_action()
self.make_workflow(action, WorkflowTrigger.WorkflowTriggerType.DOCUMENT_ADDED) self.make_workflow(action, WorkflowTrigger.WorkflowTriggerType.DOCUMENT_ADDED)
with mock.patch("documents.tasks.apply_ai_suggestions.delay") as delay: with (
mock.patch("documents.tasks.apply_ai_suggestions.delay") as delay,
self.captureOnCommitCallbacks(execute=True),
):
run_workflows( run_workflows(
WorkflowTrigger.WorkflowTriggerType.DOCUMENT_ADDED, WorkflowTrigger.WorkflowTriggerType.DOCUMENT_ADDED,
self.doc, self.doc,
) )
delay.assert_not_called()
delay.assert_called_once_with(action_id=action.pk, document_id=self.doc.pk) delay.assert_called_once_with(action_id=action.pk, document_id=self.doc.pk)
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
+6
View File
@@ -705,6 +705,12 @@ CELERY_BROKER_TRANSPORT_OPTIONS = {
CELERY_TASK_TRACK_STARTED = True CELERY_TASK_TRACK_STARTED = True
CELERY_TASK_TIME_LIMIT: Final[int] = get_int_from_env("PAPERLESS_WORKER_TIMEOUT", 1800) CELERY_TASK_TIME_LIMIT: Final[int] = get_int_from_env("PAPERLESS_WORKER_TIMEOUT", 1800)
# https://docs.celeryq.dev/en/stable/userguide/configuration.html#std-setting-task_allow_error_cb_on_chord_header
# Without this, a failing chord header never triggers the errback, so a mail
# whose attachments all fail is never recorded and is re-fetched forever.
# The errback runs once per failed header task, so it must be idempotent.
CELERY_TASK_ALLOW_ERROR_CB_ON_CHORD_HEADER = True
CELERY_CACHE_BACKEND = "default" CELERY_CACHE_BACKEND = "default"
# https://docs.celeryq.dev/en/stable/userguide/configuration.html#task-serializer # https://docs.celeryq.dev/en/stable/userguide/configuration.html#task-serializer
+11 -5
View File
@@ -334,18 +334,24 @@ def error_callback(
""" """
A shared task that is called whenever something goes wrong during A shared task that is called whenever something goes wrong during
consumption of a file. See queue_consumption_tasks. consumption of a file. See queue_consumption_tasks.
With CELERY_TASK_ALLOW_ERROR_CB_ON_CHORD_HEADER enabled this runs once per
failed header task, not once per chord, so it must be idempotent.
""" """
rule = MailRule.objects.get(pk=rule_id) rule = MailRule.objects.get(pk=rule_id)
received = make_aware(message_date) if is_naive(message_date) else message_date
ProcessedMail.objects.create( ProcessedMail.objects.get_or_create(
rule=rule, rule=rule,
folder=rule.folder, folder=rule.folder,
uid=message_uid, uid=message_uid,
uid_validity=uid_validity, uid_validity=uid_validity,
subject=message_subject, defaults={
received=make_aware(message_date) if is_naive(message_date) else message_date, "subject": message_subject,
status="FAILED", "received": received,
error=traceback.format_exc(), "status": "FAILED",
"error": traceback.format_exc(),
},
) )

Some files were not shown because too many files have changed in this diff Show More