import { Construct } from 'constructs'; import * as cdk from 'aws-cdk-lib'; import * as dynamodb from 'aws-cdk-lib/aws-dynamodb'; import * as lambda from 'aws-cdk-lib/aws-lambda'; import * as events from 'aws-cdk-lib/aws-events'; import * as targets from 'aws-cdk-lib/aws-events-targets'; import * as scheduler from 'aws-cdk-lib/aws-scheduler'; import * as iam from 'aws-cdk-lib/aws-iam'; import * as logs from 'aws-cdk-lib/aws-logs'; import { PythonFunction } from '@aws-cdk/aws-lambda-python-alpha'; 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); const account = cdk.Stack.of(this).account; const region = cdk.Stack.of(this).region; // ── DynamoDB ────────────────────────────────────────────── this.table = new dynamodb.Table(this, 'Table', { tableName: 'exec-aide', billingMode: dynamodb.BillingMode.PAY_PER_REQUEST, partitionKey: { name: 'pk', type: dynamodb.AttributeType.STRING }, sortKey: { name: 'sk', type: dynamodb.AttributeType.STRING }, timeToLiveAttribute: 'ttl', removalPolicy: cdk.RemovalPolicy.RETAIN, }); this.table.addGlobalSecondaryIndex({ indexName: 'by-date', partitionKey: { name: 'classified_date', type: dynamodb.AttributeType.STRING }, sortKey: { name: 'sk', type: dynamodb.AttributeType.STRING }, projectionType: dynamodb.ProjectionType.ALL, }); // ── Shared environment + IAM ────────────────────────────── const lambdaEnv = { TABLE_NAME: this.table.tableName, SECRET_GMAIL: 'exec-aide/gmail-oauth', SECRET_SLACK: 'exec-aide/slack-credentials', SSM_PREFIX: '/exec-aide', }; const secretsReadPolicy = new iam.PolicyStatement({ actions: ['secretsmanager:GetSecretValue'], resources: [ `arn:aws:secretsmanager:${region}:${account}:secret:exec-aide/gmail-oauth-*`, `arn:aws:secretsmanager:${region}:${account}:secret:exec-aide/slack-credentials-*`, ], }); const secretsWritePolicy = new iam.PolicyStatement({ actions: ['secretsmanager:PutSecretValue'], resources: [ `arn:aws:secretsmanager:${region}:${account}:secret:exec-aide/gmail-oauth-*`, ], }); const ssmPolicy = new iam.PolicyStatement({ actions: ['ssm:GetParametersByPath', 'ssm:GetParameter'], resources: [ `arn:aws:ssm:${region}:${account}:parameter/exec-aide`, `arn:aws:ssm:${region}:${account}:parameter/exec-aide/*`, ], }); // ── Fetch & Classify Lambda ─────────────────────────────── const fetchClassify = new PythonFunction(this, 'FetchClassify', { functionName: 'exec-aide-fetch-classify', entry: path.join(__dirname, '../../src'), index: 'fetch_classify/app.py', handler: 'lambda_handler', runtime: lambda.Runtime.PYTHON_3_12, architecture: lambda.Architecture.ARM_64, memorySize: 256, timeout: cdk.Duration.seconds(120), environment: lambdaEnv, logRetention: logs.RetentionDays.TWO_MONTHS, }); this.table.grantReadWriteData(fetchClassify); fetchClassify.addToRolePolicy(secretsReadPolicy); fetchClassify.addToRolePolicy(secretsWritePolicy); fetchClassify.addToRolePolicy(ssmPolicy); fetchClassify.addToRolePolicy(new iam.PolicyStatement({ actions: ['bedrock:InvokeModel'], resources: [ 'arn:aws:bedrock:*::foundation-model/anthropic.*', `arn:aws:bedrock:${region}:${account}:inference-profile/us.anthropic.*`, ], })); new events.Rule(this, 'PollSchedule', { ruleName: 'exec-aide-fetch-classify-poll', description: 'Poll Gmail for new messages', schedule: events.Schedule.rate(cdk.Duration.minutes(15)), targets: [new targets.LambdaFunction(fetchClassify)], }); // ── Daily Digest Lambda ─────────────────────────────────── const dailyDigest = new PythonFunction(this, 'DailyDigest', { functionName: 'exec-aide-daily-digest', entry: path.join(__dirname, '../../src'), index: 'daily_digest/app.py', handler: 'lambda_handler', runtime: lambda.Runtime.PYTHON_3_12, architecture: lambda.Architecture.ARM_64, memorySize: 256, timeout: cdk.Duration.seconds(120), environment: lambdaEnv, logRetention: logs.RetentionDays.TWO_MONTHS, }); this.table.grantReadWriteData(dailyDigest); dailyDigest.addToRolePolicy(secretsReadPolicy); dailyDigest.addToRolePolicy(secretsWritePolicy); dailyDigest.addToRolePolicy(ssmPolicy); // ── EventBridge Scheduler (DST-aware 5 PM ET) ───────────── const schedulerRole = new iam.Role(this, 'DigestSchedulerRole', { roleName: 'exec-aide-digest-scheduler', assumedBy: new iam.ServicePrincipal('scheduler.amazonaws.com'), }); dailyDigest.grantInvoke(schedulerRole); const schedule = new scheduler.CfnSchedule(this, 'DigestSchedule', { 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: dailyDigest.functionArn, roleArn: schedulerRole.roleArn, }, }); dailyDigest.addPermission('SchedulerInvoke', { 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; } }