From c40922440b1ac941f161d18ef57174c7f528f240 Mon Sep 17 00:00:00 2001 From: Trenton H <797416+stumpylog@users.noreply.github.com> Date: Wed, 9 Sep 2026 07:49:49 -0700 Subject: [PATCH] Fix: prevent overlapping mail-account processing runs (#14046) process_mail_accounts had no guard against a scheduled run still being in progress when the next one fires. Skip a run outright if another MAIL_FETCH task is already PENDING/STARTED. --- src/paperless_mail/tasks.py | 24 +++- .../test_process_mail_accounts_overlap.py | 134 ++++++++++++++++++ 2 files changed, 156 insertions(+), 2 deletions(-) create mode 100644 src/paperless_mail/tests/test_process_mail_accounts_overlap.py diff --git a/src/paperless_mail/tasks.py b/src/paperless_mail/tasks.py index df1f30d91..fce516e32 100644 --- a/src/paperless_mail/tasks.py +++ b/src/paperless_mail/tasks.py @@ -1,7 +1,9 @@ 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 @@ -10,8 +12,26 @@ from paperless_mail.models import MailRule logger = logging.getLogger("paperless.mail.tasks") -@shared_task -def process_mail_accounts(account_ids: list[int] | None = None) -> str: +@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 + # 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." + total_new_documents = 0 accounts = ( MailAccount.objects.filter(pk__in=account_ids) diff --git a/src/paperless_mail/tests/test_process_mail_accounts_overlap.py b/src/paperless_mail/tests/test_process_mail_accounts_overlap.py new file mode 100644 index 000000000..4ee60ed09 --- /dev/null +++ b/src/paperless_mail/tests/test_process_mail_accounts_overlap.py @@ -0,0 +1,134 @@ +from typing import Final + +import pytest +import pytest_mock + +from documents.models import PaperlessTask +from documents.tests.factories import PaperlessTaskFactory +from paperless_mail import tasks +from paperless_mail.tests.factories import MailAccountFactory +from paperless_mail.tests.factories import MailRuleFactory + +NO_DOCUMENTS_ADDED: Final = "No new documents were added." +SKIPPED: Final = "Skipped: mail account processing already in progress." + + +@pytest.mark.django_db +@pytest.mark.usefixtures("account_with_rule") +class TestProcessMailAccountsOverlap: + @pytest.fixture + def account_with_rule(self) -> None: + """An enabled mail account with a single enabled rule.""" + account = MailAccountFactory.create() + MailRuleFactory.create(account=account, enabled=True) + + @pytest.mark.parametrize( + ("status", "expected_result", "expected_call_count"), + [ + pytest.param( + PaperlessTask.Status.PENDING, + SKIPPED, + 0, + id="pending-task-blocks", + ), + pytest.param( + PaperlessTask.Status.STARTED, + SKIPPED, + 0, + id="started-task-blocks", + ), + pytest.param( + PaperlessTask.Status.SUCCESS, + NO_DOCUMENTS_ADDED, + 1, + id="finished-task-does-not-block", + ), + ], + ) + def test_skips_only_while_another_mail_fetch_task_runs( + self, + mocker: pytest_mock.MockerFixture, + status: PaperlessTask.Status, + expected_result: str, + expected_call_count: int, + ) -> None: + """ + GIVEN: + - An enabled mail account with a rule + - Another mail fetch task row in the given status + WHEN: + - Mail accounts are processed + THEN: + - Processing is skipped only if that other task is pending or running + """ + PaperlessTaskFactory.create( + task_type=PaperlessTask.TaskType.MAIL_FETCH, + trigger_source=PaperlessTask.TriggerSource.SCHEDULED, + status=status, + ) + + mocked_handle = mocker.patch.object( + tasks.MailAccountHandler, + "handle_mail_account", + return_value=0, + ) + + result = tasks.process_mail_accounts() + + assert mocked_handle.call_count == expected_call_count + assert result == expected_result + + def test_runs_when_no_other_mail_fetch_task_exists( + self, + mocker: pytest_mock.MockerFixture, + ) -> None: + """ + GIVEN: + - An enabled mail account with a rule + - No other mail fetch task rows + WHEN: + - Mail accounts are processed + THEN: + - The account is handled + """ + mocked_handle = mocker.patch.object( + tasks.MailAccountHandler, + "handle_mail_account", + return_value=0, + ) + + result = tasks.process_mail_accounts() + + mocked_handle.assert_called_once() + assert result == NO_DOCUMENTS_ADDED + + def test_does_not_skip_due_to_its_own_task_row( + self, + mocker: pytest_mock.MockerFixture, + ) -> None: + """ + GIVEN: + - An enabled mail account with a rule + - A running mail fetch task row belonging to this very task + WHEN: + - Mail accounts are processed under that task id + THEN: + - The task does not skip itself and handles the account + """ + PaperlessTaskFactory.create( + task_id="self-task-id", + task_type=PaperlessTask.TaskType.MAIL_FETCH, + trigger_source=PaperlessTask.TriggerSource.SCHEDULED, + status=PaperlessTask.Status.STARTED, + ) + + mocked_handle = mocker.patch.object( + tasks.MailAccountHandler, + "handle_mail_account", + return_value=0, + ) + + result = tasks.process_mail_accounts.apply(task_id="self-task-id").result + + mocked_handle.assert_called_once() + assert result == NO_DOCUMENTS_ADDED