mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-10-01 23:13:27 +00:00
174 lines
6.4 KiB
Python
174 lines
6.4 KiB
Python
|
|
"""DynamoDB persistence for the work-order email processor.
|
||
|
|
|
||
|
|
Owns the work-orders + comments table writes and the retry-idempotent event_id /
|
||
|
|
comment_id determinism (_header_date_iso stays coupled with save_event so the
|
||
|
|
key derivation is single-sourced). The boto3 DynamoDB resource is built lazily on
|
||
|
|
first use so tests can patch this module's ``dynamodb`` attribute before any real
|
||
|
|
client is constructed (moto-before-handler invariant).
|
||
|
|
"""
|
||
|
|
|
||
|
|
import hashlib
|
||
|
|
import logging
|
||
|
|
import os
|
||
|
|
from datetime import datetime, timezone
|
||
|
|
from email.utils import parsedate_to_datetime
|
||
|
|
|
||
|
|
import boto3
|
||
|
|
|
||
|
|
logger = logging.getLogger()
|
||
|
|
logger.setLevel(logging.INFO)
|
||
|
|
|
||
|
|
WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders")
|
||
|
|
COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments")
|
||
|
|
|
||
|
|
# Lazy cached DynamoDB resource. Keeps the public attribute name ``dynamodb`` so
|
||
|
|
# the test monkeypatch target changes module only, not attribute name.
|
||
|
|
dynamodb = None
|
||
|
|
|
||
|
|
|
||
|
|
def _get_dynamodb():
|
||
|
|
global dynamodb
|
||
|
|
if dynamodb is None:
|
||
|
|
dynamodb = boto3.resource("dynamodb")
|
||
|
|
return dynamodb
|
||
|
|
|
||
|
|
|
||
|
|
def save_work_order(parsed: dict, s3_key: str):
|
||
|
|
"""Create or update a work order in DynamoDB."""
|
||
|
|
table = _get_dynamodb().Table(WORK_ORDERS_TABLE)
|
||
|
|
work_order_id = parsed["work_order_id"]
|
||
|
|
now = datetime.now(timezone.utc).isoformat()
|
||
|
|
|
||
|
|
# Build update expression dynamically from non-null fields
|
||
|
|
field_map = {
|
||
|
|
"description": "description",
|
||
|
|
"status": "wo_status", # 'status' is a DynamoDB reserved word
|
||
|
|
"site_code": "site_code",
|
||
|
|
"building": "building",
|
||
|
|
"address": "address",
|
||
|
|
"severity": "severity",
|
||
|
|
"priority": "priority",
|
||
|
|
"date_reported": "date_reported",
|
||
|
|
"scheduled_start": "scheduled_start",
|
||
|
|
"due_date": "due_date",
|
||
|
|
"assigned_to": "assigned_to",
|
||
|
|
}
|
||
|
|
|
||
|
|
update_parts = ["#updated_at = :updated_at", "#source_key = :source_key"]
|
||
|
|
attr_names = {
|
||
|
|
"#updated_at": "updated_at",
|
||
|
|
"#source_key": "source_email_s3_key",
|
||
|
|
}
|
||
|
|
attr_values = {
|
||
|
|
":updated_at": now,
|
||
|
|
":source_key": s3_key,
|
||
|
|
}
|
||
|
|
|
||
|
|
for src_field, dynamo_field in field_map.items():
|
||
|
|
value = parsed.get(src_field)
|
||
|
|
if value is not None:
|
||
|
|
placeholder = f":{dynamo_field}"
|
||
|
|
name_placeholder = f"#{dynamo_field}"
|
||
|
|
update_parts.append(f"{name_placeholder} = {placeholder}")
|
||
|
|
attr_names[name_placeholder] = dynamo_field
|
||
|
|
attr_values[placeholder] = value
|
||
|
|
|
||
|
|
# For new items, set created_at
|
||
|
|
update_parts.append("#created_at = if_not_exists(#created_at, :created_at)")
|
||
|
|
attr_names["#created_at"] = "created_at"
|
||
|
|
attr_values[":created_at"] = now
|
||
|
|
|
||
|
|
# Customer is always AMAZON for now
|
||
|
|
update_parts.append("#customer = :customer")
|
||
|
|
attr_names["#customer"] = "customer"
|
||
|
|
attr_values[":customer"] = "AMAZON"
|
||
|
|
|
||
|
|
# Track the record type (new_work_order, update, comment)
|
||
|
|
email_type = parsed.get("email_type")
|
||
|
|
if email_type:
|
||
|
|
update_parts.append("#record_type = :record_type")
|
||
|
|
attr_names["#record_type"] = "record_type"
|
||
|
|
attr_values[":record_type"] = email_type
|
||
|
|
|
||
|
|
table.update_item(
|
||
|
|
Key={"work_order_id": work_order_id},
|
||
|
|
UpdateExpression="SET " + ", ".join(update_parts),
|
||
|
|
ExpressionAttributeNames=attr_names,
|
||
|
|
ExpressionAttributeValues=attr_values,
|
||
|
|
)
|
||
|
|
|
||
|
|
logger.info(f"Saved work order {work_order_id}")
|
||
|
|
|
||
|
|
|
||
|
|
def _header_date_iso(header_date):
|
||
|
|
"""Parse an RFC 2822 Date header into a UTC ISO string, or None.
|
||
|
|
|
||
|
|
Deterministic for a given raw email, unlike model output. Total: any
|
||
|
|
unparseable/out-of-range header (incl. OverflowError from extreme years,
|
||
|
|
which is NOT a ValueError) yields None, never an exception -- a crafted
|
||
|
|
Date header must not be able to fail the invocation."""
|
||
|
|
if not header_date:
|
||
|
|
return None
|
||
|
|
try:
|
||
|
|
dt = parsedate_to_datetime(header_date)
|
||
|
|
if dt is None:
|
||
|
|
return None
|
||
|
|
if dt.tzinfo is None:
|
||
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
||
|
|
return dt.astimezone(timezone.utc).isoformat()
|
||
|
|
except (TypeError, ValueError, OverflowError, OSError):
|
||
|
|
return None
|
||
|
|
|
||
|
|
|
||
|
|
def save_event(
|
||
|
|
parsed: dict,
|
||
|
|
s3_key: str,
|
||
|
|
object_key: str,
|
||
|
|
parse_method: str,
|
||
|
|
header_date: str | None,
|
||
|
|
):
|
||
|
|
"""Save an event to the events table. Every email creates an event entry.
|
||
|
|
|
||
|
|
Issue #23: the range key (attr name stays ``comment_id``) must be unique per
|
||
|
|
source email AND identical across Lambda async retries of the same S3 object.
|
||
|
|
The uniqueness suffix is a deterministic hash of the S3 object key alone (not
|
||
|
|
the s3:// URI, so it is stable across a bucket rename), and wall-clock now()
|
||
|
|
is kept OUT of the key. When no time is available we use the literal
|
||
|
|
'nocomment' segment rather than now() -- otherwise each retry would produce a
|
||
|
|
different key and duplicate the row. Two distinct emails on the same WO map
|
||
|
|
to distinct object keys -> distinct rows.
|
||
|
|
|
||
|
|
Advisory A1: the time segment may come from parsed comment_time ONLY on the
|
||
|
|
template path, where it is a pure function of the raw email. On the AI path
|
||
|
|
the model can return a different comment_time on a retry (even at
|
||
|
|
temperature 0 determinism is not guaranteed), which would fork the key and
|
||
|
|
duplicate the row -- so there the segment derives from the email's Date
|
||
|
|
header instead."""
|
||
|
|
table = _get_dynamodb().Table(COMMENTS_TABLE)
|
||
|
|
|
||
|
|
work_order_id = parsed["work_order_id"]
|
||
|
|
email_type = parsed.get("email_type", "unknown")
|
||
|
|
key_suffix = hashlib.sha256(object_key.encode("utf-8")).hexdigest()[:12]
|
||
|
|
comment_time = parsed.get("comment_time")
|
||
|
|
if parse_method == "template":
|
||
|
|
time_part = comment_time if comment_time else "nocomment"
|
||
|
|
else:
|
||
|
|
time_part = _header_date_iso(header_date) or "nocomment"
|
||
|
|
event_id = f"{work_order_id}#{time_part}#{key_suffix}"
|
||
|
|
# created_at is display-only; may fall back to now() without affecting the key.
|
||
|
|
display_time = comment_time or datetime.now(timezone.utc).isoformat()
|
||
|
|
|
||
|
|
item = {
|
||
|
|
"work_order_id": work_order_id,
|
||
|
|
"comment_id": event_id, # keeping key name for table compatibility
|
||
|
|
"record_type": email_type,
|
||
|
|
"commenter": parsed.get("commenter") or "",
|
||
|
|
"text": parsed.get("comment_text") or "",
|
||
|
|
"created_at": display_time,
|
||
|
|
"source_email_s3_key": s3_key,
|
||
|
|
"ingested_at": datetime.now(timezone.utc).isoformat(),
|
||
|
|
}
|
||
|
|
|
||
|
|
table.put_item(Item=item)
|
||
|
|
logger.info(f"Saved event {event_id} (type={email_type})")
|