import email import logging import os import re from datetime import datetime from email import policy import boto3 logger = logging.getLogger() logger.setLevel(logging.INFO) s3 = boto3.client("s3") dynamodb = boto3.resource("dynamodb") VENDOR_REPLIES_TABLE = os.environ.get("VENDOR_REPLIES_TABLE", "VendorReplies") def parse_raw_email(raw_bytes): msg = email.message_from_bytes(raw_bytes, policy=policy.default) subject = msg.get("Subject", "") sender = msg.get("From", "") to = msg.get("To", "") date = msg.get("Date", "") body = "" if msg.is_multipart(): for part in msg.walk(): content_type = part.get_content_type() if content_type == "text/plain": body = part.get_content() break elif content_type == "text/html" and not body: body = part.get_content() else: body = msg.get_content() body = strip_quoted_reply(body) return { "subject": subject, "sender": sender, "to": to, "date": date, "body": body, } def strip_quoted_reply(text): if not text: return "" lines = text.split("\n") result = [] for line in lines: if re.match(r"^On .+ wrote:$", line.strip()): break if line.strip().startswith("---Original Message---"): break if line.strip().startswith("> "): continue result.append(line) return "\n".join(result).strip() def extract_ids_from_subject(subject): wo_match = re.search(r"\[WO-(\d{8})", subject) dsp_match = re.search(r"(DSP-\d{5})", subject) wo_number = wo_match.group(1) if wo_match else None dispatch_number = dsp_match.group(1) if dsp_match else None return wo_number, dispatch_number def extract_sender_email(sender): match = re.search(r"<(.+?)>", sender) if match: return match.group(1) return sender.strip() def handler(event, context): for record in event.get("Records", []): bucket = record["s3"]["bucket"]["name"] key = record["s3"]["object"]["key"] logger.info(f"Processing vendor reply: s3://{bucket}/{key}") response = s3.get_object(Bucket=bucket, Key=key) raw_email = response["Body"].read() email_data = parse_raw_email(raw_email) logger.info(f"Subject: {email_data['subject']}") logger.info(f"From: {email_data['sender']}") wo_number, dispatch_number = extract_ids_from_subject(email_data["subject"]) if not wo_number: logger.warning(f"No WO number found in subject, skipping: {key}") continue sender_email = extract_sender_email(email_data["sender"]) now = datetime.utcnow().isoformat() reply_id = f"{wo_number}#{dispatch_number or 'none'}#{now}" table = dynamodb.Table(VENDOR_REPLIES_TABLE) table.put_item( Item={ "internal_wo_number": wo_number, "reply_id": reply_id, "dispatch_number": dispatch_number or "", "sender_email": sender_email, "reply_text": email_data["body"], "subject": email_data["subject"], "received_at": now, "source_email_s3_key": f"s3://{bucket}/{key}", "synced": False, } ) logger.info( f"Saved vendor reply: WO={wo_number} DSP={dispatch_number} from={sender_email}" ) return {"statusCode": 200, "body": "OK"}