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

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 

15 

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 

39 

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 

56 

57# Apple Mail sets multiple IMAP KEYWORD and the general "\Flagged" FLAG 

58# imaplib => conn.fetch(b"<message_id>", "FLAGS") 

59 

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))' 

68 

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} 

78 

79MAIL_FETCH_BATCH_SIZE = 500 

80 

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 

86 

87 

88class MailError(Exception): 

89 pass 

90 

91 

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. 

97 

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 """ 

102 

103 def get_criteria(self) -> dict | LogicOperator: 

104 """ 

105 Returns filtering criteria/query for this mail action. 

106 """ 

107 return {} 

108 

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 

119 

120 

121class DeleteMailAction(BaseMailAction): 

122 """ 

123 A mail action that deletes mails after processing. 

124 """ 

125 

126 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None: 

127 M.delete(message_uid) 

128 

129 

130class MarkReadMailAction(BaseMailAction): 

131 """ 

132 A mail action that marks mails as read after processing. 

133 """ 

134 

135 def get_criteria(self): 

136 return {"seen": False} 

137 

138 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None: 

139 M.flag(message_uid, [MailMessageFlags.SEEN], value=True) 

140 

141 

142class MoveMailAction(BaseMailAction): 

143 """ 

144 A mail action that moves mails to a different folder after processing. 

145 """ 

146 

147 def post_consume(self, M, message_uid, parameter) -> None: 

148 M.move(message_uid, parameter) 

149 

150 

151class FlagMailAction(BaseMailAction): 

152 """ 

153 A mail action that marks mails as important ("star") after processing. 

154 """ 

155 

156 def get_criteria(self): 

157 return {"flagged": False} 

158 

159 def post_consume(self, M: MailBox, message_uid: str, parameter: str) -> None: 

160 M.flag(message_uid, [MailMessageFlags.FLAGGED], value=True) 

161 

162 

163class TagMailAction(BaseMailAction): 

164 """ 

165 A mail action that tags mails after processing. 

166 """ 

167 

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() 

173 

174 if self.color.lower() not in APPLE_MAIL_TAG_COLORS: 

175 raise MailError("Not a valid AppleMail tag color.") 

176 

177 self.keyword = None 

178 

179 else: 

180 self.keyword = parameter 

181 self.color = None 

182 self.supports_gmail_labels = supports_gmail_labels 

183 

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.") 

195 

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) 

199 

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 ) 

208 

209 # Set new $MailFlagBits 

210 M.flag(message_uid, APPLE_MAIL_TAG_COLORS.get(self.color), value=True) 

211 

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) 

216 

217 elif self.keyword: 

218 M.flag(message_uid, [self.keyword], value=True) 

219 

220 else: 

221 raise MailError("No keyword specified.") 

222 

223 

224def mailbox_login(mailbox: MailBox, account: MailAccount) -> None: 

225 logger = logging.getLogger("paperless_mail") 

226 

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 

236 

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) 

242 

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 

251 

252 

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 """ 

267 

268 rule = MailRule.objects.get(pk=rule_id) 

269 account = MailAccount.objects.get(pk=rule.account.pk) 

270 

271 # Ensure the date is properly timezone aware 

272 if is_naive(message_date): 

273 message_date = make_aware(message_date) 

274 

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 

283 

284 mailbox_login(M, account) 

285 M.folder.set(rule.folder) 

286 

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 

296 

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 ) 

307 

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 

321 

322 

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. 

337 

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 

343 

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 ) 

357 

358 

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 """ 

370 

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() 

387 

388 

389def get_rule_action(rule: MailRule, *, supports_gmail_labels: bool) -> BaseMailAction: 

390 """ 

391 Returns a BaseMailAction instance for the given rule. 

392 """ 

393 

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 

409 

410 

411def make_criterias(rule: MailRule, *, supports_gmail_labels: bool): 

412 """ 

413 Returns criteria to be applied to MailBox.fetch for the given rule. 

414 """ 

415 

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 

428 

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) 

440 

441 

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 

446 

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 """ 

450 

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) 

455 

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}") 

467 

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) 

476 

477 

478class PinnedClientMixin: 

479 """Builds the imaplib client against the pre-resolved addresses, if any.""" 

480 

481 def __init__(self, *args, pinned_ips: list[str] | None, **kwargs) -> None: 

482 self._pinned_ips = pinned_ips 

483 super().__init__(*args, **kwargs) 

484 

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 ) 

493 

494 

495class PinnedMailBox(PinnedClientMixin, MailBox): 

496 def _get_mailbox_client(self) -> imaplib.IMAP4: 

497 return self._pinned_client(self._ssl_context) 

498 

499 

500class PinnedMailBoxUnencrypted(PinnedClientMixin, MailBoxUnencrypted): 

501 def _get_mailbox_client(self) -> imaplib.IMAP4: 

502 return self._pinned_client() 

503 

504 

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 

517 

518 

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 

529 

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 ) 

535 

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) 

539 

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 

559 

560 

561class MailAccountHandler(LoggingMixin): 

562 """ 

563 The main class that handles mail accounts. 

564 

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 """ 

572 

573 logging_name = "paperless_mail" 

574 

575 _message_preprocessor_types: list[type[MailMessagePreprocessor]] = [ 

576 MailMessageDecryptor, 

577 ] 

578 

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 

584 

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) 

589 

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}") 

600 

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 

613 

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) 

622 

623 elif rule.assign_title_from == MailRule.TitleSource.FROM_FILENAME: 

624 return unicodedata.normalize("NFC", Path(att.filename).stem) 

625 

626 elif rule.assign_title_from == MailRule.TitleSource.NONE: 

627 return None 

628 

629 else: 

630 raise NotImplementedError( 

631 "Unknown title selector.", 

632 ) # pragma: no cover 

633 

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 

648 

649 def _get_correspondent( 

650 self, 

651 message: MailMessage, 

652 rule: MailRule, 

653 ) -> Correspondent | None: 

654 c_from = rule.assign_correspondent_from 

655 

656 if c_from == MailRule.CorrespondentSource.FROM_NOTHING: 

657 return None 

658 

659 elif c_from == MailRule.CorrespondentSource.FROM_EMAIL: 

660 return self._correspondent_from_name(message.from_) 

661 

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_) 

668 

669 elif c_from == MailRule.CorrespondentSource.FROM_CUSTOM: 

670 return rule.assign_correspondent 

671 

672 else: 

673 raise NotImplementedError( 

674 "Unknown correspondent selector", 

675 ) # pragma: no cover 

676 

677 def handle_mail_account(self, account: MailAccount): 

678 """ 

679 Main entry method to handle a specific mail account. 

680 """ 

681 

682 self.renew_logging_group() 

683 

684 self.log.debug(f"Processing mail account {account}") 

685 

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 

704 

705 supports_gmail_labels = "X-GM-EXT-1" in M.client.capabilities 

706 supports_auth_plain = "AUTH=PLAIN" in M.client.capabilities 

707 

708 self.log.debug(f"GMAIL Label Support: {supports_gmail_labels}") 

709 self.log.debug(f"AUTH=PLAIN Support: {supports_auth_plain}") 

710 

711 mailbox_login(M, account) 

712 

713 self.log.debug( 

714 f"Account {account}: Processing {account.rules.count()} rule(s)", 

715 ) 

716 

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 ) 

744 

745 return total_processed_files 

746 

747 def _preprocess_message(self, message: MailMessage): 

748 for preprocessor in self._message_preprocessors: 

749 message = preprocessor.run(message) 

750 return message 

751 

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 ) 

780 

781 raise MailError( 

782 f"Rule {rule}: Folder {folder} " 

783 f"does not exist in account {rule.account}", 

784 ) from err 

785 

786 self._current_uid_validity = self._get_uid_validity(M, rule.folder) 

787 

788 criterias = make_criterias(rule, supports_gmail_labels=supports_gmail_labels) 

789 

790 self.log.debug( 

791 f"Rule {rule}: Searching folder with criteria {criterias}", 

792 ) 

793 

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 

802 

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)) 

818 

819 new_uids = all_uids - processed_uids 

820 

821 if not new_uids: 

822 self.log.debug( 

823 f"Rule {rule}: No new mail matching criteria {criterias}", 

824 ) 

825 return 0 

826 

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 

838 

839 mails_processed = 0 

840 total_processed_files = 0 

841 rule_seen_messages: set[tuple[str, str | None]] = set() 

842 

843 for message in messages: 

844 if TYPE_CHECKING: 

845 assert isinstance(message, MailMessage) 

846 

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) 

854 

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 

860 

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 

876 

877 try: 

878 processed_files = self._handle_message(message, rule) 

879 if processed_files > 0: 

880 consumed_messages.add(message_key) 

881 

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 ) 

888 

889 self.log.debug(f"Rule {rule}: Processed {mails_processed} matching mail(s)") 

890 

891 return total_processed_files 

892 

893 def _handle_message(self, message, rule: MailRule) -> int: 

894 message = self._preprocess_message(message) 

895 

896 processed_elements = 0 

897 

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 

906 

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 ) 

912 

913 tag_ids: list[int] = [tag.id for tag in rule.assign_tags.all()] 

914 doc_type = rule.assign_document_type 

915 

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 ) 

926 

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 ) 

937 

938 return processed_elements 

939 

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 ) 

959 

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 ) 

969 

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 

978 

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 ) 

988 

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 

996 

997 def _process_attachments( 

998 self, 

999 message: MailMessage, 

1000 rule: MailRule, 

1001 tag_ids, 

1002 doc_type, 

1003 ): 

1004 processed_attachments = 0 

1005 

1006 consume_tasks = [] 

1007 

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 

1020 

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 

1043 

1044 correspondent = self._get_correspondent(message, rule) 

1045 

1046 title = self._get_title(message, att, rule) 

1047 

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) 

1051 

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 ) 

1058 

1059 settings.SCRATCH_DIR.mkdir(parents=True, exist_ok=True) 

1060 

1061 temp_dir = Path( 

1062 tempfile.mkdtemp( 

1063 prefix="paperless-mail-", 

1064 dir=settings.SCRATCH_DIR, 

1065 ), 

1066 ) 

1067 

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" 

1076 

1077 temp_filename.write_bytes(att.payload) 

1078 

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 ) 

1096 

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 ) 

1105 

1106 consume_tasks.append(consume_task) 

1107 

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 ) 

1116 

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) 

1127 

1128 return processed_attachments 

1129 

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 

1161 

1162 f.write(message.obj.as_bytes()) 

1163 

1164 correspondent = self._get_correspondent(message, rule) 

1165 

1166 self.log.info( 

1167 f"Rule {rule}: " 

1168 f"Consuming eml from mail " 

1169 f"{message.subject} from {message.from_}", 

1170 ) 

1171 

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 ) 

1187 

1188 consume_task = consume_file.s( 

1189 input_doc=input_doc, 

1190 overrides=doc_overrides, 

1191 ).set(headers={"trigger_source": PaperlessTask.TriggerSource.EMAIL_CONSUME}) 

1192 

1193 queue_consumption_tasks( 

1194 consume_tasks=[consume_task], 

1195 rule=rule, 

1196 message=message, 

1197 uid_validity=self._current_uid_validity, 

1198 ) 

1199 

1200 processed_elements = 1 

1201 return processed_elements