""" Email Ingestion Service — convert inbound emails into tickets. Architecture ------------ A scheduled APScheduler job (check_inbound_email) runs every N minutes, connects to a configured IMAP mailbox, and converts unread messages into tickets. The job is wired into create_app() alongside the SLA scheduler. Sender resolution ----------------- The From address is matched against existing User.email rows. - Match found → ticket is created under that user's account. - No match found → ticket is created under a configurable fallback user (default: the system admin). A comment is prepended noting the external sender so IT staff can follow up. Duplicate suppression --------------------- Message-IDs (from the Message-ID header) are stored in the SystemSetting key email_ingested_message_ids as a comma-separated list (capped at 500 entries). Re-delivering an already-processed message is a no-op. Configuration (all stored in SystemSetting, editable from admin/settings) ---------- email_ingestion_enabled '1' / '0' email_ingestion_host IMAP server hostname email_ingestion_port IMAP port (default 993) email_ingestion_user Mailbox username / email address email_ingestion_password Mailbox password email_ingestion_interval Poll interval in minutes (default 5) email_ingestion_folder IMAP folder to watch (default INBOX) email_ingestion_move_to Folder to move processed mail into (default Processed) """ import email import imaplib import logging import os import re import uuid from datetime import datetime from email.header import decode_header from email.utils import parseaddr, getaddresses logger = logging.getLogger(__name__) _MAX_STORED_IDS = 500 # cap on the message-ID suppression list _MAX_ATTACH_BYTES = 16 * 1024 * 1024 # 16 MB per attachment — matches manual upload limit _MAX_ATTACH_COUNT = 10 # max attachments per email _ALLOWED_EXT = { # mirrors ALLOWED_EXT in tickets.py 'png', 'jpg', 'jpeg', 'gif', 'pdf', 'doc', 'docx', 'txt', 'zip', 'log', } # ── Helpers ─────────────────────────────────────────────────────────────────── def _decode_header_value(raw): """Decode an RFC-2047 encoded email header value to a plain string.""" if raw is None: return '' parts = [] for chunk, charset in decode_header(raw): if isinstance(chunk, bytes): try: parts.append(chunk.decode(charset or 'utf-8', errors='replace')) except (LookupError, UnicodeDecodeError): parts.append(chunk.decode('utf-8', errors='replace')) else: parts.append(chunk) return ''.join(parts).strip() def _extract_plain_text(msg): """Walk a MIME message and return the first text/plain part, or a stripped-down version of the first text/html part as a fallback.""" plain = None html = None if msg.is_multipart(): for part in msg.walk(): ct = part.get_content_type() cd = str(part.get('Content-Disposition', '')) if 'attachment' in cd: continue if ct == 'text/plain' and plain is None: try: plain = part.get_payload(decode=True).decode( part.get_content_charset() or 'utf-8', errors='replace' ) except Exception: pass elif ct == 'text/html' and html is None: try: html = part.get_payload(decode=True).decode( part.get_content_charset() or 'utf-8', errors='replace' ) except Exception: pass else: ct = msg.get_content_type() try: body = msg.get_payload(decode=True).decode( msg.get_content_charset() or 'utf-8', errors='replace' ) except Exception: body = '' if ct == 'text/plain': plain = body elif ct == 'text/html': html = body if plain: return plain.strip() if html: # Minimal HTML → plain text strip text = re.sub(r'', '\n', html, flags=re.IGNORECASE) text = re.sub(r'<[^>]+>', '', text) import html as html_module return html_module.unescape(text).strip() return '' def _extract_attachments(msg): """Walk a MIME message and collect all file attachments. Returns a list of (filename, content_type, data_bytes) tuples. Selection criteria ------------------ A MIME part is treated as an attachment when it meets ANY of these: - Content-Disposition is 'attachment' (standard), OR - Content-Disposition is 'inline' and the part has a filename parameter (common for images pasted inline, e.g. screenshots). Parts without a usable filename, parts that are the text body (text/plain or text/html without a filename), and parts with Content-Disposition: inline but no filename are all skipped. Limits enforced here (caller enforces the count cap separately): - Files larger than _MAX_ATTACH_BYTES are skipped with a warning. - Files whose extension is not in _ALLOWED_EXT are skipped with a warning. """ attachments = [] for part in msg.walk(): cd = part.get('Content-Disposition', '') or '' ct = part.get_content_type() or 'application/octet-stream' # Derive filename from Content-Disposition or Content-Type params filename = ( part.get_filename() or part.get_param('name') # some clients use Content-Type: name= or part.get_param('filename') ) is_attachment = 'attachment' in cd.lower() is_inline_file = 'inline' in cd.lower() and filename if not (is_attachment or is_inline_file): continue if not filename: continue # Decode RFC-2047 encoded filenames filename = _decode_header_value(filename) if not filename: continue # Extension check ext = filename.rsplit('.', 1)[-1].lower() if '.' in filename else '' if ext not in _ALLOWED_EXT: logger.warning( f'[EMAIL INGEST] Skipping attachment "{filename}" ' f'— extension ".{ext}" not in allowed list' ) continue # Read the payload try: data = part.get_payload(decode=True) except Exception as exc: logger.warning(f'[EMAIL INGEST] Could not decode attachment "{filename}": {exc}') continue if not data: continue # Size check if len(data) > _MAX_ATTACH_BYTES: size_mb = len(data) / (1024 * 1024) logger.warning( f'[EMAIL INGEST] Skipping attachment "{filename}" ' f'— size {size_mb:.1f} MB exceeds {_MAX_ATTACH_BYTES // (1024*1024)} MB limit' ) continue attachments.append((filename, ct, data)) return attachments def _save_email_attachments(attachments, ticket_id, uploader_id, upload_dir): """Persist attachment files to disk and create Attachment DB rows. This mirrors the logic in tickets.py:save_attachment() exactly, using UUID-based stored names so files never collide. Returns the number of attachments successfully saved. Caller is responsible for committing the session afterwards. """ from app.models import Attachment from app import db saved = 0 for filename, content_type, data in attachments: ext = filename.rsplit('.', 1)[-1].lower() if '.' in filename else 'bin' stored_name = f"{uuid.uuid4().hex}.{ext}" filepath = os.path.join(upload_dir, stored_name) try: with open(filepath, 'wb') as fh: fh.write(data) except OSError as exc: logger.error(f'[EMAIL INGEST] Failed to write attachment "{filename}": {exc}') continue att = Attachment( ticket_id = ticket_id, comment_id = None, filename = filename, stored_name = stored_name, file_size = len(data), mime_type = content_type, uploaded_by = uploader_id, ) db.session.add(att) saved += 1 logger.info( f'[EMAIL INGEST] Saved attachment "{filename}" ' f'({len(data)} bytes) → {stored_name} for ticket_id={ticket_id}' ) return saved def _get_setting(key, default=''): """Read a SystemSetting value inside an existing app context.""" from app.models import SystemSetting return SystemSetting.get(key, default) def _load_seen_ids(): raw = _get_setting('email_ingested_message_ids', '') return set(x.strip() for x in raw.split(',') if x.strip()) def _save_seen_ids(seen: set): from app.models import SystemSetting from app import db # Keep the most recent N IDs to prevent unbounded growth trimmed = sorted(seen)[-_MAX_STORED_IDS:] SystemSetting.set( 'email_ingested_message_ids', ','.join(trimmed), 'Message-IDs of emails already converted to tickets', ) db.session.commit() # ── Core ingestion logic ────────────────────────────────────────────────────── def check_inbound_email(app): """Entry point called by APScheduler. Wraps _run_ingestion with error isolation so a transient IMAP failure never kills the worker.""" with app.app_context(): try: if _get_setting('email_ingestion_enabled', '0') != '1': return _run_ingestion(app) except Exception as exc: logger.error(f'[EMAIL INGEST] Unhandled error: {exc}', exc_info=True) def _run_ingestion(app): from app import db from app.models import (Ticket, TicketStatus, TicketPriority, TicketCategory, User, UserRole) from app.services.notification_service import notify_new_ticket from app.services.sla_service import set_due_date host = _get_setting('email_ingestion_host', '') port = int(_get_setting('email_ingestion_port', '993')) username = _get_setting('email_ingestion_user', '') password = _get_setting('email_ingestion_password', '') folder = _get_setting('email_ingestion_folder', 'INBOX') move_to = _get_setting('email_ingestion_move_to', 'Processed') if not host or not username or not password: logger.warning('[EMAIL INGEST] Missing IMAP credentials — skipping') return seen_ids = _load_seen_ids() new_ids = set() created = 0 try: imap = imaplib.IMAP4_SSL(host, port) imap.login(username, password) except Exception as exc: logger.error(f'[EMAIL INGEST] IMAP login failed: {exc}') return try: imap.select(folder) # Search for unseen messages only status, data = imap.search(None, 'UNSEEN') if status != 'OK' or not data[0]: return msg_ids = data[0].split() logger.info(f'[EMAIL INGEST] Found {len(msg_ids)} unseen message(s) in {folder}') for num in msg_ids: try: _, raw = imap.fetch(num, '(RFC822)') msg = email.message_from_bytes(raw[0][1]) message_id = msg.get('Message-ID', '').strip() if message_id and message_id in seen_ids: logger.debug(f'[EMAIL INGEST] Skipping duplicate {message_id}') continue # ── Parse headers ───────────────────────────────────────────── subject = _decode_header_value(msg.get('Subject', '(No Subject)')) from_raw = msg.get('From', '') from_name, from_email = parseaddr(from_raw) from_email = from_email.lower().strip() body = _extract_plain_text(msg) if not body: body = f'[Email received from {from_email} with no readable body]' # Truncate very long bodies to 8000 chars if len(body) > 8000: body = body[:8000] + '\n\n[…message truncated…]' # ── Extract attachments from the MIME message ────────────── # Done before any DB work so we can log counts accurately. raw_attachments = _extract_attachments(msg) if len(raw_attachments) > _MAX_ATTACH_COUNT: logger.warning( f'[EMAIL INGEST] Email has {len(raw_attachments)} attachments ' f'— capping at {_MAX_ATTACH_COUNT}' ) raw_attachments = raw_attachments[:_MAX_ATTACH_COUNT] # ── Resolve sender to a user ────────────────────────────────── sender_user = User.query.filter_by( email=from_email, is_active=True ).first() if sender_user: created_by_id = sender_user.id external_note = None else: # Fall back to the first active admin fallback = User.query.filter( User.role == UserRole.ADMIN, User.is_active == True, ).first() if not fallback: logger.warning( f'[EMAIL INGEST] No fallback admin found, skipping: {from_email}' ) continue created_by_id = fallback.id external_note = ( f'**[Email received from unknown sender]**\n\n' f'From: {from_name} <{from_email}>\n\n' f'This ticket was automatically created from an inbound email. ' f'The sender is not a registered user — please follow up directly.' ) # ── Create the ticket ───────────────────────────────────────── ticket = Ticket( title = subject[:200], description = body, category = TicketCategory.OTHER, priority = TicketPriority.MEDIUM, status = TicketStatus.OPEN, created_by_id = created_by_id, ai_generated = False, ) ticket.ticket_number = ticket.generate_ticket_number() set_due_date(ticket, app) db.session.add(ticket) db.session.flush() # Prepend external-sender note as a comment if needed if external_note: from app.models import Comment from app.services.validation_service import render_comment_body comment = Comment( ticket_id = ticket.id, author_id = created_by_id, body = render_comment_body(external_note), is_internal= True, # IT staff only ) db.session.add(comment) # ── Save email attachments ──────────────────────────────── # ticket.id is available after flush(), so Attachment rows # can reference it. Files are written to UPLOAD_FOLDER and # flushed into the session here; the single commit() below # persists ticket + comment + all attachments atomically. upload_dir = app.config['UPLOAD_FOLDER'] os.makedirs(upload_dir, exist_ok=True) attach_count = _save_email_attachments( raw_attachments, ticket.id, created_by_id, upload_dir ) from app.services.log_service import log_action log_action( created_by_id, 'ticket_create_email', 'ticket', ticket.id, f'ticket_number={ticket.ticket_number} ' f'from_email={from_email} message_id={message_id or "none"} ' f'attachments={attach_count}' ) db.session.commit() logger.info( f'[EMAIL INGEST] Created ticket {ticket.ticket_number} ' f'from={from_email} subject="{subject[:60]}" ' f'attachments={attach_count}' ) notify_new_ticket(ticket) # Track message ID for deduplication if message_id: new_ids.add(message_id) created += 1 # Move processed message to done folder try: imap.create(move_to) except Exception: pass # folder may already exist imap.copy(num, move_to) imap.store(num, '+FLAGS', '\\Deleted') except Exception as exc: db.session.rollback() logger.error(f'[EMAIL INGEST] Failed to process message {num}: {exc}', exc_info=True) imap.expunge() finally: try: imap.logout() except Exception: pass if new_ids: _save_seen_ids(seen_ids | new_ids) if created: logger.info(f'[EMAIL INGEST] Run complete — {created} ticket(s) created')