From e7da56270a8b8d05b41ad62a27d34ac29b0aaecb Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Thu, 30 Apr 2026 17:40:58 -0400 Subject: [PATCH] Add SAM stack, Fargate listener, and shared modules - SAM template with DynamoDB, two Lambdas (fetch-classify, daily-digest), ECS/Fargate for Socket Mode listener, EventBridge Scheduler - Gmail API integration with OAuth refresh and History API sync - Bedrock Haiku classification (HIGH/NORMAL/LOW) - Slack DM delivery with Block Kit formatting and threading - Socket Mode listener for two-way conversation (placeholder) - Setup script for Gmail OAuth token exchange - Slack app manifest for app creation --- README.md | 92 ++++++++ listener/Dockerfile | 8 + listener/app.py | 76 +++++++ listener/requirements.txt | 3 + samconfig.toml.example | 9 + scripts/get_gmail_token.py | 49 ++++ slack-app-manifest.yaml | 31 +++ src/daily_digest/__init__.py | 0 src/daily_digest/app.py | 72 ++++++ src/fetch_classify/__init__.py | 0 src/fetch_classify/app.py | 76 +++++++ src/requirements.txt | 3 + src/shared/__init__.py | 0 src/shared/classify.py | 72 ++++++ src/shared/dynamo.py | 157 +++++++++++++ src/shared/gmail.py | 167 ++++++++++++++ src/shared/secrets.py | 72 ++++++ src/shared/slack.py | 224 +++++++++++++++++++ template.yaml | 397 +++++++++++++++++++++++++++++++++ 19 files changed, 1508 insertions(+) create mode 100644 README.md create mode 100644 listener/Dockerfile create mode 100644 listener/app.py create mode 100644 listener/requirements.txt create mode 100644 samconfig.toml.example create mode 100644 scripts/get_gmail_token.py create mode 100644 slack-app-manifest.yaml create mode 100644 src/daily_digest/__init__.py create mode 100644 src/daily_digest/app.py create mode 100644 src/fetch_classify/__init__.py create mode 100644 src/fetch_classify/app.py create mode 100644 src/requirements.txt create mode 100644 src/shared/__init__.py create mode 100644 src/shared/classify.py create mode 100644 src/shared/dynamo.py create mode 100644 src/shared/gmail.py create mode 100644 src/shared/secrets.py create mode 100644 src/shared/slack.py create mode 100644 template.yaml diff --git a/README.md b/README.md new file mode 100644 index 0000000..824a7b5 --- /dev/null +++ b/README.md @@ -0,0 +1,92 @@ +# Exec Aide + +Personal AI agent that monitors Adam Moussa's Gmail inbox, classifies emails by urgency using Claude Haiku (Bedrock), and delivers Slack alerts and daily digests. + +## Architecture + +``` +EventBridge (15 min) → fetch-classify Lambda → Gmail API → Bedrock Haiku → DynamoDB + Slack DM +EventBridge Scheduler (5 PM ET M-F) → daily-digest Lambda → DynamoDB + Gmail API → Slack DM +``` + +- **exec-aide-fetch-classify**: Polls Gmail via History API, classifies each new message, sends immediate Slack DM for HIGH priority items, tracks thread state for unanswered detection. +- **exec-aide-daily-digest**: Summarizes the day's high-priority items, bypassed work orders (emails not CC'ing work-orders@), and unanswered threads older than 24h. +- **DynamoDB table `exec-aide`**: Single-table design with message records, thread tracking, and sync metadata. 90-day TTL. + +## Classification Rules + +| Category | Trigger | +|---|---| +| **HIGH** | VIP sender, urgency keywords, escalation language, financial/legal matters | +| **NORMAL** | Standard business correspondence | +| **LOW** | Newsletters, automated notifications, marketing | +| **Bypassed WO** | Email to Adam without work-orders@ in TO/CC | +| **Unanswered** | Thread where Adam hasn't replied in 24h+ | + +## Prerequisites + +- AWS CLI + SAM CLI +- Python 3.12 +- Gmail API enabled in Google Cloud Console +- Slack app with `chat:write` + `im:write` scopes + +## One-Time Setup + +### 1. Gmail OAuth + +```bash +# Install the setup dependency (not needed in Lambda) +pip install google-auth-oauthlib + +# Download OAuth client JSON from Google Cloud Console, then: +python scripts/get_gmail_token.py --client-secrets-file path/to/client_secret.json +# Authorize as adam@seahavenind.com in the browser +# Copy the output into Secrets Manager +``` + +### 2. Secrets Manager + +```bash +# Gmail OAuth (output from step 1) +aws secretsmanager create-secret --name exec-aide/gmail-oauth --secret-string '...' + +# Slack bot token +aws secretsmanager create-secret --name exec-aide/slack-bot-token --secret-string 'xoxb-...' +``` + +### 3. SSM Parameters + +```bash +aws ssm put-parameter --name /exec-aide/adam-email --value adam@seahavenind.com --type String +aws ssm put-parameter --name /exec-aide/adam-slack-user-id --value U01XXXXXXXX --type String +aws ssm put-parameter --name /exec-aide/vip-senders --value '["ceo@example.com"]' --type String +aws ssm put-parameter --name /exec-aide/vip-domains --value '["example.com"]' --type String +aws ssm put-parameter --name /exec-aide/work-order-addresses --value '["work-orders@seahavenind.com","work-orders@seahaven.com"]' --type String +aws ssm put-parameter --name /exec-aide/unanswered-threshold-hours --value 24 --type String +``` + +## Deploy + +```bash +cp samconfig.toml.example samconfig.toml # edit if needed +sam build +sam deploy +``` + +## Manual Testing + +```bash +# Trigger a poll +aws lambda invoke --function-name exec-aide-fetch-classify /dev/stdout + +# Trigger the daily digest +aws lambda invoke --function-name exec-aide-daily-digest /dev/stdout +``` + +## Updating VIP List or Keywords + +Edit the SSM parameters directly — changes take effect on the next Lambda cold start (or within ~15 minutes): + +```bash +aws ssm put-parameter --name /exec-aide/vip-senders --value '["new@example.com"]' --type String --overwrite +``` diff --git a/listener/Dockerfile b/listener/Dockerfile new file mode 100644 index 0000000..28df66f --- /dev/null +++ b/listener/Dockerfile @@ -0,0 +1,8 @@ +FROM python:3.12-slim + +WORKDIR /app +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt +COPY app.py . + +CMD ["python", "app.py"] diff --git a/listener/app.py b/listener/app.py new file mode 100644 index 0000000..5d3688d --- /dev/null +++ b/listener/app.py @@ -0,0 +1,76 @@ +import json +import logging +import os +import time + +import boto3 +from slack_bolt import App +from slack_bolt.adapter.socket_mode import SocketModeHandler + +logging.basicConfig(level=logging.INFO) +logger = logging.getLogger("exec-aide") + +secrets_client = boto3.client("secretsmanager") +ssm_client = boto3.client("ssm") + +_config = None + + +def get_slack_credentials(): + resp = secrets_client.get_secret_value(SecretId=os.environ["SECRET_SLACK"]) + return json.loads(resp["SecretString"]) + + +def get_config(): + global _config + if _config is not None: + return _config + + prefix = os.environ["SSM_PREFIX"] + params = {} + paginator = ssm_client.get_paginator("get_parameters_by_path") + for page in paginator.paginate(Path=prefix, Recursive=True): + for p in page["Parameters"]: + key = p["Name"].split("/")[-1] + params[key] = p["Value"] + + _config = { + "adam_email": params.get("adam-email", ""), + "adam_slack_user_id": params.get("adam-slack-user-id", ""), + } + return _config + + +creds = get_slack_credentials() +app = App(token=creds["botToken"]) + + +@app.event("message") +def handle_message(event, say): + if event.get("bot_id") or event.get("subtype"): + return + + config = get_config() + if event.get("user") != config["adam_slack_user_id"]: + return + + text = event.get("text", "").strip() + thread_ts = event.get("thread_ts") or event["ts"] + + response = ( + "I received your message. Right now I'm focused on monitoring your inbox — " + "I'll send you real-time alerts for high-priority emails and a daily digest at 5 PM.\n\n" + "More conversational features are coming soon." + ) + + say(text=response, thread_ts=thread_ts) + + +def main(): + handler = SocketModeHandler(app, creds["appToken"]) + logger.info("Exec Aide listener starting (Socket Mode)") + handler.start() + + +if __name__ == "__main__": + main() diff --git a/listener/requirements.txt b/listener/requirements.txt new file mode 100644 index 0000000..5e1612b --- /dev/null +++ b/listener/requirements.txt @@ -0,0 +1,3 @@ +slack-bolt>=1.18 +slack-sdk>=3.27 +boto3>=1.34 diff --git a/samconfig.toml.example b/samconfig.toml.example new file mode 100644 index 0000000..4241e92 --- /dev/null +++ b/samconfig.toml.example @@ -0,0 +1,9 @@ +version = 0.1 + +[default.deploy.parameters] +stack_name = "exec-aide" +resolve_s3 = true +s3_prefix = "exec-aide" +region = "us-east-1" +confirm_changeset = true +capabilities = "CAPABILITY_NAMED_IAM" diff --git a/scripts/get_gmail_token.py b/scripts/get_gmail_token.py new file mode 100644 index 0000000..46b5c5d --- /dev/null +++ b/scripts/get_gmail_token.py @@ -0,0 +1,49 @@ +#!/usr/bin/env python3 +"""One-time OAuth token exchange for Gmail API access. + +Usage: + 1. Download OAuth client JSON from Google Cloud Console + 2. Run: python scripts/get_gmail_token.py --client-secrets-file path/to/client_secret.json + 3. Authorize as adam@seahavenind.com in the browser + 4. Copy the output JSON into Secrets Manager: + aws secretsmanager create-secret --name exec-aide/gmail-oauth --secret-string '' +""" +import argparse +import json + +from google_auth_oauthlib.flow import InstalledAppFlow + +SCOPES = ["https://www.googleapis.com/auth/gmail.readonly"] + + +def main(): + parser = argparse.ArgumentParser(description="Get Gmail OAuth refresh token") + parser.add_argument( + "--client-secrets-file", + required=True, + help="Path to the OAuth client secrets JSON downloaded from Google Cloud Console", + ) + args = parser.parse_args() + + flow = InstalledAppFlow.from_client_secrets_file(args.client_secrets_file, SCOPES) + creds = flow.run_local_server(port=8080) + + secret = { + "client_id": creds.client_id, + "client_secret": creds.client_secret, + "refresh_token": creds.refresh_token, + "access_token": creds.token, + "token_expiry": creds.expiry.isoformat() if creds.expiry else "", + } + + print("\nStore this in Secrets Manager as exec-aide/gmail-oauth:\n") + print(json.dumps(secret, indent=2)) + print( + "\nCommand:\n" + f"aws secretsmanager create-secret --name exec-aide/gmail-oauth " + f"--secret-string '{json.dumps(secret)}'" + ) + + +if __name__ == "__main__": + main() diff --git a/slack-app-manifest.yaml b/slack-app-manifest.yaml new file mode 100644 index 0000000..8c098d3 --- /dev/null +++ b/slack-app-manifest.yaml @@ -0,0 +1,31 @@ +display_information: + name: Exec Aide + description: Adam's personal executive assistant — inbox monitoring, triage, and more. + background_color: "#1a1a2e" + +features: + bot_user: + display_name: Exec Aide + always_online: true + app_home: + home_tab_enabled: false + messages_tab_enabled: true + messages_tab_read_only_enabled: false + +oauth_config: + scopes: + bot: + - chat:write + - im:history + - im:read + - im:write + +settings: + event_subscriptions: + bot_events: + - message.im + org_deploy_enabled: false + socket_mode_enabled: true + token_rotation_enabled: false + interactivity: + is_enabled: false diff --git a/src/daily_digest/__init__.py b/src/daily_digest/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/daily_digest/app.py b/src/daily_digest/app.py new file mode 100644 index 0000000..3c154fd --- /dev/null +++ b/src/daily_digest/app.py @@ -0,0 +1,72 @@ +import logging +from datetime import datetime, timezone + +from shared.secrets import get_config, get_slack_token +from shared.gmail import get_authenticated_service, fetch_thread +from shared.dynamo import get_todays_messages, get_unanswered_threads, save_thread_state +from shared.slack import send_dm, build_daily_digest + +logger = logging.getLogger() +logger.setLevel(logging.INFO) + + +def lambda_handler(event, context): + config = get_config() + service = get_authenticated_service() + + today = datetime.now(timezone.utc).strftime("%Y-%m-%d") + messages = get_todays_messages(today) + + high_items = [m for m in messages if m.get("classification") == "HIGH"] + bypassed = [m for m in messages if m.get("bypassed_wo")] + + unanswered_raw = get_unanswered_threads(config["unanswered_threshold_hours"]) + unanswered = [] + for thread in unanswered_raw: + thread_id = thread["pk"].replace("THD#", "") + try: + thread_msgs = fetch_thread(service, thread_id) + adam_replied_after = any( + m["from_email"] == config["adam_email"] + and m["internal_date"] > _iso_to_epoch_ms(thread.get("unanswered_since", "")) + for m in thread_msgs + ) + if adam_replied_after: + save_thread_state(thread_id, thread_msgs[-1], config["adam_email"]) + else: + unanswered.append(thread) + except Exception: + logger.exception("Failed to verify thread %s", thread_id) + unanswered.append(thread) + + normal_count = sum(1 for m in messages if m.get("classification") == "NORMAL") + low_count = sum(1 for m in messages if m.get("classification") == "LOW") + + blocks = build_daily_digest(high_items, bypassed, unanswered, normal_count, low_count) + send_dm( + config["adam_slack_user_id"], + blocks, + "Daily Email Digest", + get_slack_token(), + ) + + result = { + "sent": True, + "high": len(high_items), + "bypassed": len(bypassed), + "unanswered": len(unanswered), + "normal": normal_count, + "low": low_count, + } + logger.info("Digest sent: %s", result) + return result + + +def _iso_to_epoch_ms(iso_str): + if not iso_str: + return 0 + try: + dt = datetime.fromisoformat(iso_str) + return int(dt.timestamp() * 1000) + except (ValueError, TypeError): + return 0 diff --git a/src/fetch_classify/__init__.py b/src/fetch_classify/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/fetch_classify/app.py b/src/fetch_classify/app.py new file mode 100644 index 0000000..a9edac9 --- /dev/null +++ b/src/fetch_classify/app.py @@ -0,0 +1,76 @@ +import logging +from datetime import datetime, timezone + +from shared.secrets import get_config, get_slack_token +from shared.gmail import get_authenticated_service, fetch_history, fetch_message, initial_sync +from shared.dynamo import message_exists, save_message, save_thread_state, get_sync_metadata, save_sync_metadata +from shared.classify import classify_email, check_bypassed_workorder +from shared.slack import send_dm, build_high_priority_alert + +logger = logging.getLogger() +logger.setLevel(logging.INFO) + + +def lambda_handler(event, context): + config = get_config() + service = get_authenticated_service() + + sync_meta = get_sync_metadata() + + if sync_meta is None: + logger.info("No sync metadata found, running initial sync") + msg_ids, history_id = initial_sync(service) + else: + msg_ids, history_id = fetch_history(service, sync_meta["history_id"]) + if msg_ids is None: + logger.warning("History expired, falling back to initial sync") + msg_ids, history_id = initial_sync(service) + + today = datetime.now(timezone.utc).strftime("%Y-%m-%d") + results = {"processed": 0, "high": 0, "skipped": 0, "errors": 0} + + for msg_id in msg_ids: + try: + if message_exists(msg_id): + results["skipped"] += 1 + continue + + msg = fetch_message(service, msg_id) + classification = classify_email( + msg, config["vip_senders"], config["vip_domains"] + ) + bypassed = check_bypassed_workorder(msg, config["work_order_addresses"]) + + slack_ts = "" + if classification["classification"] == "HIGH": + blocks = build_high_priority_alert({ + **msg, + "classification_reason": classification["reason"], + }) + slack_ts = send_dm( + config["adam_slack_user_id"], + blocks, + f"HIGH: {msg['subject']}", + get_slack_token(), + ) or "" + results["high"] += 1 + + save_message({ + **msg, + "classification": classification["classification"], + "classification_reason": classification["reason"], + "bypassed_wo": bypassed, + "classified_date": today, + "notified": bool(slack_ts), + "slack_ts": slack_ts, + }) + save_thread_state(msg["thread_id"], msg, config["adam_email"]) + + results["processed"] += 1 + except Exception: + logger.exception("Failed to process message %s", msg_id) + results["errors"] += 1 + + save_sync_metadata(history_id) + logger.info("Poll complete: %s", results) + return results diff --git a/src/requirements.txt b/src/requirements.txt new file mode 100644 index 0000000..871567c --- /dev/null +++ b/src/requirements.txt @@ -0,0 +1,3 @@ +google-api-python-client>=2.0 +google-auth>=2.0 +requests>=2.31 diff --git a/src/shared/__init__.py b/src/shared/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/shared/classify.py b/src/shared/classify.py new file mode 100644 index 0000000..fdcef9b --- /dev/null +++ b/src/shared/classify.py @@ -0,0 +1,72 @@ +import json +import logging + +import boto3 + +logger = logging.getLogger(__name__) + +_bedrock = boto3.client("bedrock-runtime", region_name="us-east-1") + +MODEL_ID = "us.anthropic.claude-haiku-4-5-20251001" + +SYSTEM_PROMPT = """\ +You are an email triage assistant for Adam Moussa, President of Sea Haven Industries \ +(a facilities maintenance company). + +Classify this email's urgency. Return ONLY valid JSON — no markdown, no explanation: + +{"classification": "HIGH" or "NORMAL" or "LOW", "reason": "one sentence"} + +HIGH: +- Sender is on the VIP list +- Urgency keywords: urgent, asap, emergency, critical, immediate, time-sensitive, deadline today, EOD +- Escalation language: threats, legal mentions, contract termination, safety incidents, regulatory compliance +- Direct asks from C-suite, major clients, or key vendors requiring immediate response +- Financial approvals over $10,000 + +NORMAL: +- Standard business correspondence requiring response but not time-critical +- Meeting invites, project updates, routine vendor communications + +LOW: +- Newsletters, marketing, automated notifications, bulk mailings +- FYI-only messages, calendar confirmations, read receipts\ +""" + + +def classify_email(message, vip_senders, vip_domains): + sender_email = message["from_email"] + sender_domain = sender_email.split("@")[-1] if "@" in sender_email else "" + is_vip = sender_email in vip_senders or sender_domain in vip_domains + + user_msg = ( + f"From: {message['from_name']} <{sender_email}>" + f"{' [VIP SENDER]' if is_vip else ''}\n" + f"Subject: {message['subject']}\n" + f"Preview: {message['snippet']}" + ) + + try: + resp = _bedrock.converse( + modelId=MODEL_ID, + system=[{"text": SYSTEM_PROMPT}], + messages=[{"role": "user", "content": [{"text": user_msg}]}], + inferenceConfig={"temperature": 0.0, "maxTokens": 256}, + ) + raw = resp["output"]["message"]["content"][0]["text"] + result = json.loads(raw) + return { + "classification": result.get("classification", "NORMAL"), + "reason": result.get("reason", ""), + } + except Exception: + logger.exception("Classification failed, defaulting to NORMAL") + return {"classification": "NORMAL", "reason": "Classification unavailable"} + + +def check_bypassed_workorder(message, work_order_addresses): + all_recipients = [ + addr.lower() + for addr in message["to_emails"] + message["cc_emails"] + ] + return not any(wo in all_recipients for wo in work_order_addresses) diff --git a/src/shared/dynamo.py b/src/shared/dynamo.py new file mode 100644 index 0000000..6937cf3 --- /dev/null +++ b/src/shared/dynamo.py @@ -0,0 +1,157 @@ +import os +import time +from datetime import datetime, timezone + +import boto3 +from botocore.exceptions import ClientError + +_table = None + + +def _get_table(): + global _table + if _table is None: + _table = boto3.resource("dynamodb").Table(os.environ["TABLE_NAME"]) + return _table + + +def _ttl_90_days(): + return int(time.time()) + (90 * 24 * 60 * 60) + + +def message_exists(message_id): + resp = _get_table().get_item( + Key={"pk": f"MSG#{message_id}", "sk": "MSG"}, + ProjectionExpression="pk", + ) + return "Item" in resp + + +def save_message(msg_data): + item = { + "pk": f"MSG#{msg_data['message_id']}", + "sk": "MSG", + "thread_id": msg_data["thread_id"], + "from_email": msg_data["from_email"], + "from_name": msg_data["from_name"], + "to_emails": msg_data["to_emails"], + "cc_emails": msg_data["cc_emails"], + "subject": msg_data["subject"], + "snippet": msg_data["snippet"], + "internal_date": msg_data["internal_date"], + "classification": msg_data["classification"], + "classification_reason": msg_data["classification_reason"], + "bypassed_wo": msg_data["bypassed_wo"], + "notified": msg_data.get("notified", False), + "classified_date": msg_data["classified_date"], + "slack_ts": msg_data.get("slack_ts", ""), + "ttl": _ttl_90_days(), + } + try: + _get_table().put_item( + Item=item, + ConditionExpression="attribute_not_exists(pk)", + ) + except ClientError as e: + if e.response["Error"]["Code"] == "ConditionalCheckFailedException": + return + raise + + +def save_thread_state(thread_id, message, adam_email): + now_iso = datetime.now(timezone.utc).isoformat() + is_adam = message["from_email"] == adam_email + + if is_adam: + _get_table().update_item( + Key={"pk": f"THD#{thread_id}", "sk": "THD"}, + UpdateExpression=( + "SET last_message_from = :from_email, " + "last_message_date = :now, " + "subject = :subject, " + "adam_last_reply_date = :now, " + "ttl = :ttl " + "REMOVE unanswered_since" + ), + ExpressionAttributeValues={ + ":from_email": message["from_email"], + ":now": now_iso, + ":subject": message["subject"], + ":ttl": _ttl_90_days(), + }, + ) + else: + _get_table().update_item( + Key={"pk": f"THD#{thread_id}", "sk": "THD"}, + UpdateExpression=( + "SET last_message_from = :from_email, " + "last_message_date = :now, " + "subject = :subject, " + "unanswered_since = if_not_exists(unanswered_since, :now), " + "ttl = :ttl" + ), + ExpressionAttributeValues={ + ":from_email": message["from_email"], + ":now": now_iso, + ":subject": message["subject"], + ":ttl": _ttl_90_days(), + }, + ) + + +def get_sync_metadata(): + resp = _get_table().get_item( + Key={"pk": "META#sync", "sk": "META"} + ) + return resp.get("Item") + + +def save_sync_metadata(history_id): + _get_table().put_item( + Item={ + "pk": "META#sync", + "sk": "META", + "history_id": str(history_id), + "last_sync_at": datetime.now(timezone.utc).isoformat(), + "initial_sync_done": True, + } + ) + + +def get_todays_messages(date_str): + items = [] + resp = _get_table().query( + IndexName="by-date", + KeyConditionExpression="classified_date = :d AND sk = :sk", + ExpressionAttributeValues={":d": date_str, ":sk": "MSG"}, + ) + items.extend(resp.get("Items", [])) + while "LastEvaluatedKey" in resp: + resp = _get_table().query( + IndexName="by-date", + KeyConditionExpression="classified_date = :d AND sk = :sk", + ExpressionAttributeValues={":d": date_str, ":sk": "MSG"}, + ExclusiveStartKey=resp["LastEvaluatedKey"], + ) + items.extend(resp.get("Items", [])) + return items + + +def get_unanswered_threads(threshold_hours): + cutoff = datetime.now(timezone.utc).timestamp() - (threshold_hours * 3600) + items = [] + scan_kwargs = { + "FilterExpression": "begins_with(pk, :prefix) AND attribute_exists(unanswered_since)", + "ExpressionAttributeValues": {":prefix": "THD#"}, + } + resp = _get_table().scan(**scan_kwargs) + items.extend(resp.get("Items", [])) + while "LastEvaluatedKey" in resp: + scan_kwargs["ExclusiveStartKey"] = resp["LastEvaluatedKey"] + resp = _get_table().scan(**scan_kwargs) + items.extend(resp.get("Items", [])) + + return [ + item for item in items + if datetime.fromisoformat(item["unanswered_since"]).timestamp() < cutoff + ] diff --git a/src/shared/gmail.py b/src/shared/gmail.py new file mode 100644 index 0000000..b621da2 --- /dev/null +++ b/src/shared/gmail.py @@ -0,0 +1,167 @@ +import logging +import re +from datetime import datetime, timezone + +import google.auth.transport.requests +from google.oauth2.credentials import Credentials +from googleapiclient.discovery import build +from googleapiclient.errors import HttpError + +from shared.secrets import get_gmail_oauth, save_gmail_tokens + +logger = logging.getLogger(__name__) + +_service = None + +SCOPES = ["https://www.googleapis.com/auth/gmail.readonly"] + + +def get_authenticated_service(): + global _service + oauth = get_gmail_oauth() + creds = Credentials( + token=oauth.get("access_token"), + refresh_token=oauth["refresh_token"], + token_uri="https://oauth2.googleapis.com/token", + client_id=oauth["client_id"], + client_secret=oauth["client_secret"], + scopes=SCOPES, + ) + if not creds.valid: + creds.refresh(google.auth.transport.requests.Request()) + save_gmail_tokens( + creds.token, + creds.expiry.isoformat() if creds.expiry else "", + ) + _service = build("gmail", "v1", credentials=creds, cache_discovery=False) + return _service + + +def fetch_history(service, history_id): + message_ids = [] + new_history_id = history_id + try: + page_token = None + while True: + resp = ( + service.users() + .history() + .list( + userId="me", + startHistoryId=history_id, + historyTypes=["messageAdded"], + labelIds=["INBOX"], + pageToken=page_token, + ) + .execute() + ) + new_history_id = resp.get("historyId", new_history_id) + for record in resp.get("history", []): + for added in record.get("messagesAdded", []): + message_ids.append(added["message"]["id"]) + page_token = resp.get("nextPageToken") + if not page_token: + break + except HttpError as e: + if e.resp.status == 404: + logger.warning("History ID expired, signaling re-sync") + return None, None + raise + return message_ids, new_history_id + + +def fetch_message(service, message_id): + msg = ( + service.users() + .messages() + .get( + userId="me", + id=message_id, + format="metadata", + metadataHeaders=["From", "To", "Cc", "Subject", "Date"], + ) + .execute() + ) + return _parse_message(msg) + + +def fetch_thread(service, thread_id): + thread = ( + service.users() + .threads() + .get(userId="me", id=thread_id, format="metadata") + .execute() + ) + return [_parse_message(m) for m in thread.get("messages", [])] + + +def initial_sync(service): + message_ids = [] + page_token = None + while True: + resp = ( + service.users() + .messages() + .list( + userId="me", + q="newer_than:2d", + labelIds=["INBOX"], + pageToken=page_token, + ) + .execute() + ) + for m in resp.get("messages", []): + message_ids.append(m["id"]) + page_token = resp.get("nextPageToken") + if not page_token: + break + + if not message_ids: + profile = service.users().getProfile(userId="me").execute() + return [], str(profile["historyId"]) + + first_msg = ( + service.users() + .messages() + .get(userId="me", id=message_ids[0], format="minimal") + .execute() + ) + return message_ids, str(first_msg["historyId"]) + + +def _parse_message(msg): + headers = {h["name"]: h["value"] for h in msg.get("payload", {}).get("headers", [])} + from_raw = headers.get("From", "") + from_name, from_email = _parse_from(from_raw) + return { + "message_id": msg["id"], + "thread_id": msg["threadId"], + "from_email": from_email, + "from_name": from_name, + "to_emails": _parse_address_list(headers.get("To", "")), + "cc_emails": _parse_address_list(headers.get("Cc", "")), + "subject": headers.get("Subject", "(no subject)"), + "snippet": msg.get("snippet", ""), + "internal_date": int(msg.get("internalDate", 0)), + } + + +def _parse_from(raw): + match = re.match(r"^(.+?)\s*<(.+?)>$", raw) + if match: + return match.group(1).strip().strip('"'), match.group(2).lower() + return "", raw.strip().lower() + + +def _parse_address_list(raw): + if not raw: + return [] + addresses = [] + for part in raw.split(","): + part = part.strip() + match = re.match(r".*<(.+?)>", part) + if match: + addresses.append(match.group(1).lower()) + elif "@" in part: + addresses.append(part.lower()) + return addresses diff --git a/src/shared/secrets.py b/src/shared/secrets.py new file mode 100644 index 0000000..4e4e6f1 --- /dev/null +++ b/src/shared/secrets.py @@ -0,0 +1,72 @@ +import json +import os + +import boto3 + +_secrets_client = boto3.client("secretsmanager") +_ssm_client = boto3.client("ssm") + +_gmail_oauth = None +_slack_creds = None +_config = None + + +def get_gmail_oauth(): + global _gmail_oauth + if _gmail_oauth is None: + resp = _secrets_client.get_secret_value( + SecretId=os.environ["SECRET_GMAIL"] + ) + _gmail_oauth = json.loads(resp["SecretString"]) + return _gmail_oauth + + +def save_gmail_tokens(access_token, token_expiry): + global _gmail_oauth + oauth = get_gmail_oauth() + oauth["access_token"] = access_token + oauth["token_expiry"] = token_expiry + _secrets_client.put_secret_value( + SecretId=os.environ["SECRET_GMAIL"], + SecretString=json.dumps(oauth), + ) + _gmail_oauth = oauth + + +def get_slack_credentials(): + """Returns { botToken, signingSecret, appToken }.""" + global _slack_creds + if _slack_creds is None: + resp = _secrets_client.get_secret_value( + SecretId=os.environ["SECRET_SLACK"] + ) + _slack_creds = json.loads(resp["SecretString"]) + return _slack_creds + + +def get_slack_token(): + return get_slack_credentials()["botToken"] + + +def get_config(): + global _config + if _config is not None: + return _config + + prefix = os.environ["SSM_PREFIX"] + params = {} + paginator = _ssm_client.get_paginator("get_parameters_by_path") + for page in paginator.paginate(Path=prefix, Recursive=True): + for p in page["Parameters"]: + key = p["Name"].split("/")[-1] + params[key] = p["Value"] + + _config = { + "adam_email": params.get("adam-email", ""), + "adam_slack_user_id": params.get("adam-slack-user-id", ""), + "vip_senders": json.loads(params.get("vip-senders", "[]")), + "vip_domains": json.loads(params.get("vip-domains", "[]")), + "work_order_addresses": json.loads(params.get("work-order-addresses", "[]")), + "unanswered_threshold_hours": int(params.get("unanswered-threshold-hours", "24")), + } + return _config diff --git a/src/shared/slack.py b/src/shared/slack.py new file mode 100644 index 0000000..a0fb090 --- /dev/null +++ b/src/shared/slack.py @@ -0,0 +1,224 @@ +import json +import logging +from datetime import datetime, timezone + +import requests + +logger = logging.getLogger(__name__) + +_dm_channel_id = None + + +def send_dm(user_id, blocks, text, token): + global _dm_channel_id + if _dm_channel_id is None: + resp = requests.post( + "https://slack.com/api/conversations.open", + headers={"Authorization": f"Bearer {token}"}, + json={"users": user_id}, + ) + resp.raise_for_status() + data = resp.json() + if not data.get("ok"): + raise RuntimeError(f"conversations.open failed: {data.get('error')}") + _dm_channel_id = data["channel"]["id"] + + resp = requests.post( + "https://slack.com/api/chat.postMessage", + headers={"Authorization": f"Bearer {token}"}, + json={ + "channel": _dm_channel_id, + "blocks": blocks, + "text": text, + }, + ) + resp.raise_for_status() + data = resp.json() + if not data.get("ok"): + raise RuntimeError(f"chat.postMessage failed: {data.get('error')}") + return data.get("ts") + + +def build_high_priority_alert(message): + time_ago = _format_time_ago(message["internal_date"]) + gmail_url = f"https://mail.google.com/mail/u/0/#inbox/{message['message_id']}" + + return [ + { + "type": "header", + "text": {"type": "plain_text", "text": ":rotating_light: High-Priority Email"}, + }, + { + "type": "section", + "fields": [ + {"type": "mrkdwn", "text": f"*From*\n{message['from_name'] or message['from_email']}"}, + {"type": "mrkdwn", "text": f"*Received*\n{time_ago}"}, + ], + }, + { + "type": "section", + "text": {"type": "mrkdwn", "text": f"*{message['subject']}*\n{message['snippet'][:300]}"}, + }, + { + "type": "context", + "elements": [ + {"type": "mrkdwn", "text": f":brain: _{message.get('classification_reason', '')}_"}, + ], + }, + { + "type": "actions", + "elements": [ + { + "type": "button", + "text": {"type": "plain_text", "text": "Open in Gmail"}, + "url": gmail_url, + "action_id": "open_gmail", + } + ], + }, + ] + + +def build_daily_digest(high_items, bypassed, unanswered, normal_count, low_count): + today_str = datetime.now(timezone.utc).strftime("%A, %B %d") + total = len(high_items) + normal_count + low_count + blocks = [ + { + "type": "header", + "text": {"type": "plain_text", "text": ":newspaper: Daily Inbox Digest"}, + }, + { + "type": "context", + "elements": [ + {"type": "mrkdwn", "text": f"{today_str} · {total} emails processed today"}, + ], + }, + {"type": "divider"}, + ] + + if high_items: + blocks.append({ + "type": "section", + "text": {"type": "mrkdwn", "text": f":rotating_light: *High Priority ({len(high_items)})*"}, + }) + for item in high_items: + time_ago = _format_time_ago(item["internal_date"]) + gmail_url = f"https://mail.google.com/mail/u/0/#inbox/{item['pk'].replace('MSG#', '')}" + blocks.append({ + "type": "section", + "text": { + "type": "mrkdwn", + "text": ( + f"*{item['subject']}*\n" + f"From: {item.get('from_name') or item['from_email']} · {time_ago}\n" + f"_{item.get('classification_reason', '')}_" + ), + }, + "accessory": { + "type": "button", + "text": {"type": "plain_text", "text": "Open"}, + "url": gmail_url, + "action_id": f"open_{item['pk']}", + }, + }) + blocks.append({"type": "divider"}) + + if bypassed: + blocks.append({ + "type": "section", + "text": { + "type": "mrkdwn", + "text": ( + f":warning: *Bypassed Work Orders ({len(bypassed)})*\n" + "_Sent directly to you without CC'ing work-orders@_" + ), + }, + }) + for item in bypassed: + blocks.append({ + "type": "section", + "text": { + "type": "mrkdwn", + "text": f" · *{item['subject']}* from {item.get('from_name') or item['from_email']}", + }, + }) + blocks.append({"type": "divider"}) + + if unanswered: + blocks.append({ + "type": "section", + "text": { + "type": "mrkdwn", + "text": ( + f":hourglass: *Unanswered Threads ({len(unanswered)})*\n" + "_Awaiting your reply for 24h+_" + ), + }, + }) + for thread in unanswered: + hours = _hours_since_iso(thread.get("unanswered_since", "")) + blocks.append({ + "type": "section", + "text": { + "type": "mrkdwn", + "text": ( + f" · *{thread.get('subject', '(no subject)')}*" + f" — last from {thread.get('last_message_from', 'unknown')}" + f" · {hours}h waiting" + ), + }, + }) + blocks.append({"type": "divider"}) + + blocks.append({ + "type": "context", + "elements": [ + { + "type": "mrkdwn", + "text": ( + f":bar_chart: Today: {len(high_items)} high · " + f"{normal_count} normal · {low_count} low · " + f"{len(bypassed)} bypassed WO · {len(unanswered)} unanswered" + ), + } + ], + }) + + return blocks + + +def send_message_to_channel(channel, text, token): + resp = requests.post( + "https://slack.com/api/chat.postMessage", + headers={"Authorization": f"Bearer {token}"}, + json={"channel": channel, "text": text}, + ) + resp.raise_for_status() + data = resp.json() + if not data.get("ok"): + logger.error("chat.postMessage failed: %s", data.get("error")) + + +def _format_time_ago(internal_date): + if not internal_date: + return "unknown" + msg_time = datetime.fromtimestamp(int(internal_date) / 1000, tz=timezone.utc) + delta = datetime.now(timezone.utc) - msg_time + hours = int(delta.total_seconds() / 3600) + if hours < 1: + mins = int(delta.total_seconds() / 60) + return f"{mins}m ago" + if hours < 24: + return f"{hours}h ago" + return f"{hours // 24}d ago" + + +def _hours_since_iso(iso_str): + if not iso_str: + return 0 + try: + dt = datetime.fromisoformat(iso_str) + delta = datetime.now(timezone.utc) - dt + return int(delta.total_seconds() / 3600) + except (ValueError, TypeError): + return 0 diff --git a/template.yaml b/template.yaml new file mode 100644 index 0000000..58e26e0 --- /dev/null +++ b/template.yaml @@ -0,0 +1,397 @@ +AWSTemplateFormatVersion: '2010-09-09' +Transform: AWS::Serverless-2016-10-31 +Description: > + exec-aide — Gmail inbox monitor with AI classification and Slack alerts + +Globals: + Function: + Runtime: python3.12 + Architectures: + - arm64 + MemorySize: 256 + Timeout: 120 + Environment: + Variables: + TABLE_NAME: !Ref ExecAideTable + SECRET_GMAIL: exec-aide/gmail-oauth + SECRET_SLACK: exec-aide/slack-credentials + SSM_PREFIX: /exec-aide + +Resources: + + # ── DynamoDB ────────────────────────────────────────────── + + ExecAideTable: + Type: AWS::DynamoDB::Table + Properties: + TableName: exec-aide + BillingMode: PAY_PER_REQUEST + AttributeDefinitions: + - AttributeName: pk + AttributeType: S + - AttributeName: sk + AttributeType: S + - AttributeName: classified_date + AttributeType: S + KeySchema: + - AttributeName: pk + KeyType: HASH + - AttributeName: sk + KeyType: RANGE + GlobalSecondaryIndexes: + - IndexName: by-date + KeySchema: + - AttributeName: classified_date + KeyType: HASH + - AttributeName: sk + KeyType: RANGE + Projection: + ProjectionType: ALL + TimeToLiveSpecification: + AttributeName: ttl + Enabled: true + + # ── CloudWatch Log Groups ──────────────────────────────── + + FetchClassifyLogGroup: + Type: AWS::Logs::LogGroup + Properties: + LogGroupName: /aws/lambda/exec-aide-fetch-classify + RetentionInDays: 60 + + DailyDigestLogGroup: + Type: AWS::Logs::LogGroup + Properties: + LogGroupName: /aws/lambda/exec-aide-daily-digest + RetentionInDays: 60 + + # ── Lambda Functions ───────────────────────────────────── + + FetchClassifyFunction: + Type: AWS::Serverless::Function + DependsOn: FetchClassifyLogGroup + Properties: + FunctionName: exec-aide-fetch-classify + Handler: fetch_classify.app.lambda_handler + CodeUri: src/ + Events: + PollSchedule: + Type: Schedule + Properties: + Schedule: rate(15 minutes) + Description: Poll Gmail for new messages + Enabled: true + Policies: + - DynamoDBCrudPolicy: + TableName: !Ref ExecAideTable + - Statement: + - Effect: Allow + Action: + - secretsmanager:GetSecretValue + Resource: + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/gmail-oauth-* + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/slack-credentials-* + - Effect: Allow + Action: + - secretsmanager:PutSecretValue + Resource: + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/gmail-oauth-* + - Effect: Allow + Action: + - ssm:GetParametersByPath + - ssm:GetParameter + Resource: + - !Sub arn:aws:ssm:${AWS::Region}:${AWS::AccountId}:parameter/exec-aide/* + - Effect: Allow + Action: + - bedrock:InvokeModel + Resource: + - arn:aws:bedrock:*::foundation-model/anthropic.* + - !Sub arn:aws:bedrock:${AWS::Region}:${AWS::AccountId}:inference-profile/us.anthropic.* + + DailyDigestFunction: + Type: AWS::Serverless::Function + DependsOn: DailyDigestLogGroup + Properties: + FunctionName: exec-aide-daily-digest + Handler: daily_digest.app.lambda_handler + CodeUri: src/ + Policies: + - DynamoDBCrudPolicy: + TableName: !Ref ExecAideTable + - Statement: + - Effect: Allow + Action: + - secretsmanager:GetSecretValue + Resource: + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/gmail-oauth-* + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/slack-credentials-* + - Effect: Allow + Action: + - secretsmanager:PutSecretValue + Resource: + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/gmail-oauth-* + - Effect: Allow + Action: + - ssm:GetParametersByPath + - ssm:GetParameter + Resource: + - !Sub arn:aws:ssm:${AWS::Region}:${AWS::AccountId}:parameter/exec-aide/* + + # ── ECS / Fargate (Socket Mode listener) ────────────────── + + ListenerLogGroup: + Type: AWS::Logs::LogGroup + Properties: + LogGroupName: /ecs/exec-aide-listener + RetentionInDays: 60 + + EcsCluster: + Type: AWS::ECS::Cluster + Properties: + ClusterName: exec-aide + + ListenerTaskDefinition: + Type: AWS::ECS::TaskDefinition + Properties: + Family: exec-aide-listener + Cpu: "256" + Memory: "512" + NetworkMode: awsvpc + RequiresCompatibilities: + - FARGATE + RuntimePlatform: + CpuArchitecture: ARM64 + OperatingSystemFamily: LINUX + ExecutionRoleArn: !GetAtt ListenerExecutionRole.Arn + TaskRoleArn: !GetAtt ListenerTaskRole.Arn + ContainerDefinitions: + - Name: listener + Image: !Sub ${AWS::AccountId}.dkr.ecr.${AWS::Region}.amazonaws.com/exec-aide-listener:latest + Essential: true + Environment: + - Name: TABLE_NAME + Value: !Ref ExecAideTable + - Name: SECRET_SLACK + Value: exec-aide/slack-credentials + - Name: SSM_PREFIX + Value: /exec-aide + LogConfiguration: + LogDriver: awslogs + Options: + awslogs-group: /ecs/exec-aide-listener + awslogs-region: !Ref AWS::Region + awslogs-stream-prefix: listener + + ListenerService: + Type: AWS::ECS::Service + Properties: + ServiceName: exec-aide-listener + Cluster: !Ref EcsCluster + TaskDefinition: !Ref ListenerTaskDefinition + DesiredCount: 1 + LaunchType: FARGATE + NetworkConfiguration: + AwsvpcConfiguration: + AssignPublicIp: ENABLED + Subnets: + - !Ref ListenerSubnet + SecurityGroups: + - !Ref ListenerSecurityGroup + + ListenerSubnet: + Type: AWS::EC2::Subnet + Properties: + VpcId: !Ref ListenerVpc + CidrBlock: 10.30.0.0/24 + MapPublicIpOnLaunch: true + AvailabilityZone: !Select [0, !GetAZs ""] + Tags: + - Key: Name + Value: exec-aide-listener + + ListenerVpc: + Type: AWS::EC2::VPC + Properties: + CidrBlock: 10.30.0.0/16 + EnableDnsHostnames: true + EnableDnsSupport: true + Tags: + - Key: Name + Value: exec-aide + + ListenerIgw: + Type: AWS::EC2::InternetGateway + Properties: + Tags: + - Key: Name + Value: exec-aide + + ListenerIgwAttachment: + Type: AWS::EC2::VPCGatewayAttachment + Properties: + VpcId: !Ref ListenerVpc + InternetGatewayId: !Ref ListenerIgw + + ListenerRouteTable: + Type: AWS::EC2::RouteTable + Properties: + VpcId: !Ref ListenerVpc + + ListenerRoute: + Type: AWS::EC2::Route + DependsOn: ListenerIgwAttachment + Properties: + RouteTableId: !Ref ListenerRouteTable + DestinationCidrBlock: 0.0.0.0/0 + GatewayId: !Ref ListenerIgw + + ListenerSubnetRouteTableAssoc: + Type: AWS::EC2::SubnetRouteTableAssociation + Properties: + SubnetId: !Ref ListenerSubnet + RouteTableId: !Ref ListenerRouteTable + + ListenerSecurityGroup: + Type: AWS::EC2::SecurityGroup + Properties: + GroupDescription: exec-aide listener — outbound only + VpcId: !Ref ListenerVpc + SecurityGroupEgress: + - IpProtocol: "-1" + CidrIp: 0.0.0.0/0 + + ListenerExecutionRole: + Type: AWS::IAM::Role + Properties: + RoleName: exec-aide-listener-execution + AssumeRolePolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Principal: + Service: ecs-tasks.amazonaws.com + Action: sts:AssumeRole + ManagedPolicyArns: + - arn:aws:iam::aws:policy/service-role/AmazonECSTaskExecutionRolePolicy + + ListenerTaskRole: + Type: AWS::IAM::Role + Properties: + RoleName: exec-aide-listener-task + AssumeRolePolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Principal: + Service: ecs-tasks.amazonaws.com + Action: sts:AssumeRole + Policies: + - PolicyName: ExecAideListenerAccess + PolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Action: + - secretsmanager:GetSecretValue + Resource: + - !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:exec-aide/slack-credentials-* + - Effect: Allow + Action: + - ssm:GetParametersByPath + - ssm:GetParameter + Resource: + - !Sub arn:aws:ssm:${AWS::Region}:${AWS::AccountId}:parameter/exec-aide/* + - Effect: Allow + Action: + - dynamodb:GetItem + - dynamodb:PutItem + - dynamodb:UpdateItem + - dynamodb:Query + Resource: + - !GetAtt ExecAideTable.Arn + - !Sub ${ExecAideTable.Arn}/index/* + + ListenerEcrRepo: + Type: AWS::ECR::Repository + Properties: + RepositoryName: exec-aide-listener + ImageScanningConfiguration: + ScanOnPush: true + LifecyclePolicy: + LifecyclePolicyText: | + { + "rules": [ + { + "rulePriority": 1, + "description": "Keep last 5 images", + "selection": { + "tagStatus": "any", + "countType": "imageCountMoreThan", + "countNumber": 5 + }, + "action": { "type": "expire" } + } + ] + } + + # ── EventBridge Scheduler (daily digest, DST-aware) ────── + + DailyDigestSchedule: + Type: AWS::Scheduler::Schedule + Properties: + Name: exec-aide-daily-digest + Description: Daily 5 PM ET inbox digest + ScheduleExpression: cron(0 17 ? * MON-FRI *) + ScheduleExpressionTimezone: America/New_York + FlexibleTimeWindow: + Mode: "OFF" + State: ENABLED + Target: + Arn: !GetAtt DailyDigestFunction.Arn + RoleArn: !GetAtt DailyDigestSchedulerRole.Arn + + DailyDigestSchedulerRole: + Type: AWS::IAM::Role + Properties: + RoleName: exec-aide-digest-scheduler + AssumeRolePolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Principal: + Service: scheduler.amazonaws.com + Action: sts:AssumeRole + Policies: + - PolicyName: InvokeLambda + PolicyDocument: + Version: "2012-10-17" + Statement: + - Effect: Allow + Action: lambda:InvokeFunction + Resource: !GetAtt DailyDigestFunction.Arn + + DailyDigestSchedulePermission: + Type: AWS::Lambda::Permission + Properties: + FunctionName: !Ref DailyDigestFunction + Action: lambda:InvokeFunction + Principal: scheduler.amazonaws.com + SourceArn: !GetAtt DailyDigestSchedule.Arn + +Outputs: + FetchClassifyFunctionArn: + Description: Fetch & classify Lambda ARN + Value: !GetAtt FetchClassifyFunction.Arn + DailyDigestFunctionArn: + Description: Daily digest Lambda ARN + Value: !GetAtt DailyDigestFunction.Arn + ExecAideTableName: + Description: DynamoDB table name + Value: !Ref ExecAideTable + ListenerEcrRepoUri: + Description: ECR repo for the Socket Mode listener + Value: !GetAtt ListenerEcrRepo.RepositoryUri + EcsClusterName: + Description: ECS cluster name + Value: !Ref EcsCluster