Coverage for paperless_mail/mail.py: 22%
503 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 09:07 +0000
1import datetime
2import imaplib
3import itertools
4import logging
5import socket
6import ssl
7import tempfile
8import traceback
9import unicodedata
10from datetime import date
11from datetime import timedelta
12from fnmatch import fnmatch
13from pathlib import Path
14from typing import TYPE_CHECKING
16import magic
17import pathvalidate
18from celery import chord
19from celery import shared_task
20from celery.canvas import Signature
21from django.conf import settings
22from django.db import DatabaseError
23from django.db.models import Q
24from django.utils import timezone
25from django.utils.timezone import is_naive
26from django.utils.timezone import make_aware
27from imap_tools import AND
28from imap_tools import NOT
29from imap_tools import MailAttachment
30from imap_tools import MailBox
31from imap_tools import MailboxFolderSelectError
32from imap_tools import MailBoxUnencrypted
33from imap_tools import MailMessage
34from imap_tools import MailMessageFlags
35from imap_tools import errors
36from imap_tools.mailbox import MailBoxStartTls
37from imap_tools.query import LogicOperator
38from imap_tools.utils import check_command_status
40from documents.data_models import ConsumableDocument
41from documents.data_models import DocumentMetadataOverrides
42from documents.data_models import DocumentSource
43from documents.loggers import LoggingMixin
44from documents.models import Correspondent
45from documents.models import PaperlessTask
46from documents.parsers import is_mime_type_supported
47from documents.tasks import consume_file
48from paperless.network import is_public_ip
49from paperless.network import resolve_hostname_ips
50from paperless_mail.models import MailAccount
51from paperless_mail.models import MailRule
52from paperless_mail.models import ProcessedMail
53from paperless_mail.oauth import PaperlessMailOAuth2Manager
54from paperless_mail.preprocessor import MailMessageDecryptor
55from paperless_mail.preprocessor import MailMessagePreprocessor
57# Apple Mail sets multiple IMAP KEYWORD and the general "\Flagged" FLAG
58# imaplib => conn.fetch(b"<message_id>", "FLAGS")
60# no flag - (FLAGS (\\Seen $NotJunk NotJunk))'
61# red - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk))'
62# orange - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk $MailFlagBit0))'
63# yellow - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk $MailFlagBit1))'
64# blue - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk $MailFlagBit2))'
65# green - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk $MailFlagBit0 $MailFlagBit1))'
66# violet - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk $MailFlagBit0 $MailFlagBit2))'
67# grey - (FLAGS (\\Flagged \\Seen $NotJunk NotJunk $MailFlagBit1 $MailFlagBit2))'
69APPLE_MAIL_TAG_COLORS = {
70 "red": [],
71 "orange": ["$MailFlagBit0"],
72 "yellow": ["$MailFlagBit1"],
73 "blue": ["$MailFlagBit2"],
74 "green": ["$MailFlagBit0", "$MailFlagBit1"],
75 "violet": ["$MailFlagBit0", "$MailFlagBit2"],
76 "grey": ["$MailFlagBit1", "$MailFlagBit2"],
77}
79MAIL_FETCH_BATCH_SIZE = 500
81# SQLite's default SQLITE_MAX_VARIABLE_NUMBER has been 32766 since 3.32.0
82# (2020), but older/custom builds and other backends may allow fewer, so
83# stay comfortably under that ceiling for `uid__in` queries against
84# ProcessedMail.
85PROCESSED_UID_QUERY_BATCH_SIZE = 10_000
88class MailError(Exception):
89 pass
92class BaseMailAction:
93 """
94 Base class for mail actions. A mail action is performed on a mail after
95 consumption of the document is complete and is used to signal to the user
96 that this mail was processed by paperless via the mail client.
98 Furthermore, mail actions reduce the amount of mails to be analyzed by
99 excluding mails on which the action was already performed (i.e., excluding
100 read mails when the action is to mark mails as read).
101 """
103 def get_criteria(self) -> dict | LogicOperator:
104 """
105 Returns filtering criteria/query for this mail action.
106 """
107 return {}
109 def post_consume(
110 self,
111 M: MailBox,
112 message_uid: str,
113 parameter: str,
114 ): # pragma: no cover
115 """
116 Perform mail action on the given mail uid in the mailbox.
117 """
118 raise NotImplementedError
121class DeleteMailAction(BaseMailAction):
122 """
123 A mail action that deletes mails after processing.
124 """
126 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None:
127 M.delete(message_uid)
130class MarkReadMailAction(BaseMailAction):
131 """
132 A mail action that marks mails as read after processing.
133 """
135 def get_criteria(self):
136 return {"seen": False}
138 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None:
139 M.flag(message_uid, [MailMessageFlags.SEEN], value=True)
142class MoveMailAction(BaseMailAction):
143 """
144 A mail action that moves mails to a different folder after processing.
145 """
147 def post_consume(self, M, message_uid, parameter) -> None:
148 M.move(message_uid, parameter)
151class FlagMailAction(BaseMailAction):
152 """
153 A mail action that marks mails as important ("star") after processing.
154 """
156 def get_criteria(self):
157 return {"flagged": False}
159 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None:
160 M.flag(message_uid, [MailMessageFlags.FLAGGED], value=True)
163class TagMailAction(BaseMailAction):
164 """
165 A mail action that tags mails after processing.
166 """
168 def __init__(self, parameter: str, *, supports_gmail_labels: bool) -> None:
169 # The custom tag should look like "apple:<color>"
170 if "apple:" in parameter.lower():
171 _, self.color = parameter.split(":")
172 self.color = self.color.strip()
174 if self.color.lower() not in APPLE_MAIL_TAG_COLORS:
175 raise MailError("Not a valid AppleMail tag color.")
177 self.keyword = None
179 else:
180 self.keyword = parameter
181 self.color = None
182 self.supports_gmail_labels = supports_gmail_labels
184 def get_criteria(self):
185 # AppleMail: We only need to check if mails are \Flagged
186 if self.color:
187 return {"flagged": False}
188 elif self.keyword:
189 if self.supports_gmail_labels:
190 return AND(NOT(gmail_label=self.keyword), no_keyword=self.keyword)
191 else:
192 return {"no_keyword": self.keyword}
193 else: # pragma: no cover
194 raise ValueError("This should never happen.")
196 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None:
197 if self.supports_gmail_labels:
198 M.client.uid("STORE", message_uid, "+X-GM-LABELS", self.keyword)
200 # AppleMail
201 elif self.color:
202 # Remove all existing $MailFlagBits
203 M.flag(
204 message_uid,
205 set(itertools.chain(*APPLE_MAIL_TAG_COLORS.values())),
206 value=False,
207 )
209 # Set new $MailFlagBits
210 M.flag(message_uid, APPLE_MAIL_TAG_COLORS.get(self.color), value=True)
212 # Set the general \Flagged
213 # This defaults to the "red" flag in AppleMail and
214 # "stars" in Thunderbird or GMail
215 M.flag(message_uid, [MailMessageFlags.FLAGGED], value=True)
217 elif self.keyword:
218 M.flag(message_uid, [self.keyword], value=True)
220 else:
221 raise MailError("No keyword specified.")
224def mailbox_login(mailbox: MailBox, account: MailAccount) -> None:
225 logger = logging.getLogger("paperless_mail")
227 try:
228 if account.is_token:
229 mailbox.xoauth2(account.username, account.password)
230 else:
231 try:
232 _ = account.password.encode("ascii")
233 use_ascii_login = True
234 except UnicodeEncodeError:
235 use_ascii_login = False
237 if use_ascii_login:
238 mailbox.login(account.username, account.password)
239 else:
240 logger.debug("Falling back to AUTH=PLAIN")
241 mailbox.login_utf8(account.username, account.password)
243 except Exception as e:
244 logger.error(
245 f"Error while authenticating account {account}: {e}",
246 exc_info=False,
247 )
248 raise MailError(
249 f"Error while authenticating account {account}",
250 ) from e
253@shared_task
254def apply_mail_action(
255 result: list,
256 rule_id: int,
257 message_uid: str,
258 message_subject: str,
259 message_date: datetime.datetime,
260 uid_validity: str | None = None,
261) -> None:
262 """
263 This shared task applies the mail action of a particular mail rule to the
264 given mail. Creates a ProcessedMail object, so that the mail won't be
265 processed in the future.
266 """
268 rule = MailRule.objects.get(pk=rule_id)
269 account = MailAccount.objects.get(pk=rule.account.pk)
271 # Ensure the date is properly timezone aware
272 if is_naive(message_date):
273 message_date = make_aware(message_date)
275 try:
276 with get_mailbox(
277 server=account.imap_server,
278 port=account.imap_port,
279 security=account.imap_security,
280 ) as M:
281 # Need to know the support for the possible tagging
282 supports_gmail_labels = "X-GM-EXT-1" in M.client.capabilities
284 mailbox_login(M, account)
285 M.folder.set(rule.folder)
287 action = get_rule_action(rule, supports_gmail_labels=supports_gmail_labels)
288 try:
289 action.post_consume(M, message_uid, rule.action_parameter)
290 except errors.ImapToolsError:
291 logger = logging.getLogger("paperless_mail")
292 logger.exception(
293 "Error while processing mail action during post_consume",
294 )
295 raise
297 ProcessedMail.objects.create(
298 owner=rule.owner,
299 rule=rule,
300 folder=rule.folder,
301 uid=message_uid,
302 uid_validity=uid_validity,
303 subject=message_subject[:256],
304 received=message_date,
305 status="SUCCESS",
306 )
308 except Exception:
309 ProcessedMail.objects.create(
310 owner=rule.owner,
311 rule=rule,
312 folder=rule.folder,
313 uid=message_uid,
314 uid_validity=uid_validity,
315 subject=message_subject[:256],
316 received=message_date,
317 status="FAILED",
318 error=traceback.format_exc(),
319 )
320 raise
323@shared_task
324def error_callback(
325 request,
326 exc,
327 tb,
328 rule_id: int,
329 message_uid: str,
330 message_subject: str,
331 message_date: datetime.datetime,
332 uid_validity: str | None = None,
333) -> None:
334 """
335 A shared task that is called whenever something goes wrong during
336 consumption of a file. See queue_consumption_tasks.
338 With CELERY_TASK_ALLOW_ERROR_CB_ON_CHORD_HEADER enabled this runs once per
339 failed header task, not once per chord, so it must be idempotent.
340 """
341 rule = MailRule.objects.get(pk=rule_id)
342 received = make_aware(message_date) if is_naive(message_date) else message_date
344 ProcessedMail.objects.get_or_create(
345 rule=rule,
346 folder=rule.folder,
347 uid=message_uid,
348 uid_validity=uid_validity,
349 defaults={
350 "owner": rule.owner,
351 "subject": message_subject[:256],
352 "received": received,
353 "status": "FAILED",
354 "error": traceback.format_exc(),
355 },
356 )
359def queue_consumption_tasks(
360 *,
361 consume_tasks: list[Signature],
362 rule: MailRule,
363 message: MailMessage,
364 uid_validity: str | None,
365) -> None:
366 """
367 Queue a list of consumption tasks (Signatures for the consume_file shared
368 task) with celery.
369 """
371 mail_action_task = apply_mail_action.s(
372 rule_id=rule.pk,
373 message_uid=message.uid,
374 message_subject=message.subject,
375 message_date=message.date,
376 uid_validity=uid_validity,
377 )
378 chord(header=consume_tasks, body=mail_action_task).on_error(
379 error_callback.s(
380 rule_id=rule.pk,
381 message_uid=message.uid,
382 message_subject=message.subject,
383 message_date=message.date,
384 uid_validity=uid_validity,
385 ),
386 ).delay()
389def get_rule_action(rule: MailRule, *, supports_gmail_labels: bool) -> BaseMailAction:
390 """
391 Returns a BaseMailAction instance for the given rule.
392 """
394 if rule.action == MailRule.MailAction.FLAG:
395 return FlagMailAction()
396 elif rule.action == MailRule.MailAction.DELETE:
397 return DeleteMailAction()
398 elif rule.action == MailRule.MailAction.MOVE:
399 return MoveMailAction()
400 elif rule.action == MailRule.MailAction.MARK_READ:
401 return MarkReadMailAction()
402 elif rule.action == MailRule.MailAction.TAG:
403 return TagMailAction(
404 rule.action_parameter,
405 supports_gmail_labels=supports_gmail_labels,
406 )
407 else:
408 raise NotImplementedError("Unknown action.") # pragma: no cover
411def make_criterias(rule: MailRule, *, supports_gmail_labels: bool):
412 """
413 Returns criteria to be applied to MailBox.fetch for the given rule.
414 """
416 maximum_age = date.today() - timedelta(days=rule.maximum_age)
417 criterias = {}
418 if rule.maximum_age > 0:
419 criterias["date_gte"] = maximum_age
420 if rule.filter_from:
421 criterias["from_"] = rule.filter_from
422 if rule.filter_to:
423 criterias["to"] = rule.filter_to
424 if rule.filter_subject:
425 criterias["subject"] = rule.filter_subject
426 if rule.filter_body:
427 criterias["body"] = rule.filter_body
429 rule_query = get_rule_action(
430 rule,
431 supports_gmail_labels=supports_gmail_labels,
432 ).get_criteria()
433 if isinstance(rule_query, dict):
434 if len(rule_query) or criterias:
435 return AND(**rule_query, **criterias)
436 else:
437 return "ALL"
438 else:
439 return AND(rule_query, **criterias)
442class PinnedIMAP4(imaplib.IMAP4):
443 """
444 IMAP4 client which connects to addresses that have already been resolved.
445 ``self.host`` keeps the original hostname for TLS SNI and cert verification
447 Without pinned addresses, and with the ssl_context of the matching imaplib
448 class, this behaves exactly like imaplib.IMAP4 / imaplib.IMAP4_SSL.
449 """
451 def __init__(self, host, port, pinned_ips, ssl_context=None, timeout=None) -> None:
452 self._pinned_ips = pinned_ips
453 self.ssl_context = ssl_context
454 super().__init__(host, port, timeout=timeout)
456 def _connect_pinned(self, timeout):
457 last_error: OSError | None = None
458 for ip_str in self._pinned_ips:
459 try:
460 address = (ip_str, self.port)
461 if timeout is not None:
462 return socket.create_connection(address, timeout)
463 return socket.create_connection(address)
464 except OSError as e:
465 last_error = e
466 raise last_error or OSError(f"Could not connect to {self.host}")
468 def _create_socket(self, timeout):
469 if self._pinned_ips: 469 ↛ 470line 469 didn't jump to line 470 because the condition on line 469 was never true
470 sock = self._connect_pinned(timeout)
471 else:
472 sock = super()._create_socket(timeout)
473 if self.ssl_context is None:
474 return sock
475 return self.ssl_context.wrap_socket(sock, server_hostname=self.host)
478class PinnedClientMixin:
479 """Builds the imaplib client against the pre-resolved addresses, if any."""
481 def __init__(self, *args, pinned_ips: list[str] | None, **kwargs) -> None:
482 self._pinned_ips = pinned_ips
483 super().__init__(*args, **kwargs)
485 def _pinned_client(self, ssl_context=None) -> imaplib.IMAP4:
486 return PinnedIMAP4(
487 self._host,
488 self._port,
489 self._pinned_ips,
490 ssl_context=ssl_context,
491 timeout=self._timeout,
492 )
495class PinnedMailBox(PinnedClientMixin, MailBox):
496 def _get_mailbox_client(self) -> imaplib.IMAP4:
497 return self._pinned_client(self._ssl_context)
500class PinnedMailBoxUnencrypted(PinnedClientMixin, MailBoxUnencrypted):
501 def _get_mailbox_client(self) -> imaplib.IMAP4:
502 return self._pinned_client()
505class PinnedMailBoxStartTls(PinnedClientMixin, MailBoxStartTls):
506 def _get_mailbox_client(self) -> imaplib.IMAP4:
507 if self._port == 993: 507 ↛ 508line 507 didn't jump to line 508 because the condition on line 507 was never true
508 raise ValueError(
509 "Port 993 requires IMAP4_SSL. Use MailBox class for SSL/TLS connection.",
510 )
511 client = self._pinned_client()
512 check_command_status(
513 client.starttls(self._ssl_context),
514 errors.MailboxStarttlsError,
515 )
516 return client
519def get_mailbox(server, port, security) -> MailBox:
520 """
521 Returns the correct MailBox instance for the given configuration.
522 """
523 pinned_ips: list[str] | None = None
524 if not settings.EMAIL_ALLOW_INTERNAL_HOSTS: 524 ↛ 525line 524 didn't jump to line 525 because the condition on line 524 was never true
525 try:
526 pinned_ips = resolve_hostname_ips(server)
527 except ValueError as e:
528 raise MailError(str(e)) from e
530 for ip_str in pinned_ips:
531 if not is_public_ip(ip_str):
532 raise MailError(
533 f"Connection blocked: {server} resolves to a non-public address",
534 )
536 ssl_context = ssl.create_default_context()
537 if settings.EMAIL_CERTIFICATE_FILE is not None: # pragma: no cover 537 ↛ 538line 537 didn't jump to line 538 because the condition on line 537 was never true
538 ssl_context.load_verify_locations(cafile=settings.EMAIL_CERTIFICATE_FILE)
540 if security == MailAccount.ImapSecurity.NONE:
541 mailbox = PinnedMailBoxUnencrypted(server, port, pinned_ips=pinned_ips)
542 elif security == MailAccount.ImapSecurity.STARTTLS:
543 mailbox = PinnedMailBoxStartTls(
544 server,
545 port,
546 ssl_context=ssl_context,
547 pinned_ips=pinned_ips,
548 )
549 elif security == MailAccount.ImapSecurity.SSL: 549 ↛ 557line 549 didn't jump to line 557 because the condition on line 549 was always true
550 mailbox = PinnedMailBox(
551 server,
552 port,
553 ssl_context=ssl_context,
554 pinned_ips=pinned_ips,
555 )
556 else:
557 raise NotImplementedError("Unknown IMAP security") # pragma: no cover
558 return mailbox
561class MailAccountHandler(LoggingMixin):
562 """
563 The main class that handles mail accounts.
565 * processes all rules for a given mail account
566 * for each mail rule, fetches relevant mails, and queues documents from
567 matching mails for consumption
568 * marks processed mails in the database, so that they won't be processed
569 again
570 * runs mail actions on the mail server, when consumption is completed
571 """
573 logging_name = "paperless_mail"
575 _message_preprocessor_types: list[type[MailMessagePreprocessor]] = [
576 MailMessageDecryptor,
577 ]
579 def __init__(self) -> None:
580 super().__init__()
581 self.renew_logging_group()
582 self._init_preprocessors()
583 self._current_uid_validity: str | None = None
585 def _init_preprocessors(self) -> None:
586 self._message_preprocessors: list[MailMessagePreprocessor] = []
587 for preprocessor_type in self._message_preprocessor_types:
588 self._init_preprocessor(preprocessor_type)
590 def _init_preprocessor(self, preprocessor_type) -> None:
591 if preprocessor_type.able_to_run():
592 try:
593 self._message_preprocessors.append(preprocessor_type())
594 except Exception as e:
595 self.log.warning(
596 f"Error while initializing preprocessor {preprocessor_type.NAME}: {e}",
597 )
598 else:
599 self.log.debug(f"Skipping mail preprocessor {preprocessor_type.NAME}")
601 def _correspondent_from_name(self, name: str) -> Correspondent | None:
602 try:
603 return Correspondent.objects.get_or_create(
604 name=name,
605 defaults={
606 "match": name,
607 "matching_algorithm": Correspondent.MATCH_LITERAL,
608 },
609 )[0]
610 except DatabaseError as e:
611 self.log.error(f"Error while retrieving correspondent {name}: {e}")
612 return None
614 def _get_title(
615 self,
616 message: MailMessage,
617 att: MailAttachment,
618 rule: MailRule,
619 ) -> str | None:
620 if rule.assign_title_from == MailRule.TitleSource.FROM_SUBJECT:
621 return unicodedata.normalize("NFC", message.subject)
623 elif rule.assign_title_from == MailRule.TitleSource.FROM_FILENAME:
624 return unicodedata.normalize("NFC", Path(att.filename).stem)
626 elif rule.assign_title_from == MailRule.TitleSource.NONE:
627 return None
629 else:
630 raise NotImplementedError(
631 "Unknown title selector.",
632 ) # pragma: no cover
634 def _get_uid_validity(self, M: MailBox, folder: str) -> str | None:
635 try:
636 uid_validity = M.folder.status(folder, ["UIDVALIDITY"]).get("UIDVALIDITY")
637 if uid_validity is not None:
638 return str(uid_validity)
639 except errors.MailboxFolderStatusError as e:
640 self.log.warning(
641 f"Server does not support retrieving UIDVALIDITY for folder {folder}: {e}",
642 )
643 except Exception as e:
644 self.log.warning(
645 f"Unable to retrieve UIDVALIDITY for folder {folder}: {e}",
646 )
647 return None
649 def _get_correspondent(
650 self,
651 message: MailMessage,
652 rule: MailRule,
653 ) -> Correspondent | None:
654 c_from = rule.assign_correspondent_from
656 if c_from == MailRule.CorrespondentSource.FROM_NOTHING:
657 return None
659 elif c_from == MailRule.CorrespondentSource.FROM_EMAIL:
660 return self._correspondent_from_name(message.from_)
662 elif c_from == MailRule.CorrespondentSource.FROM_NAME:
663 from_values = message.from_values
664 if from_values is not None and len(from_values.name) > 0:
665 return self._correspondent_from_name(from_values.name)
666 else:
667 return self._correspondent_from_name(message.from_)
669 elif c_from == MailRule.CorrespondentSource.FROM_CUSTOM:
670 return rule.assign_correspondent
672 else:
673 raise NotImplementedError(
674 "Unknown correspondent selector",
675 ) # pragma: no cover
677 def handle_mail_account(self, account: MailAccount):
678 """
679 Main entry method to handle a specific mail account.
680 """
682 self.renew_logging_group()
684 self.log.debug(f"Processing mail account {account}")
686 total_processed_files = 0
687 consumed_messages: set[tuple[str, str | None]] = set()
688 try:
689 with get_mailbox(
690 account.imap_server,
691 account.imap_port,
692 account.imap_security,
693 ) as M:
694 if (
695 account.is_token
696 and account.expiration is not None
697 and account.expiration < timezone.now()
698 ):
699 manager = PaperlessMailOAuth2Manager()
700 if manager.refresh_account_oauth_token(account):
701 account.refresh_from_db()
702 else:
703 return total_processed_files
705 supports_gmail_labels = "X-GM-EXT-1" in M.client.capabilities
706 supports_auth_plain = "AUTH=PLAIN" in M.client.capabilities
708 self.log.debug(f"GMAIL Label Support: {supports_gmail_labels}")
709 self.log.debug(f"AUTH=PLAIN Support: {supports_auth_plain}")
711 mailbox_login(M, account)
713 self.log.debug(
714 f"Account {account}: Processing {account.rules.count()} rule(s)",
715 )
717 for rule in account.rules.order_by("order"):
718 if not rule.enabled:
719 self.log.debug(f"Rule {rule}: Skipping disabled rule")
720 continue
721 try:
722 total_processed_files += self._handle_mail_rule(
723 M,
724 rule,
725 supports_gmail_labels=supports_gmail_labels,
726 consumed_messages=consumed_messages,
727 )
728 if total_processed_files > 0 and rule.stop_processing:
729 self.log.debug(
730 f"Rule {rule}: Stopping processing rules due to stop_processing flag",
731 )
732 break
733 except Exception as e:
734 self.log.exception(
735 f"Rule {rule}: Error while processing rule: {e}",
736 )
737 except MailError:
738 raise
739 except Exception as e:
740 self.log.error(
741 f"Error while retrieving mailbox {account}: {e}",
742 exc_info=False,
743 )
745 return total_processed_files
747 def _preprocess_message(self, message: MailMessage):
748 for preprocessor in self._message_preprocessors:
749 message = preprocessor.run(message)
750 return message
752 def _handle_mail_rule(
753 self,
754 M: MailBox,
755 rule: MailRule,
756 *,
757 supports_gmail_labels: bool,
758 consumed_messages: set[tuple[str, str | None]],
759 ) -> int:
760 folders = [rule.folder]
761 # In case of MOVE, make sure also the destination exists
762 if rule.action == MailRule.MailAction.MOVE:
763 folders.insert(0, rule.action_parameter)
764 try:
765 for folder in folders:
766 self.log.debug(f"Rule {rule}: Selecting folder {folder}")
767 M.folder.set(folder)
768 except MailboxFolderSelectError as err:
769 self.log.error(
770 f"Unable to access folder {folder}, attempting folder listing",
771 )
772 try:
773 for folder_info in M.folder.list():
774 self.log.info(f"Located folder: {folder_info.name}")
775 except Exception as e:
776 self.log.error(
777 "Exception during folder listing, unable to provide list folders: "
778 + str(e),
779 )
781 raise MailError(
782 f"Rule {rule}: Folder {folder} "
783 f"does not exist in account {rule.account}",
784 ) from err
786 self._current_uid_validity = self._get_uid_validity(M, rule.folder)
788 criterias = make_criterias(rule, supports_gmail_labels=supports_gmail_labels)
790 self.log.debug(
791 f"Rule {rule}: Searching folder with criteria {criterias}",
792 )
794 try:
795 all_uids = set(
796 M.uids(criteria=criterias, charset=rule.account.character_set),
797 )
798 except Exception as err:
799 raise MailError(
800 f"Rule {rule}: Error while searching folder {rule.folder}",
801 ) from err
803 all_uids_list = list(all_uids)
804 processed_uids: set[str] = set()
805 for i in range(0, len(all_uids_list), PROCESSED_UID_QUERY_BATCH_SIZE):
806 uid_chunk = all_uids_list[i : i + PROCESSED_UID_QUERY_BATCH_SIZE]
807 processed_uids_qs = ProcessedMail.objects.filter(
808 rule=rule,
809 folder=rule.folder,
810 uid__in=uid_chunk,
811 )
812 if self._current_uid_validity is not None:
813 processed_uids_qs = processed_uids_qs.filter(
814 Q(uid_validity=self._current_uid_validity)
815 | Q(uid_validity__isnull=True),
816 )
817 processed_uids.update(processed_uids_qs.values_list("uid", flat=True))
819 new_uids = all_uids - processed_uids
821 if not new_uids:
822 self.log.debug(
823 f"Rule {rule}: No new mail matching criteria {criterias}",
824 )
825 return 0
827 sorted_new_uids = sorted(new_uids, key=int)
828 try:
829 messages = M.fetch(
830 uid_list=sorted_new_uids,
831 mark_seen=False,
832 bulk=MAIL_FETCH_BATCH_SIZE,
833 )
834 except Exception as err:
835 raise MailError(
836 f"Rule {rule}: Error while fetching folder {rule.folder}",
837 ) from err
839 mails_processed = 0
840 total_processed_files = 0
841 rule_seen_messages: set[tuple[str, str | None]] = set()
843 for message in messages:
844 if TYPE_CHECKING:
845 assert isinstance(message, MailMessage)
847 message_key = (rule.folder, message.uid)
848 if message_key in rule_seen_messages:
849 self.log.debug(
850 f"Skipping duplicate fetched mail '{message.uid}' subject '{message.subject}' from '{message.from_}'.",
851 )
852 continue
853 rule_seen_messages.add(message_key)
855 if message_key in consumed_messages:
856 self.log.debug(
857 f"Skipping mail '{message.uid}' subject '{message.subject}' from '{message.from_}', already queued by a previous rule in this run.",
858 )
859 continue
861 already_processed = ProcessedMail.objects.filter(
862 rule=rule,
863 uid=message.uid,
864 folder=rule.folder,
865 )
866 if self._current_uid_validity is not None:
867 already_processed = already_processed.filter(
868 Q(uid_validity=self._current_uid_validity)
869 | Q(uid_validity__isnull=True),
870 )
871 if already_processed.exists():
872 self.log.debug(
873 f"Skipping mail '{message.uid}' subject '{message.subject}' from '{message.from_}', already processed.",
874 )
875 continue
877 try:
878 processed_files = self._handle_message(message, rule)
879 if processed_files > 0:
880 consumed_messages.add(message_key)
882 total_processed_files += processed_files
883 mails_processed += 1
884 except Exception as e:
885 self.log.exception(
886 f"Rule {rule}: Error while processing mail {message.uid}: {e}",
887 )
889 self.log.debug(f"Rule {rule}: Processed {mails_processed} matching mail(s)")
891 return total_processed_files
893 def _handle_message(self, message, rule: MailRule) -> int:
894 message = self._preprocess_message(message)
896 processed_elements = 0
898 # Skip Message handling when only attachments are to be processed but
899 # message doesn't have any.
900 if (
901 not message.attachments
902 and rule.consumption_scope == MailRule.ConsumptionScope.ATTACHMENTS_ONLY
903 ):
904 self._record_processed_without_consumption(message, rule)
905 return processed_elements
907 self.log.debug(
908 f"Rule {rule}: "
909 f"Processing mail {message.subject} from {message.from_} with "
910 f"{len(message.attachments)} attachment(s)",
911 )
913 tag_ids: list[int] = [tag.id for tag in rule.assign_tags.all()]
914 doc_type = rule.assign_document_type
916 if (
917 rule.consumption_scope == MailRule.ConsumptionScope.EML_ONLY
918 or rule.consumption_scope == MailRule.ConsumptionScope.EVERYTHING
919 ):
920 processed_elements += self._process_eml(
921 message,
922 rule,
923 tag_ids,
924 doc_type,
925 )
927 if (
928 rule.consumption_scope == MailRule.ConsumptionScope.ATTACHMENTS_ONLY
929 or rule.consumption_scope == MailRule.ConsumptionScope.EVERYTHING
930 ):
931 processed_elements += self._process_attachments(
932 message,
933 rule,
934 tag_ids,
935 doc_type,
936 )
938 return processed_elements
940 def _record_processed_without_consumption(
941 self,
942 message: MailMessage,
943 rule: MailRule,
944 ) -> None:
945 ProcessedMail.objects.get_or_create(
946 rule=rule,
947 uid=message.uid,
948 folder=rule.folder,
949 uid_validity=self._current_uid_validity,
950 defaults={
951 "owner": rule.owner,
952 "subject": message.subject[:256],
953 "received": make_aware(message.date)
954 if is_naive(message.date)
955 else message.date,
956 "status": "PROCESSED_WO_CONSUMPTION",
957 },
958 )
960 def filename_inclusion_matches(
961 self,
962 filter_attachment_filename_include: str | None,
963 filename: str,
964 ) -> bool:
965 if filter_attachment_filename_include:
966 filter_attachment_filename_inclusions = (
967 filter_attachment_filename_include.split(",")
968 )
970 # Force the filename and pattern to the lowercase
971 # as this is system dependent otherwise
972 filename = filename.lower()
973 for filename_include in filter_attachment_filename_inclusions:
974 if filename_include and fnmatch(filename, filename_include.lower()):
975 return True
976 return False
977 return True
979 def filename_exclusion_matches(
980 self,
981 filter_attachment_filename_exclude: str | None,
982 filename: str,
983 ) -> bool:
984 if filter_attachment_filename_exclude:
985 filter_attachment_filename_exclusions = (
986 filter_attachment_filename_exclude.split(",")
987 )
989 # Force the filename and pattern to the lowercase
990 # as this is system dependent otherwise
991 filename = filename.lower()
992 for filename_exclude in filter_attachment_filename_exclusions:
993 if filename_exclude and fnmatch(filename, filename_exclude.lower()):
994 return True
995 return False
997 def _process_attachments(
998 self,
999 message: MailMessage,
1000 rule: MailRule,
1001 tag_ids,
1002 doc_type,
1003 ):
1004 processed_attachments = 0
1006 consume_tasks = []
1008 for att in message.attachments:
1009 if (
1010 att.content_disposition != "attachment"
1011 and rule.attachment_type
1012 == MailRule.AttachmentProcessing.ATTACHMENTS_ONLY
1013 ):
1014 self.log.debug(
1015 f"Rule {rule}: "
1016 f"Skipping attachment {att.filename} "
1017 f"with content disposition {att.content_disposition}",
1018 )
1019 continue
1021 if not self.filename_inclusion_matches(
1022 rule.filter_attachment_filename_include,
1023 att.filename,
1024 ):
1025 # Force the filename and pattern to the lowercase
1026 # as this is system dependent otherwise
1027 self.log.debug(
1028 f"Rule {rule}: "
1029 f"Skipping attachment {att.filename} "
1030 f"does not match pattern {rule.filter_attachment_filename_include}",
1031 )
1032 continue
1033 elif self.filename_exclusion_matches(
1034 rule.filter_attachment_filename_exclude,
1035 att.filename,
1036 ):
1037 self.log.debug(
1038 f"Rule {rule}: "
1039 f"Skipping attachment {att.filename} "
1040 f"does match pattern {rule.filter_attachment_filename_exclude}",
1041 )
1042 continue
1044 correspondent = self._get_correspondent(message, rule)
1046 title = self._get_title(message, att, rule)
1048 # don't trust the content type of the attachment. Could be
1049 # generic application/octet-stream.
1050 mime_type = magic.from_buffer(att.payload, mime=True)
1052 if is_mime_type_supported(mime_type):
1053 self.log.info(
1054 f"Rule {rule}: "
1055 f"Consuming attachment {att.filename} from mail "
1056 f"{message.subject} from {message.from_}",
1057 )
1059 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True)
1061 temp_dir = Path(
1062 tempfile.mkdtemp(
1063 prefix="paperless-mail-",
1064 dir=settings.SCRATCH_DIR,
1065 ),
1066 )
1068 attachment_name = pathvalidate.sanitize_filename(
1069 unicodedata.normalize("NFC", att.filename),
1070 )
1071 if attachment_name:
1072 temp_filename = temp_dir / attachment_name
1073 else: # pragma: no cover
1074 # Some cases may have no name (generally inline)
1075 temp_filename = temp_dir / "no-name-attachment"
1077 temp_filename.write_bytes(att.payload)
1079 input_doc = ConsumableDocument(
1080 source=DocumentSource.MailFetch,
1081 original_file=temp_filename,
1082 mailrule_id=rule.pk,
1083 )
1084 doc_overrides = DocumentMetadataOverrides(
1085 title=title,
1086 filename=attachment_name,
1087 correspondent_id=correspondent.id if correspondent else None,
1088 document_type_id=doc_type.id if doc_type else None,
1089 tag_ids=tag_ids,
1090 owner_id=(
1091 rule.owner.id
1092 if (rule.assign_owner_from_rule and rule.owner)
1093 else None
1094 ),
1095 )
1097 consume_task = consume_file.s(
1098 input_doc=input_doc,
1099 overrides=doc_overrides,
1100 ).set(
1101 headers={
1102 "trigger_source": PaperlessTask.TriggerSource.EMAIL_CONSUME,
1103 },
1104 )
1106 consume_tasks.append(consume_task)
1108 processed_attachments += 1
1109 else:
1110 self.log.debug(
1111 f"Rule {rule}: "
1112 f"Skipping attachment {att.filename} "
1113 f"since guessed mime type {mime_type} is not supported "
1114 f"by paperless",
1115 )
1117 if len(consume_tasks) > 0:
1118 queue_consumption_tasks(
1119 consume_tasks=consume_tasks,
1120 rule=rule,
1121 message=message,
1122 uid_validity=self._current_uid_validity,
1123 )
1124 else:
1125 # No files to consume, just mark as processed if it wasn't by .eml processing
1126 self._record_processed_without_consumption(message, rule)
1128 return processed_attachments
1130 def _process_eml(
1131 self,
1132 message: MailMessage,
1133 rule: MailRule,
1134 tag_ids,
1135 doc_type,
1136 ):
1137 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True)
1138 _, temp_filename = tempfile.mkstemp(
1139 prefix="paperless-mail-",
1140 dir=settings.SCRATCH_DIR,
1141 suffix=".eml",
1142 )
1143 with Path(temp_filename).open("wb") as f:
1144 # Move "From"-header to beginning of file
1145 # TODO: This ugly workaround is needed because the parser is
1146 # chosen only by the mime_type detected via magic
1147 # (see documents/consumer.py "mime_type = magic.from_file")
1148 # Unfortunately magic sometimes fails to detect the mime
1149 # type of .eml files correctly as message/rfc822 and instead
1150 # detects text/plain.
1151 # This also effects direct file consumption of .eml files
1152 # which are not treated with this workaround.
1153 from_element = None
1154 for i, header in enumerate(message.obj._headers):
1155 if header[0] == "From":
1156 from_element = i
1157 if from_element:
1158 new_headers = [message.obj._headers.pop(from_element)]
1159 new_headers += message.obj._headers
1160 message.obj._headers = new_headers
1162 f.write(message.obj.as_bytes())
1164 correspondent = self._get_correspondent(message, rule)
1166 self.log.info(
1167 f"Rule {rule}: "
1168 f"Consuming eml from mail "
1169 f"{message.subject} from {message.from_}",
1170 )
1172 input_doc = ConsumableDocument(
1173 source=DocumentSource.MailFetch,
1174 original_file=temp_filename,
1175 mailrule_id=rule.pk,
1176 )
1177 doc_overrides = DocumentMetadataOverrides(
1178 title=message.subject,
1179 filename=pathvalidate.sanitize_filename(
1180 unicodedata.normalize("NFC", f"{message.subject}.eml"),
1181 ),
1182 correspondent_id=correspondent.id if correspondent else None,
1183 document_type_id=doc_type.id if doc_type else None,
1184 tag_ids=tag_ids,
1185 owner_id=rule.owner.id if rule.owner else None,
1186 )
1188 consume_task = consume_file.s(
1189 input_doc=input_doc,
1190 overrides=doc_overrides,
1191 ).set(headers={"trigger_source": PaperlessTask.TriggerSource.EMAIL_CONSUME})
1193 queue_consumption_tasks(
1194 consume_tasks=[consume_task],
1195 rule=rule,
1196 message=message,
1197 uid_validity=self._current_uid_validity,
1198 )
1200 processed_elements = 1
1201 return processed_elements