incomming mail processing
This commit is contained in:
@@ -0,0 +1,308 @@
|
||||
## Copyright © 2026 Olaf Kolkman
|
||||
## SPDX-License-Identifier: GPL-3.0-or-later
|
||||
|
||||
from email import policy
|
||||
from email.header import decode_header, make_header
|
||||
from email.parser import BytesParser
|
||||
from email.utils import getaddresses
|
||||
from html.parser import HTMLParser
|
||||
import imaplib
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
from threading import Event
|
||||
|
||||
from backend.app.core.config import settings
|
||||
from backend.app.database import get_connection
|
||||
from backend.app.services.link_service import create_link, normalize_tags
|
||||
from backend.app.services.secret_store import decrypt_secret, encrypt_secret
|
||||
|
||||
URL_PATTERN = re.compile(r'https?://[^\s<>()"\']+')
|
||||
HASHTAG_PATTERN = re.compile(r'(?<!\w)#([A-Za-z0-9][\w-]*)')
|
||||
|
||||
|
||||
class _HTMLTextExtractor(HTMLParser):
|
||||
def __init__(self) -> None:
|
||||
super().__init__()
|
||||
self.text_parts: list[str] = []
|
||||
self.hrefs: list[str] = []
|
||||
|
||||
def handle_data(self, data: str) -> None:
|
||||
if data:
|
||||
self.text_parts.append(data)
|
||||
|
||||
def handle_starttag(self, tag: str, attrs: list[tuple[str, str | None]]) -> None:
|
||||
if tag.lower() != 'a':
|
||||
return
|
||||
for name, value in attrs:
|
||||
if name.lower() == 'href' and value:
|
||||
self.hrefs.append(value)
|
||||
|
||||
|
||||
def _decode_header_value(value: str | None) -> str:
|
||||
if not value:
|
||||
return ''
|
||||
return str(make_header(decode_header(value))).strip()
|
||||
|
||||
|
||||
def _extract_message_text(message) -> tuple[str, list[str]]:
|
||||
plain_parts: list[str] = []
|
||||
html_parts: list[str] = []
|
||||
html_hrefs: list[str] = []
|
||||
|
||||
parts = message.walk() if message.is_multipart() else [message]
|
||||
for part in parts:
|
||||
if part.get_content_maintype() == 'multipart':
|
||||
continue
|
||||
content_type = part.get_content_type()
|
||||
try:
|
||||
content = part.get_content()
|
||||
except (LookupError, UnicodeDecodeError):
|
||||
continue
|
||||
if not isinstance(content, str):
|
||||
continue
|
||||
if content_type == 'text/plain':
|
||||
plain_parts.append(content)
|
||||
elif content_type == 'text/html':
|
||||
parser = _HTMLTextExtractor()
|
||||
parser.feed(content)
|
||||
html_parts.append(' '.join(parser.text_parts))
|
||||
html_hrefs.extend(parser.hrefs)
|
||||
|
||||
text = '\n\n'.join(part.strip() for part in (plain_parts or html_parts) if part.strip())
|
||||
urls = list(dict.fromkeys([*html_hrefs, *URL_PATTERN.findall(text)]))
|
||||
return text, urls
|
||||
|
||||
|
||||
def _extract_comment(text: str) -> str:
|
||||
without_links = URL_PATTERN.sub(' ', text)
|
||||
without_tags = HASHTAG_PATTERN.sub(' ', without_links)
|
||||
lines = []
|
||||
for raw_line in without_tags.splitlines():
|
||||
line = ' '.join(raw_line.split()).strip()
|
||||
if line:
|
||||
lines.append(line)
|
||||
return '\n'.join(lines)
|
||||
|
||||
|
||||
def _extract_tags(*values: str) -> list[str]:
|
||||
return normalize_tags([match.group(0) for value in values for match in HASHTAG_PATTERN.finditer(value or '')])
|
||||
|
||||
|
||||
def _resolve_sender_email(message) -> str:
|
||||
addresses = getaddresses(message.get_all('from', []))
|
||||
for _, email in addresses:
|
||||
if email:
|
||||
return email.strip().lower()
|
||||
return ''
|
||||
|
||||
|
||||
def parse_incoming_email(raw_message: bytes) -> dict:
|
||||
message = BytesParser(policy=policy.default).parsebytes(raw_message)
|
||||
subject = _decode_header_value(message.get('Subject'))
|
||||
text, urls = _extract_message_text(message)
|
||||
subject_without_tags = HASHTAG_PATTERN.sub(' ', subject)
|
||||
subject_without_link = URL_PATTERN.sub(' ', subject_without_tags)
|
||||
title = ' '.join(subject_without_link.split()).strip()
|
||||
tags = _extract_tags(subject, text)
|
||||
return {
|
||||
'message_id': _decode_header_value(message.get('Message-ID')),
|
||||
'from_email': _resolve_sender_email(message),
|
||||
'subject': subject,
|
||||
'title': title,
|
||||
'comment': _extract_comment(text),
|
||||
'tags': tags,
|
||||
'url': urls[0] if urls else None,
|
||||
}
|
||||
|
||||
|
||||
def get_imap_settings() -> dict:
|
||||
values = {
|
||||
'imap_host': settings.imap_host,
|
||||
'imap_port': settings.imap_port,
|
||||
'imap_username': settings.imap_username,
|
||||
'imap_password': settings.imap_password,
|
||||
'imap_mailbox': settings.imap_mailbox,
|
||||
'imap_use_ssl': settings.imap_use_ssl,
|
||||
'imap_poll_interval_seconds': settings.imap_poll_interval_seconds,
|
||||
}
|
||||
with get_connection() as conn:
|
||||
row = conn.execute('SELECT value FROM app_settings WHERE name = ?', ('imap',)).fetchone()
|
||||
if row:
|
||||
values.update(json.loads(row['value']))
|
||||
values['imap_password'] = decrypt_secret(values['imap_password'])
|
||||
return values
|
||||
|
||||
|
||||
def save_imap_settings(values: dict) -> None:
|
||||
stored_values = {
|
||||
**values,
|
||||
'imap_password': encrypt_secret(values['imap_password']),
|
||||
}
|
||||
with get_connection() as conn:
|
||||
conn.execute(
|
||||
'''INSERT INTO app_settings (name, value, updated_at) VALUES (?, ?, CURRENT_TIMESTAMP)
|
||||
ON CONFLICT(name) DO UPDATE SET value = excluded.value, updated_at = CURRENT_TIMESTAMP''',
|
||||
('imap', json.dumps(stored_values)),
|
||||
)
|
||||
conn.commit()
|
||||
|
||||
|
||||
def imap_configured(imap_values: dict | None = None) -> bool:
|
||||
imap = imap_values or get_imap_settings()
|
||||
return bool(imap['imap_host'] and imap['imap_username'] and imap['imap_password'] and imap['imap_mailbox'])
|
||||
|
||||
|
||||
def _find_user_by_email(address: str) -> dict | None:
|
||||
email = address.strip().lower()
|
||||
if not email:
|
||||
return None
|
||||
with get_connection() as conn:
|
||||
row = conn.execute(
|
||||
'SELECT id, username FROM users WHERE lower(email) = ? AND email_verified = 1',
|
||||
(email,),
|
||||
).fetchone()
|
||||
if row is None:
|
||||
row = conn.execute(
|
||||
'''SELECT users.id, users.username
|
||||
FROM user_email_addresses
|
||||
JOIN users ON users.id = user_email_addresses.user_id
|
||||
WHERE lower(user_email_addresses.email) = ? AND user_email_addresses.verified = 1''',
|
||||
(email,),
|
||||
).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
|
||||
def _already_processed(mailbox: str, uid: str) -> bool:
|
||||
with get_connection() as conn:
|
||||
row = conn.execute(
|
||||
'SELECT 1 FROM inbound_email_messages WHERE mailbox = ? AND uid = ? LIMIT 1',
|
||||
(mailbox, uid),
|
||||
).fetchone()
|
||||
return row is not None
|
||||
|
||||
|
||||
def _record_message(mailbox: str, uid: str, status: str, parsed: dict, user_id: str | None = None, link_id: str | None = None, details: dict | None = None) -> None:
|
||||
payload = details or {}
|
||||
with get_connection() as conn:
|
||||
conn.execute(
|
||||
'''INSERT INTO inbound_email_messages (mailbox, uid, message_id, user_id, link_id, status, details)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||
ON CONFLICT(mailbox, uid) DO UPDATE SET
|
||||
message_id = excluded.message_id,
|
||||
user_id = excluded.user_id,
|
||||
link_id = excluded.link_id,
|
||||
status = excluded.status,
|
||||
details = excluded.details''',
|
||||
(mailbox, uid, parsed.get('message_id'), user_id, link_id, status, json.dumps(payload)),
|
||||
)
|
||||
conn.commit()
|
||||
|
||||
|
||||
def ingest_message(mailbox: str, uid: str, raw_message: bytes) -> dict:
|
||||
if _already_processed(mailbox, uid):
|
||||
return {'status': 'duplicate', 'mailbox': mailbox, 'uid': uid}
|
||||
|
||||
parsed = parse_incoming_email(raw_message)
|
||||
sender = parsed['from_email']
|
||||
if not sender:
|
||||
_record_message(mailbox, uid, 'ignored_missing_sender', parsed)
|
||||
return {'status': 'ignored_missing_sender', 'mailbox': mailbox, 'uid': uid}
|
||||
|
||||
user = _find_user_by_email(sender)
|
||||
if user is None:
|
||||
_record_message(mailbox, uid, 'ignored_unknown_sender', parsed, details={'from_email': sender})
|
||||
return {'status': 'ignored_unknown_sender', 'mailbox': mailbox, 'uid': uid, 'from_email': sender}
|
||||
|
||||
if not parsed['url']:
|
||||
_record_message(mailbox, uid, 'ignored_no_link', parsed, user_id=user['id'])
|
||||
return {'status': 'ignored_no_link', 'mailbox': mailbox, 'uid': uid, 'user_id': user['id']}
|
||||
|
||||
title = parsed['title'] or parsed['url']
|
||||
record = create_link(user['id'], title, parsed['url'], parsed['comment'], None, parsed['tags'])
|
||||
_record_message(mailbox, uid, 'stored', parsed, user_id=user['id'], link_id=record['id'])
|
||||
return {
|
||||
'status': 'stored',
|
||||
'mailbox': mailbox,
|
||||
'uid': uid,
|
||||
'user_id': user['id'],
|
||||
'link_id': record['id'],
|
||||
'link': record,
|
||||
}
|
||||
|
||||
|
||||
def _open_imap_client(imap_values: dict):
|
||||
client_class = imaplib.IMAP4_SSL if imap_values['imap_use_ssl'] else imaplib.IMAP4
|
||||
return client_class(imap_values['imap_host'], imap_values['imap_port'])
|
||||
|
||||
|
||||
def _fetch_message_bytes(response_data) -> bytes | None:
|
||||
for item in response_data or []:
|
||||
if isinstance(item, tuple) and len(item) > 1 and isinstance(item[1], (bytes, bytearray)):
|
||||
return bytes(item[1])
|
||||
return None
|
||||
|
||||
|
||||
def poll_inbox_once(imap_values: dict | None = None, client_factory=None) -> dict:
|
||||
values = imap_values or get_imap_settings()
|
||||
if not imap_configured(values):
|
||||
return {'status': 'not_configured', 'processed': 0, 'stored': 0, 'ignored': 0}
|
||||
|
||||
mailbox = values['imap_mailbox']
|
||||
client = (client_factory or _open_imap_client)(values)
|
||||
processed = 0
|
||||
stored = 0
|
||||
ignored = 0
|
||||
try:
|
||||
client.login(values['imap_username'], values['imap_password'])
|
||||
status, _ = client.select(mailbox)
|
||||
if status != 'OK':
|
||||
raise RuntimeError(f'Could not select IMAP mailbox {mailbox}')
|
||||
status, data = client.uid('search', None, 'UNSEEN')
|
||||
if status != 'OK':
|
||||
raise RuntimeError('Could not list unseen IMAP messages')
|
||||
raw_uids = data[0].split() if data and data[0] else []
|
||||
for uid in raw_uids:
|
||||
uid_value = uid.decode('utf-8') if isinstance(uid, bytes) else str(uid)
|
||||
status, message_data = client.uid('fetch', uid, '(BODY.PEEK[])')
|
||||
if status != 'OK':
|
||||
raise RuntimeError(f'Could not fetch IMAP message {uid_value}')
|
||||
raw_message = _fetch_message_bytes(message_data)
|
||||
if raw_message is None:
|
||||
continue
|
||||
result = ingest_message(mailbox, uid_value, raw_message)
|
||||
processed += 1
|
||||
if result['status'] == 'stored':
|
||||
stored += 1
|
||||
else:
|
||||
ignored += 1
|
||||
if result['status'] in {
|
||||
'stored', 'duplicate', 'ignored_missing_sender', 'ignored_unknown_sender', 'ignored_no_link',
|
||||
}:
|
||||
client.uid('store', uid, '+FLAGS', '(\\Seen)')
|
||||
return {'status': 'ok', 'processed': processed, 'stored': stored, 'ignored': ignored}
|
||||
finally:
|
||||
try:
|
||||
client.logout()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def run_imap_polling(stop_event: Event, logger: logging.Logger | None = None) -> None:
|
||||
active_logger = logger or logging.getLogger(__name__)
|
||||
while not stop_event.is_set():
|
||||
values = get_imap_settings()
|
||||
interval = max(5, int(values.get('imap_poll_interval_seconds', settings.imap_poll_interval_seconds or 60)))
|
||||
if not imap_configured(values):
|
||||
stop_event.wait(interval)
|
||||
continue
|
||||
try:
|
||||
summary = poll_inbox_once(values)
|
||||
if summary['processed']:
|
||||
active_logger.info(
|
||||
'IMAP poll processed=%s stored=%s ignored=%s mailbox=%s',
|
||||
summary['processed'], summary['stored'], summary['ignored'], values['imap_mailbox'],
|
||||
)
|
||||
except Exception as error:
|
||||
active_logger.warning('IMAP polling failed mailbox=%s error=%s', values.get('imap_mailbox'), error)
|
||||
stop_event.wait(interval)
|
||||
Reference in New Issue
Block a user