mirror of
https://github.com/Sea-Haven-Industries/workorder-ingest.git
synced 2026-05-18 20:20:12 +00:00
126 lines
3.5 KiB
Python
126 lines
3.5 KiB
Python
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"}
|