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