Add conversational assistant via Bedrock Sonnet tool-use

Listener posts "Thinking..." placeholder and async-invokes a new
exec-aide-conversation Lambda that runs a Bedrock Sonnet tool-use loop
over 7 email tools, then updates the Slack message with the response.
This commit is contained in:
Adam Moussa 2026-05-01 11:31:29 -04:00
parent df56bd6e9a
commit aebf6c43c6
10 changed files with 491 additions and 12 deletions

View file

@ -1,13 +1,13 @@
# 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 via a dedicated Slack app.
Personal AI executive assistant for Adam Moussa. Monitors Gmail inbox, classifies emails by urgency using Claude Haiku (Bedrock), delivers Slack alerts and daily digests, and supports two-way conversation via Slack DM.
## 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
Fargate (Socket Mode) → Slack listener → two-way conversation with Adam
Fargate (Socket Mode) → Slack listener → conversation Lambda → Bedrock Sonnet + tools → Slack reply
```
**Stack:** CDK (TypeScript), stack name `exec-aide`, region us-east-1
@ -16,12 +16,27 @@ Fargate (Socket Mode) → Slack listener → two-way conversation with Adam
- **exec-aide-fetch-classify** (Lambda): Polls Gmail via History API every 15 minutes, classifies each new message via Bedrock Haiku, filters out marketing, sends immediate Slack DM for HIGH priority, tracks thread state for unanswered detection.
- **exec-aide-daily-digest** (Lambda): Runs weekdays at 5 PM ET (DST-aware). Summarizes high-priority items, bypassed work orders, and unanswered threads older than 24h.
- **exec-aide-listener** (Fargate): Socket Mode Slack bot for two-way conversation. Runs in a dedicated VPC (10.30.0.0/16) on ARM64.
- **DynamoDB `exec-aide`**: Single-table design (MSG#, THD#, META# prefixes). GSI `by-date` for daily digest queries. 90-day TTL.
- **exec-aide-conversation** (Lambda): Bedrock Sonnet-powered conversational assistant. Handles multi-turn conversations with tool use for inbox queries, email search, thread details, and on-demand digest triggers. Invoked asynchronously by the listener.
- **exec-aide-listener** (Fargate): Socket Mode Slack bot for two-way conversation. Receives DMs, posts "Thinking..." placeholder, invokes conversation Lambda async, which updates the placeholder with the response. Runs in a dedicated VPC (10.30.0.0/16) on ARM64.
- **DynamoDB `exec-aide`**: Single-table design (MSG#, THD#, META#, CONV# prefixes). GSI `by-date` for daily digest queries. 90-day TTL.
### Conversation Tools
The conversation Lambda has access to these tools:
| Tool | Description |
|---|---|
| `get_inbox_summary` | Today's email stats + recent high-priority items |
| `get_high_priority_emails` | HIGH-classified emails for a given date |
| `get_unanswered_threads` | Threads awaiting reply for 24h+ |
| `get_bypassed_work_orders` | Emails sent without CC'ing work-orders@ |
| `search_emails_by_sender` | Find emails from a specific sender |
| `get_email_thread_detail` | Full thread context from Gmail |
| `trigger_daily_digest` | Send a digest on demand |
### CDK Constructs
- `lib/constructs/email-pipeline.ts` — DynamoDB table, both Lambda functions, EventBridge schedules
- `lib/constructs/email-pipeline.ts` — DynamoDB table, all three Lambda functions, EventBridge schedules
- `lib/constructs/socket-mode.ts` — VPC, ECS cluster, Fargate service, ECR repo (image built automatically via `ContainerImage.fromAsset()`)
## Classification Rules
@ -85,6 +100,8 @@ aws lambda invoke --function-name exec-aide-fetch-classify /dev/stdout
aws lambda invoke --function-name exec-aide-daily-digest /dev/stdout
```
Or DM the Exec Aide bot in Slack to test the conversational assistant.
## Updating VIP List
Edit the SSM parameters — changes take effect on the next Lambda cold start:

View file

@ -12,6 +12,7 @@ import * as path from 'path';
export class EmailPipelineConstruct extends Construct {
public readonly table: dynamodb.Table;
public readonly conversationFn: lambda.IFunction;
constructor(scope: Construct, id: string) {
super(scope, id);
@ -148,5 +149,35 @@ export class EmailPipelineConstruct extends Construct {
principal: new iam.ServicePrincipal('scheduler.amazonaws.com'),
sourceArn: `arn:aws:scheduler:${region}:${account}:schedule/default/${schedule.name}`,
});
// ── Conversation Lambda ───────────────────────────────────
const conversation = new PythonFunction(this, 'Conversation', {
functionName: 'exec-aide-conversation',
entry: path.join(__dirname, '../../src'),
index: 'conversation/app.py',
handler: 'lambda_handler',
runtime: lambda.Runtime.PYTHON_3_12,
architecture: lambda.Architecture.ARM_64,
memorySize: 512,
timeout: cdk.Duration.seconds(180),
environment: lambdaEnv,
logRetention: logs.RetentionDays.TWO_MONTHS,
});
this.table.grantReadWriteData(conversation);
conversation.addToRolePolicy(secretsReadPolicy);
conversation.addToRolePolicy(secretsWritePolicy);
conversation.addToRolePolicy(ssmPolicy);
conversation.addToRolePolicy(new iam.PolicyStatement({
actions: ['bedrock:InvokeModel'],
resources: [
'arn:aws:bedrock:*::foundation-model/anthropic.*',
`arn:aws:bedrock:${region}:${account}:inference-profile/us.anthropic.*`,
],
}));
dailyDigest.grantInvoke(conversation);
this.conversationFn = conversation;
}
}

View file

@ -10,6 +10,7 @@ import * as path from 'path';
export interface SocketModeProps {
table: dynamodb.ITable;
conversationFnArn: string;
}
export class SocketModeConstruct extends Construct {
@ -81,6 +82,7 @@ export class SocketModeConstruct extends Construct {
TABLE_NAME: props.table.tableName,
SECRET_SLACK: 'exec-aide/slack-credentials',
SSM_PREFIX: '/exec-aide',
CONVERSATION_FN_ARN: props.conversationFnArn,
},
logging: ecs.LogDrivers.awsLogs({
streamPrefix: 'listener',
@ -110,6 +112,11 @@ export class SocketModeConstruct extends Construct {
],
}));
taskDef.taskRole.addToPrincipalPolicy(new iam.PolicyStatement({
actions: ['lambda:InvokeFunction'],
resources: [props.conversationFnArn],
}));
// ── Fargate Service ────────────────────────────────────────
new ecs.FargateService(this, 'Service', {

View file

@ -11,6 +11,7 @@ export class ExecAideStack extends cdk.Stack {
new SocketModeConstruct(this, 'SocketMode', {
table: emailPipeline.table,
conversationFnArn: emailPipeline.conversationFn.functionArn,
});
new cdk.CfnOutput(this, 'TableName', {

View file

@ -1,7 +1,6 @@
import json
import logging
import os
import time
import boto3
from slack_bolt import App
@ -12,6 +11,7 @@ logger = logging.getLogger("exec-aide")
secrets_client = boto3.client("secretsmanager")
ssm_client = boto3.client("ssm")
lambda_client = boto3.client("lambda")
_config = None
@ -55,16 +55,26 @@ def handle_message(event, say):
return
text = event.get("text", "").strip()
if not text:
return
thread_ts = event.get("thread_ts") or event["ts"]
channel = event["channel"]
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."
result = say(text="Thinking...", thread_ts=thread_ts)
placeholder_ts = result["ts"]
lambda_client.invoke(
FunctionName=os.environ["CONVERSATION_FN_ARN"],
InvocationType="Event",
Payload=json.dumps({
"user_message": text,
"thread_ts": thread_ts,
"channel": channel,
"placeholder_ts": placeholder_ts,
}),
)
say(text=response, thread_ts=thread_ts)
def main():
handler = SocketModeHandler(app, creds["appToken"])

View file

104
src/conversation/app.py Normal file
View file

@ -0,0 +1,104 @@
import json
import logging
import boto3
from shared.dynamo import get_conversation, save_conversation
from shared.secrets import get_slack_token
from shared.slack import update_message
from conversation.tools import TOOLS, dispatch_tool
logger = logging.getLogger()
logger.setLevel(logging.INFO)
_bedrock = boto3.client("bedrock-runtime", region_name="us-east-1")
MODEL_ID = "us.anthropic.claude-sonnet-4-6"
MAX_TOOL_ITERATIONS = 10
SYSTEM_PROMPT = """\
You are Exec Aide, a personal executive assistant for Adam Moussa, President \
of Sea Haven Industries (a facilities maintenance company).
You have access to Adam's email inbox. Use the available tools to answer \
questions about emails, find information, and provide inbox summaries.
Be concise and direct. Adam is busy - give him the key information without \
unnecessary preamble. Use bullet points for lists of emails.
When referencing emails, always include the sender name and subject line. \
For time references, use relative terms (e.g., "2 hours ago", "yesterday").
If Adam asks about something outside your current capabilities (calendar, \
tasks, etc.), acknowledge it and let him know it's not available yet.\
"""
def lambda_handler(event, context):
user_message = event["user_message"]
thread_ts = event["thread_ts"]
channel = event["channel"]
placeholder_ts = event["placeholder_ts"]
try:
messages = get_conversation(thread_ts)
messages.append({"role": "user", "content": [{"text": user_message}]})
assistant_message = _run_conversation(messages)
final_text = _extract_text(assistant_message)
save_conversation(thread_ts, messages)
update_message(channel, placeholder_ts, final_text, get_slack_token())
logger.info("Conversation complete: thread=%s, turns=%d", thread_ts, len(messages))
except Exception:
logger.exception("Conversation failed for thread %s", thread_ts)
try:
update_message(
channel, placeholder_ts,
"Sorry, I hit an error processing that. Try again in a moment.",
get_slack_token(),
)
except Exception:
logger.exception("Failed to update error placeholder")
def _run_conversation(messages):
assistant_message = None
for _ in range(MAX_TOOL_ITERATIONS):
response = _bedrock.converse(
modelId=MODEL_ID,
system=[{"text": SYSTEM_PROMPT}],
messages=messages,
toolConfig={"tools": TOOLS},
inferenceConfig={"maxTokens": 4096},
)
assistant_message = response["output"]["message"]
messages.append(assistant_message)
if response["stopReason"] != "tool_use":
break
tool_results = []
for block in assistant_message["content"]:
tool_use = block.get("toolUse")
if tool_use:
logger.info("Tool call: %s(%s)", tool_use["name"], json.dumps(tool_use["input"]))
result = dispatch_tool(tool_use["name"], tool_use["input"])
tool_results.append({
"toolResult": {
"toolUseId": tool_use["toolUseId"],
"content": [{"text": json.dumps(result)}],
}
})
messages.append({"role": "user", "content": tool_results})
return assistant_message
def _extract_text(message):
parts = []
for block in message.get("content", []):
if "text" in block:
parts.append(block["text"])
return "\n".join(parts) if parts else "I processed your request but have nothing to report."

279
src/conversation/tools.py Normal file
View file

@ -0,0 +1,279 @@
import json
import logging
import os
from datetime import datetime, timezone
import boto3
from shared.dynamo import get_todays_messages, get_unanswered_threads
from shared.gmail import get_authenticated_service, fetch_thread
from shared.secrets import get_config
logger = logging.getLogger(__name__)
_lambda_client = None
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": {},
}
},
}
},
]
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": 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,
}
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."}

View file

@ -137,6 +137,24 @@ def get_todays_messages(date_str):
return items
def get_conversation(thread_ts):
resp = _get_table().get_item(
Key={"pk": f"CONV#{thread_ts}", "sk": "CONV"}
)
item = resp.get("Item")
return item["messages"] if item else []
def save_conversation(thread_ts, messages):
_get_table().put_item(Item={
"pk": f"CONV#{thread_ts}",
"sk": "CONV",
"messages": messages,
"updated_at": datetime.now(timezone.utc).isoformat(),
"ttl": _ttl_90_days(),
})
def get_unanswered_threads(threshold_hours):
cutoff = datetime.now(timezone.utc).timestamp() - (threshold_hours * 3600)
items = []

View file

@ -169,6 +169,18 @@ def send_message_to_channel(channel, text, token):
logger.error("chat.postMessage failed: %s", data.get("error"))
def update_message(channel, ts, text, token):
resp = requests.post(
"https://slack.com/api/chat.update",
headers={"Authorization": f"Bearer {token}"},
json={"channel": channel, "ts": ts, "text": text},
)
resp.raise_for_status()
data = resp.json()
if not data.get("ok"):
raise RuntimeError(f"chat.update failed: {data.get('error')}")
def _format_time_ago(internal_date):
if not internal_date:
return "unknown"