mirror of
https://github.com/Sea-Haven-Industries/workorder-ingest.git
synced 2026-05-18 20:20:12 +00:00
Add vendor reply Lambda for inbound email processing
- Parses raw email from S3 vendor-inbound/ prefix - Extracts WO number and dispatch number from subject [WO-XXXXXXXX-DSP-XXXXX] - Strips quoted reply text (On...wrote:, >, ---Original Message---) - Writes to DynamoDB VendorReplies table - Triggered by S3 ObjectCreated events via SES receipt rule
This commit is contained in:
parent
df65d34eb0
commit
2071017aef
1 changed files with 125 additions and 0 deletions
125
lambdas/vendor_reply/handler.py
Normal file
125
lambdas/vendor_reply/handler.py
Normal file
|
|
@ -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"}
|
||||
Loading…
Add table
Reference in a new issue