import json import logging import os from datetime import datetime, timezone import boto3 from shared.dynamo import ( get_todays_messages, get_unanswered_threads, save_coordination, ) from shared.gmail import get_authenticated_service, fetch_thread from shared.secrets import get_config, get_slack_token from shared.slack import post_message logger = logging.getLogger(__name__) _lambda_client = None _tool_context = {} def set_tool_context(channel, thread_ts): global _tool_context _tool_context = {"channel": channel, "thread_ts": thread_ts} TOOLS = [ { "toolSpec": { "name": "get_inbox_summary", "description": "Get a summary of today's inbox: total emails processed, counts by classification (HIGH, NORMAL, LOW), and a list of the most recent high-priority items.", "inputSchema": { "json": { "type": "object", "properties": { "date": { "type": "string", "description": "Date in YYYY-MM-DD format. Defaults to today.", } }, } }, } }, { "toolSpec": { "name": "get_high_priority_emails", "description": "Get all HIGH-priority emails for a given date.", "inputSchema": { "json": { "type": "object", "properties": { "date": { "type": "string", "description": "Date in YYYY-MM-DD format. Defaults to today.", } }, } }, } }, { "toolSpec": { "name": "get_unanswered_threads", "description": "Get email threads where Adam hasn't replied in 24+ hours.", "inputSchema": { "json": { "type": "object", "properties": {}, } }, } }, { "toolSpec": { "name": "get_bypassed_work_orders", "description": "Get emails sent directly to Adam without CC'ing the work-orders address. These should have gone through the work order system.", "inputSchema": { "json": { "type": "object", "properties": { "date": { "type": "string", "description": "Date in YYYY-MM-DD format. Defaults to today.", } }, } }, } }, { "toolSpec": { "name": "search_emails_by_sender", "description": "Search for emails from a specific sender by name or email address. Searches across all classified emails in the database.", "inputSchema": { "json": { "type": "object", "properties": { "sender": { "type": "string", "description": "Sender name or email address to search for (case-insensitive partial match).", }, "limit": { "type": "integer", "description": "Maximum number of results to return. Defaults to 10.", }, }, "required": ["sender"], } }, } }, { "toolSpec": { "name": "get_email_thread_detail", "description": "Get the full message history of a specific email thread from Gmail, including all participants and message snippets.", "inputSchema": { "json": { "type": "object", "properties": { "thread_id": { "type": "string", "description": "Gmail thread ID.", } }, "required": ["thread_id"], } }, } }, { "toolSpec": { "name": "trigger_daily_digest", "description": "Send Adam his daily inbox digest right now, regardless of the scheduled time.", "inputSchema": { "json": { "type": "object", "properties": {}, } }, } }, { "toolSpec": { "name": "search_inbox", "description": "Search Adam's Gmail inbox with a query. Supports Gmail search syntax: from:, to:, subject:, after:, before:, is:unread, label:, has:attachment, etc.", "inputSchema": { "json": { "type": "object", "properties": { "query": { "type": "string", "description": "Gmail search query (e.g., 'from:jane after:2024/01/01 subject:invoice').", }, "max_results": { "type": "integer", "description": "Maximum results to return. Defaults to 10.", }, }, "required": ["query"], } }, } }, { "toolSpec": { "name": "coordinate_with_user", "description": "Start a coordination thread to interact with someone in the channel. Posts a message in the channel and starts a threaded conversation with the mentioned person. Use when Adam asks you to coordinate, schedule, or communicate with someone.", "inputSchema": { "json": { "type": "object", "properties": { "flat_message": { "type": "string", "description": "Brief message to post flat in the channel (e.g., 'Coordinating a meeting with Jane').", }, "thread_message": { "type": "string", "description": "The message to post in the thread, @mentioning the person (use Slack format <@USER_ID>).", }, }, "required": ["flat_message", "thread_message"], } }, } }, ] def _get_lambda_client(): global _lambda_client if _lambda_client is None: _lambda_client = boto3.client("lambda") return _lambda_client def _today(): return datetime.now(timezone.utc).strftime("%Y-%m-%d") def _format_email(msg): return { "message_id": msg.get("message_id") or msg.get("pk", "").replace("MSG#", ""), "thread_id": msg.get("thread_id", ""), "from": msg.get("from_name") or msg.get("from_email", ""), "from_email": msg.get("from_email", ""), "subject": msg.get("subject", ""), "snippet": msg.get("snippet", ""), "classification": msg.get("classification", ""), "classification_reason": msg.get("classification_reason", ""), "internal_date": int(msg.get("internal_date", 0)), } def dispatch_tool(tool_name, tool_input): handlers = { "get_inbox_summary": _handle_inbox_summary, "get_high_priority_emails": _handle_high_priority, "get_unanswered_threads": _handle_unanswered, "get_bypassed_work_orders": _handle_bypassed, "search_emails_by_sender": _handle_search_sender, "get_email_thread_detail": _handle_thread_detail, "trigger_daily_digest": _handle_trigger_digest, "search_inbox": _handle_search_inbox, "coordinate_with_user": _handle_coordinate, } handler = handlers.get(tool_name) if not handler: return {"error": f"Unknown tool: {tool_name}"} try: return handler(tool_input) except Exception: logger.exception("Tool %s failed", tool_name) return {"error": f"Tool {tool_name} failed"} def _handle_inbox_summary(inputs): date = inputs.get("date") or _today() messages = get_todays_messages(date) by_class = {} for msg in messages: c = msg.get("classification", "UNKNOWN") by_class[c] = by_class.get(c, 0) + 1 high = [_format_email(m) for m in messages if m.get("classification") == "HIGH"] return { "date": date, "total": len(messages), "by_classification": by_class, "high_priority_emails": high[:10], } def _handle_high_priority(inputs): date = inputs.get("date") or _today() messages = get_todays_messages(date) high = [_format_email(m) for m in messages if m.get("classification") == "HIGH"] return {"date": date, "count": len(high), "emails": high} def _handle_unanswered(inputs): config = get_config() threads = get_unanswered_threads(config["unanswered_threshold_hours"]) return { "count": len(threads), "threads": [ { "thread_id": t.get("pk", "").replace("THD#", ""), "subject": t.get("subject", ""), "last_message_from": t.get("last_message_from", ""), "unanswered_since": t.get("unanswered_since", ""), } for t in threads ], } def _handle_bypassed(inputs): date = inputs.get("date") or _today() messages = get_todays_messages(date) bypassed = [_format_email(m) for m in messages if m.get("bypassed_wo")] return {"date": date, "count": len(bypassed), "emails": bypassed} def _handle_search_sender(inputs): from shared.dynamo import _get_table sender = inputs["sender"].lower() limit = inputs.get("limit", 10) results = [] scan_kwargs = { "FilterExpression": "sk = :sk AND (contains(from_email, :sender) OR contains(from_name, :sender))", "ExpressionAttributeValues": {":sk": "MSG", ":sender": sender}, "Limit": 500, } resp = _get_table().scan(**scan_kwargs) results.extend(resp.get("Items", [])) while "LastEvaluatedKey" in resp and len(results) < limit: scan_kwargs["ExclusiveStartKey"] = resp["LastEvaluatedKey"] resp = _get_table().scan(**scan_kwargs) results.extend(resp.get("Items", [])) results.sort(key=lambda m: int(m.get("internal_date", 0)), reverse=True) return { "sender_query": inputs["sender"], "count": len(results[:limit]), "emails": [_format_email(m) for m in results[:limit]], } def _handle_thread_detail(inputs): service = get_authenticated_service() messages = fetch_thread(service, inputs["thread_id"]) return { "thread_id": inputs["thread_id"], "message_count": len(messages), "messages": [ { "from": m.get("from_name") or m.get("from_email", ""), "from_email": m.get("from_email", ""), "subject": m.get("subject", ""), "snippet": m.get("snippet", ""), "internal_date": m.get("internal_date", 0), } for m in messages ], } def _handle_trigger_digest(inputs): _get_lambda_client().invoke( FunctionName="exec-aide-daily-digest", InvocationType="Event", ) return {"status": "triggered", "message": "Daily digest has been triggered and will arrive shortly."} def _handle_search_inbox(inputs): service = get_authenticated_service() query = inputs["query"] max_results = inputs.get("max_results", 10) resp = service.users().messages().list( userId="me", q=query, maxResults=max_results ).execute() messages = [] for msg_meta in resp.get("messages", []): msg = service.users().messages().get( userId="me", id=msg_meta["id"], format="metadata", metadataHeaders=["From", "To", "Subject", "Date"], ).execute() headers = {h["name"]: h["value"] for h in msg["payload"]["headers"]} messages.append({ "id": msg["id"], "thread_id": msg["threadId"], "from": headers.get("From", ""), "to": headers.get("To", ""), "subject": headers.get("Subject", ""), "date": headers.get("Date", ""), "snippet": msg.get("snippet", ""), }) return {"query": query, "count": len(messages), "messages": messages} def _handle_coordinate(inputs): ctx = _tool_context if not ctx.get("channel"): return {"error": "Coordination is only available in channels, not DMs."} token = get_slack_token() channel = ctx["channel"] coord_ts = post_message(channel, inputs["flat_message"], token) post_message(channel, inputs["thread_message"], token, thread_ts=coord_ts) save_coordination( channel=channel, coord_thread_ts=coord_ts, conversation_thread_ts=ctx.get("thread_ts", ""), started_by="lauren", ) return { "status": "coordinating", "message": "Thread started. I'll follow up when they respond.", }