commit df65d34eb09a9805265918df204d5cd03354ebb3 Author: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Fri Apr 3 16:37:04 2026 -0400 Initial commit: work order email ingestion pipeline Serverless AWS pipeline that receives Amazon APM work order emails via SES, parses them with Claude AI, and stores structured data in DynamoDB. Includes a web UI dashboard for viewing work orders and event history. Co-Authored-By: Claude Opus 4.6 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..5294260 --- /dev/null +++ b/.gitignore @@ -0,0 +1,11 @@ +__pycache__/ +*.py[cod] +*.egg-info/ +dist/ +build/ +.venv/ +venv/ +node_modules/ +cdk.out/ +.env +*.eml diff --git a/README.md b/README.md new file mode 100644 index 0000000..86159b6 --- /dev/null +++ b/README.md @@ -0,0 +1,70 @@ +# Work Order Ingest + +AWS serverless pipeline that ingests work order emails from Amazon's APM system (Hexagon EAM / HxGN SmartCloud), parses them with Claude AI, and stores structured data in DynamoDB. + +## Architecture + +``` +Email (amazon@seahavenind.com) + → Gmail filter (from: noreply@hxgnsmartcloud.com) + → Forwards to apm@int.seahaven.com (SES) + → S3 bucket (raw email storage) + → Lambda (email processor) + → Claude Haiku (structured extraction) + → DynamoDB (WorkOrders + WorkOrderComments) +``` + +## Components + +- **Email Processor Lambda** (`lambdas/email_processor/`) - Parses raw emails, sends to Claude for structured extraction, writes to DynamoDB +- **Web UI Lambda** (`lambdas/web_ui/`) - Server-rendered HTML dashboard for viewing work orders and event history +- **CDK Stack** (`cdk/`) - Infrastructure as code for all AWS resources +- **Shared Models** (`shared/`) - Data models for work orders, comments, and parsed emails + +## Extracted Fields + +- Work Order ID, Description, Status +- Site Code, Building, Address +- Severity, Priority +- Due Date, Date Reported, Scheduled Start +- Assigned To, Commenter, Comment Text +- Record Type (new_work_order, update, comment, cancellation) + +## DynamoDB Tables + +- **WorkOrders** - Latest state of each work order (PK: `work_order_id`) +- **WorkOrderComments** - Event history per work order (PK: `work_order_id`, SK: `comment_id`) + +## Local Testing + +```bash +python3 -m venv .venv +source .venv/bin/activate +pip install anthropic boto3 +export ANTHROPIC_API_KEY=sk-ant-... +python test_local.py +``` + +## Deployment + +```bash +cd cdk +pip install aws-cdk-lib constructs +npx cdk bootstrap # first time only +npx cdk deploy +``` + +After deploying, store your Anthropic API key in Secrets Manager: + +```bash +aws secretsmanager put-secret-value \ + --secret-id "workorder-ingest/anthropic-api-key" \ + --secret-string "sk-ant-..." \ + --region us-east-1 +``` + +## SES Configuration + +- Domain: `int.seahaven.com` (MX record pointing to `inbound-smtp.us-east-1.amazonaws.com`) +- Receipt Rule Set: `INBOUND_MAIL` +- Recipient: `apm@int.seahaven.com` diff --git a/cdk/__init__.py b/cdk/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/cdk/app.py b/cdk/app.py new file mode 100644 index 0000000..4063128 --- /dev/null +++ b/cdk/app.py @@ -0,0 +1,9 @@ +#!/usr/bin/env python3 +import aws_cdk as cdk +from stack import WorkorderIngestStack + +app = cdk.App() +WorkorderIngestStack(app, "WorkorderIngestStack", + env=cdk.Environment(region="us-east-1"), +) +app.synth() diff --git a/cdk/cdk.json b/cdk/cdk.json new file mode 100644 index 0000000..f28e013 --- /dev/null +++ b/cdk/cdk.json @@ -0,0 +1,6 @@ +{ + "app": "python3 app.py", + "context": { + "@aws-cdk/core:bootstrapQualifier": "hnb659fds" + } +} diff --git a/cdk/stack.py b/cdk/stack.py new file mode 100644 index 0000000..2770fc1 --- /dev/null +++ b/cdk/stack.py @@ -0,0 +1,172 @@ +"""CDK stack for the work order email ingestion pipeline.""" + +import aws_cdk as cdk +from aws_cdk import ( + Duration, + RemovalPolicy, + Stack, + aws_dynamodb as dynamodb, + aws_lambda as lambda_, + aws_s3 as s3, + aws_s3_notifications as s3n, + aws_ses as ses, + aws_ses_actions as ses_actions, + aws_secretsmanager as secretsmanager, + aws_iam as iam, +) +from constructs import Construct + + +class WorkorderIngestStack(Stack): + def __init__(self, scope: Construct, construct_id: str, **kwargs): + super().__init__(scope, construct_id, **kwargs) + + # --- S3 bucket for raw emails --- + email_bucket = s3.Bucket( + self, "EmailBucket", + bucket_name=f"workorder-ingest-emails-{self.account}", + removal_policy=RemovalPolicy.RETAIN, + lifecycle_rules=[ + s3.LifecycleRule(expiration=Duration.days(90)), + ], + ) + + # --- DynamoDB tables --- + work_orders_table = dynamodb.Table( + self, "WorkOrdersTable", + table_name="WorkOrders", + partition_key=dynamodb.Attribute( + name="work_order_id", + type=dynamodb.AttributeType.STRING, + ), + billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST, + removal_policy=RemovalPolicy.RETAIN, + ) + # GSI for querying by site code + work_orders_table.add_global_secondary_index( + index_name="site-code-index", + partition_key=dynamodb.Attribute( + name="site_code", + type=dynamodb.AttributeType.STRING, + ), + sort_key=dynamodb.Attribute( + name="updated_at", + type=dynamodb.AttributeType.STRING, + ), + ) + # GSI for querying by status + work_orders_table.add_global_secondary_index( + index_name="status-index", + partition_key=dynamodb.Attribute( + name="wo_status", + type=dynamodb.AttributeType.STRING, + ), + sort_key=dynamodb.Attribute( + name="updated_at", + type=dynamodb.AttributeType.STRING, + ), + ) + + comments_table = dynamodb.Table( + self, "CommentsTable", + table_name="WorkOrderComments", + partition_key=dynamodb.Attribute( + name="work_order_id", + type=dynamodb.AttributeType.STRING, + ), + sort_key=dynamodb.Attribute( + name="comment_id", + type=dynamodb.AttributeType.STRING, + ), + billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST, + removal_policy=RemovalPolicy.RETAIN, + ) + + # --- Secrets Manager for Anthropic API key --- + anthropic_secret = secretsmanager.Secret( + self, "AnthropicApiKey", + secret_name="workorder-ingest/anthropic-api-key", + description="Anthropic API key for work order email parsing", + ) + + # --- Lambda function --- + email_processor = lambda_.Function( + self, "EmailProcessor", + function_name="workorder-email-processor", + runtime=lambda_.Runtime.PYTHON_3_12, + handler="handler.handler", + code=lambda_.Code.from_asset( + "../lambdas/email_processor", + bundling=cdk.BundlingOptions( + image=lambda_.Runtime.PYTHON_3_12.bundling_image, + platform="linux/amd64", + command=[ + "bash", "-c", + "pip install --platform manylinux2014_x86_64 --only-binary=:all: " + "-r requirements.txt -t /asset-output && " + "cp -r . /asset-output/" + ], + ), + ), + timeout=Duration.seconds(60), + memory_size=256, + environment={ + "WORK_ORDERS_TABLE": work_orders_table.table_name, + "COMMENTS_TABLE": comments_table.table_name, + "ANTHROPIC_API_KEY_SECRET_ARN": anthropic_secret.secret_arn, + }, + ) + + # Grant permissions + email_bucket.grant_read(email_processor) + work_orders_table.grant_read_write_data(email_processor) + comments_table.grant_read_write_data(email_processor) + anthropic_secret.grant_read(email_processor) + + # S3 event notification → Lambda + email_bucket.add_event_notification( + s3.EventType.OBJECT_CREATED, + s3n.LambdaDestination(email_processor), + s3.NotificationKeyFilter(prefix="inbound/"), + ) + + # --- SES Receipt Rule --- + rule_set = ses.ReceiptRuleSet.from_receipt_rule_set_name( + self, "ExistingRuleSet", "INBOUND_MAIL", + ) + + rule_set.add_rule( + "WorkorderEmailRule", + recipients=["apm@int.seahaven.com"], + actions=[ + ses_actions.S3( + bucket=email_bucket, + object_key_prefix="inbound/", + ), + ], + ) + + # --- Web UI Lambda --- + web_ui = lambda_.Function( + self, "WebUI", + function_name="workorder-web-ui", + runtime=lambda_.Runtime.PYTHON_3_12, + handler="handler.handler", + code=lambda_.Code.from_asset("../lambdas/web_ui"), + timeout=Duration.seconds(15), + memory_size=128, + environment={ + "WORK_ORDERS_TABLE": work_orders_table.table_name, + "COMMENTS_TABLE": comments_table.table_name, + }, + ) + + work_orders_table.grant_read_data(web_ui) + comments_table.grant_read_data(web_ui) + + # Function URL for direct access + web_url = web_ui.add_function_url( + auth_type=lambda_.FunctionUrlAuthType.NONE, + ) + + cdk.CfnOutput(self, "WebUIUrl", value=web_url.url, description="Work Order Dashboard URL") diff --git a/lambdas/__init__.py b/lambdas/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/lambdas/email_processor/__init__.py b/lambdas/email_processor/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/lambdas/email_processor/handler.py b/lambdas/email_processor/handler.py new file mode 100644 index 0000000..ad78dc7 --- /dev/null +++ b/lambdas/email_processor/handler.py @@ -0,0 +1,270 @@ +""" +Email processor Lambda. + +Triggered by S3 events when SES delivers an email. +Parses the raw email, sends it to Claude for structured extraction, +then writes the result to DynamoDB. +""" + +import email +import json +import logging +import os +import re +from datetime import datetime +from email import policy +from typing import Optional + +import anthropic +import boto3 + +logger = logging.getLogger() +logger.setLevel(logging.INFO) + +s3 = boto3.client("s3") +dynamodb = boto3.resource("dynamodb") + +WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders") +COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments") +ANTHROPIC_API_KEY_SECRET_ARN = os.environ.get("ANTHROPIC_API_KEY_SECRET_ARN") + +EXTRACTION_PROMPT = """\ +You are an email parser for a facilities maintenance work order system. +The emails come from Amazon's APM system (via Hexagon EAM / HxGN SmartCloud). + +Analyze the following email and extract structured data. Return ONLY valid JSON with these fields: + +{ + "email_type": "new_work_order" | "update" | "comment" | "cancellation", + "work_order_id": "string or null", + "description": "work order description or null", + "status": "new" | "assigned" | "in_progress" | "on_hold" | "completed" | "cancelled" | "unknown", + "site_code": "building/site code like WIL1, ZDL8, etc. or null", + "building": "full building identifier or null", + "address": "physical address or null", + "severity": "severity level or null", + "priority": "priority level or null", + "date_reported": "ISO 8601 date or null", + "scheduled_start": "ISO 8601 date or null", + "due_date": "ISO 8601 date or null", + "assigned_to": "person/team assigned or null", + "commenter": "person who left a comment or null", + "comment_text": "the comment text or null", + "comment_time": "ISO 8601 datetime of the comment or null" +} + +Rules: +- "email_type" detection: + - "new_work_order": email announces a new WO assignment + - "comment": email contains a new comment on an existing WO + - "cancellation": email announces a WO has been cancelled + - "update": any other update to an existing WO (status change, reassignment, etc.) +- Extract the site_code from the building field (e.g., "WIL1" from "building WIL1") +- Dates should be converted to ISO 8601 format +- If a field is not present in the email, set it to null +- Do NOT invent or infer data that is not explicitly in the email +""" + + +def get_anthropic_client() -> anthropic.Anthropic: + """Create Anthropic client, fetching API key from Secrets Manager if configured.""" + if ANTHROPIC_API_KEY_SECRET_ARN: + secrets = boto3.client("secretsmanager") + secret = secrets.get_secret_value(SecretId=ANTHROPIC_API_KEY_SECRET_ARN) + api_key = secret["SecretString"] + return anthropic.Anthropic(api_key=api_key) + # Fall back to ANTHROPIC_API_KEY env var (for local testing) + return anthropic.Anthropic() + + +def parse_raw_email(raw_bytes: bytes) -> dict: + """Parse a raw email into subject, sender, body text.""" + msg = email.message_from_bytes(raw_bytes, policy=policy.default) + + subject = msg.get("Subject", "") + sender = msg.get("From", "") + to = msg.get("To", "") + cc = msg.get("Cc", "") + 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() + + return { + "subject": subject, + "sender": sender, + "to": to, + "cc": cc, + "date": date, + "body": body, + } + + +def extract_with_claude(email_data: dict) -> dict: + """Send parsed email to Claude for structured extraction.""" + client = get_anthropic_client() + + email_text = ( + f"Subject: {email_data['subject']}\n" + f"From: {email_data['sender']}\n" + f"To: {email_data['to']}\n" + f"CC: {email_data['cc']}\n" + f"Date: {email_data['date']}\n" + f"\n---\n\n" + f"{email_data['body']}" + ) + + response = client.messages.create( + model="claude-haiku-4-5-20251001", + max_tokens=1024, + messages=[ + { + "role": "user", + "content": f"{EXTRACTION_PROMPT}\n\nEMAIL:\n{email_text}", + } + ], + ) + + response_text = response.content[0].text + + # Extract JSON from response (handle markdown code blocks) + json_match = re.search(r"```(?:json)?\s*(.*?)```", response_text, re.DOTALL) + if json_match: + response_text = json_match.group(1) + + return json.loads(response_text.strip()) + + +def save_work_order(parsed: dict, s3_key: str): + """Create or update a work order in DynamoDB.""" + table = dynamodb.Table(WORK_ORDERS_TABLE) + work_order_id = parsed["work_order_id"] + now = datetime.utcnow().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 save_event(parsed: dict, s3_key: str): + """Save an event to the events table. Every email creates an event entry.""" + table = dynamodb.Table(COMMENTS_TABLE) + + work_order_id = parsed["work_order_id"] + email_type = parsed.get("email_type", "unknown") + event_time = parsed.get("comment_time") or datetime.utcnow().isoformat() + event_id = f"{work_order_id}#{event_time}" + + 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": event_time, + "source_email_s3_key": s3_key, + "ingested_at": datetime.utcnow().isoformat(), + } + + table.put_item(Item=item) + logger.info(f"Saved event {event_id} (type={email_type})") + + +def handler(event, context): + """Lambda entry point. Triggered by S3 ObjectCreated events.""" + for record in event.get("Records", []): + bucket = record["s3"]["bucket"]["name"] + key = record["s3"]["object"]["key"] + + logger.info(f"Processing email: s3://{bucket}/{key}") + + # Fetch raw email from S3 + response = s3.get_object(Bucket=bucket, Key=key) + raw_email = response["Body"].read() + + # Parse the raw email + email_data = parse_raw_email(raw_email) + logger.info(f"Subject: {email_data['subject']}") + + # Extract structured data with Claude + parsed = extract_with_claude(email_data) + logger.info(f"Parsed: type={parsed.get('email_type')}, wo={parsed.get('work_order_id')}") + + if not parsed.get("work_order_id"): + logger.warning(f"No work order ID found in email, skipping: {key}") + continue + + s3_key = f"s3://{bucket}/{key}" + + # Always upsert the work order with any new info + save_work_order(parsed, s3_key) + + # Save every email as an event for history tracking + save_event(parsed, s3_key) + + return {"statusCode": 200, "body": "OK"} diff --git a/lambdas/email_processor/requirements.txt b/lambdas/email_processor/requirements.txt new file mode 100644 index 0000000..acef34e --- /dev/null +++ b/lambdas/email_processor/requirements.txt @@ -0,0 +1,2 @@ +anthropic>=0.42.0 +boto3>=1.35.0 diff --git a/lambdas/web_ui/__init__.py b/lambdas/web_ui/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/lambdas/web_ui/handler.py b/lambdas/web_ui/handler.py new file mode 100644 index 0000000..3abc7de --- /dev/null +++ b/lambdas/web_ui/handler.py @@ -0,0 +1,241 @@ +""" +Web UI Lambda. + +Serves a simple HTML dashboard for viewing work orders and comments. +Accessed via Lambda Function URL. +""" + +import json +import os +import urllib.parse + +import boto3 + +dynamodb = boto3.resource("dynamodb") + +WORK_ORDERS_TABLE = os.environ.get("WORK_ORDERS_TABLE", "WorkOrders") +COMMENTS_TABLE = os.environ.get("COMMENTS_TABLE", "WorkOrderComments") + + +def get_work_orders(): + table = dynamodb.Table(WORK_ORDERS_TABLE) + response = table.scan() + items = response.get("Items", []) + # Sort by updated_at descending + items.sort(key=lambda x: x.get("updated_at", ""), reverse=True) + return items + + +def get_comments(work_order_id): + table = dynamodb.Table(COMMENTS_TABLE) + response = table.query( + KeyConditionExpression="work_order_id = :woid", + ExpressionAttributeValues={":woid": work_order_id}, + ) + items = response.get("Items", []) + items.sort(key=lambda x: x.get("created_at", ""), reverse=True) + return items + + +def render_badge(value, color_map): + color = color_map.get(value, "#9ca3af") + label = value.replace("_", " ").title() + return f'{label}' + + +STATUS_COLORS = { + "new": "#3b82f6", + "assigned": "#8b5cf6", + "in_progress": "#f59e0b", + "on_hold": "#6b7280", + "completed": "#10b981", + "cancelled": "#ef4444", + "unknown": "#9ca3af", +} + +RECORD_TYPE_COLORS = { + "new_work_order": "#3b82f6", + "comment": "#8b5cf6", + "update": "#f59e0b", + "cancellation": "#ef4444", + "unknown": "#9ca3af", +} + + +def render_work_order_detail(wo, events): + events_html = "" + if events: + for e in events: + record_type = e.get("record_type", "unknown") + commenter = e.get("commenter", "") + created = e.get("created_at", "") + text = e.get("text", "") + badge = render_badge(record_type, RECORD_TYPE_COLORS) + commenter_str = f"{commenter} — " if commenter else "" + + border_colors = { + "new_work_order": "#3b82f6", + "comment": "#8b5cf6", + "update": "#f59e0b", + "cancellation": "#ef4444", + } + border = border_colors.get(record_type, "#94a3b8") + + events_html += f""" +
+
+ {badge} {commenter_str}{created} +
+ {"
" + text + "
" if text else ""} +
""" + else: + events_html = '

No events yet.

' + + wo_id = wo.get("work_order_id", "") + fields = [ + ("Description", wo.get("description")), + ("Status", render_badge(wo.get("wo_status", "unknown"), STATUS_COLORS)), + ("Record Type", render_badge(wo.get("record_type", "unknown"), RECORD_TYPE_COLORS)), + ("Site Code", wo.get("site_code")), + ("Building", wo.get("building")), + ("Address", wo.get("address")), + ("Severity", wo.get("severity")), + ("Priority", wo.get("priority")), + ("Due Date", wo.get("due_date")), + ("Date Reported", wo.get("date_reported")), + ("Scheduled Start", wo.get("scheduled_start")), + ("Assigned To", wo.get("assigned_to")), + ("Created", wo.get("created_at")), + ("Last Updated", wo.get("updated_at")), + ] + + details_html = "" + for label, value in fields: + if value: + details_html += f""" +
+
{label}
+
{value}
+
""" + + return f""" + + + + + WO {wo_id} - Sea Haven + + + +
+
+ ← All Work Orders +
+
+

Work Order {wo_id}

+ {details_html} +
+
+

Events ({len(events)})

+ {events_html} +
+
+ +""" + + +def render_work_orders_list(work_orders): + rows = "" + for wo in work_orders: + wo_id = wo.get("work_order_id", "") + desc = wo.get("description", "") + site = wo.get("site_code", "") + status = wo.get("wo_status", "unknown") + record_type = wo.get("record_type", "unknown") + due = wo.get("due_date", "") + updated = wo.get("updated_at", "")[:16] + + rows += f""" + + {wo_id} + {desc} + {site} + {render_badge(record_type, RECORD_TYPE_COLORS)} + {render_badge(status, STATUS_COLORS)} + {due} + {updated} + """ + + return f""" + + + + + Work Orders - Sea Haven + + + +
+
+

Work Orders

+ {len(work_orders)} total +
+
+ + + + + + + + + + + + + + {rows if rows else ''} + +
WO #DescriptionSiteLast ActionStatusDue DateUpdated
No work orders yet.
+
+
+ +""" + + +def handler(event, context): + path = event.get("rawPath", "/") + qs = event.get("queryStringParameters") or {} + + if path == "/wo" and "id" in qs: + wo_id = qs["id"] + # Get work order + table = dynamodb.Table(WORK_ORDERS_TABLE) + result = table.get_item(Key={"work_order_id": wo_id}) + wo = result.get("Item") + if not wo: + return { + "statusCode": 404, + "headers": {"Content-Type": "text/html"}, + "body": "

Work order not found

", + } + events = get_comments(wo_id) + html = render_work_order_detail(wo, events) + else: + work_orders = get_work_orders() + html = render_work_orders_list(work_orders) + + return { + "statusCode": 200, + "headers": {"Content-Type": "text/html"}, + "body": html, + } diff --git a/lambdas/web_ui/requirements.txt b/lambdas/web_ui/requirements.txt new file mode 100644 index 0000000..3f3a438 --- /dev/null +++ b/lambdas/web_ui/requirements.txt @@ -0,0 +1 @@ +boto3>=1.35.0 diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..32b3387 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,2 @@ +aws-cdk-lib>=2.170.0 +constructs>=10.0.0 diff --git a/shared/__init__.py b/shared/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/shared/models.py b/shared/models.py new file mode 100644 index 0000000..cb50e29 --- /dev/null +++ b/shared/models.py @@ -0,0 +1,82 @@ +"""Data models for work order ingestion.""" + +from dataclasses import dataclass, field, asdict +from datetime import datetime +from enum import Enum +from typing import Optional + + +class EmailType(str, Enum): + NEW_WORK_ORDER = "new_work_order" + UPDATE = "update" + COMMENT = "comment" + UNKNOWN = "unknown" + + +class WorkOrderStatus(str, Enum): + NEW = "new" + ASSIGNED = "assigned" + IN_PROGRESS = "in_progress" + ON_HOLD = "on_hold" + COMPLETED = "completed" + CANCELLED = "cancelled" + UNKNOWN = "unknown" + + +@dataclass +class Comment: + work_order_id: str + comment_id: str # generated: {work_order_id}#{timestamp} + commenter: str + text: str + created_at: str # ISO 8601 + source_email_s3_key: str + ingested_at: str = field(default_factory=lambda: datetime.utcnow().isoformat()) + + def to_dynamo_item(self) -> dict: + return {k: v for k, v in asdict(self).items() if v is not None} + + +@dataclass +class WorkOrder: + work_order_id: str + description: Optional[str] = None + status: str = WorkOrderStatus.UNKNOWN.value + customer: str = "AMAZON" + site_code: Optional[str] = None + building: Optional[str] = None + address: Optional[str] = None + severity: Optional[str] = None + priority: Optional[str] = None + date_reported: Optional[str] = None + scheduled_start: Optional[str] = None + due_date: Optional[str] = None + assigned_to: Optional[str] = None + source_email_s3_key: Optional[str] = None + created_at: str = field(default_factory=lambda: datetime.utcnow().isoformat()) + updated_at: str = field(default_factory=lambda: datetime.utcnow().isoformat()) + + def to_dynamo_item(self) -> dict: + return {k: v for k, v in asdict(self).items() if v is not None} + + +@dataclass +class ParsedEmail: + """Result of AI parsing an inbound email.""" + email_type: str # EmailType value + work_order_id: Optional[str] = None + description: Optional[str] = None + status: Optional[str] = None + site_code: Optional[str] = None + building: Optional[str] = None + address: Optional[str] = None + severity: Optional[str] = None + priority: Optional[str] = None + date_reported: Optional[str] = None + scheduled_start: Optional[str] = None + due_date: Optional[str] = None + assigned_to: Optional[str] = None + commenter: Optional[str] = None + comment_text: Optional[str] = None + comment_time: Optional[str] = None + raw_subject: Optional[str] = None diff --git a/test_local.py b/test_local.py new file mode 100644 index 0000000..4098ff3 --- /dev/null +++ b/test_local.py @@ -0,0 +1,48 @@ +#!/usr/bin/env python3 +""" +Local test script - parses sample emails through Claude without AWS. + +Usage: + export ANTHROPIC_API_KEY=sk-ant-... + python test_local.py +""" + +import json +import sys +from pathlib import Path + +# Add project root to path so we can import the handler's parsing logic +sys.path.insert(0, str(Path(__file__).parent / "lambdas" / "email_processor")) + +from handler import parse_raw_email, extract_with_claude + + +def main(): + samples_dir = Path(__file__).parent / "samples" + eml_files = list(samples_dir.glob("*.eml")) + + if not eml_files: + print("No .eml files found in samples/") + return + + for eml_path in eml_files: + print(f"\n{'='*80}") + print(f"FILE: {eml_path.name}") + print(f"{'='*80}") + + raw = eml_path.read_bytes() + email_data = parse_raw_email(raw) + + print(f"Subject: {email_data['subject']}") + print(f"From: {email_data['sender']}") + print(f"Date: {email_data['date']}") + print(f"Body preview: {email_data['body'][:200]}...") + print() + + print("Sending to Claude for extraction...") + parsed = extract_with_claude(email_data) + print(json.dumps(parsed, indent=2)) + + +if __name__ == "__main__": + main()