diff --git a/lambdas/vendor_reply/handler.py b/lambdas/vendor_reply/handler.py new file mode 100644 index 0000000..72ee313 --- /dev/null +++ b/lambdas/vendor_reply/handler.py @@ -0,0 +1,125 @@ +import email +import json +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"}