import { Construct } from 'constructs'; import * as cdk from 'aws-cdk-lib'; import * as dynamodb from 'aws-cdk-lib/aws-dynamodb'; import * as kms from 'aws-cdk-lib/aws-kms'; import * as ssm from 'aws-cdk-lib/aws-ssm'; 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 * as cloudwatch from 'aws-cdk-lib/aws-cloudwatch'; import * as cloudwatchActions from 'aws-cdk-lib/aws-cloudwatch-actions'; import * as sns from 'aws-cdk-lib/aws-sns'; import { PythonFunction } from '@aws-cdk/aws-lambda-python-alpha'; import * as path from 'path'; export interface EmailPipelineProps { /** * Shared site-alerts SNS topic for CloudWatch ALARM actions. Imported once at * the stack level (sns.Topic.fromTopicArn) and injected here, mirroring how * the table is wired into the socket-mode construct. */ alarmTopic: sns.ITopic; } export class EmailPipelineConstruct extends Construct { public readonly table: dynamodb.Table; public readonly conversationFn: lambda.IFunction; constructor(scope: Construct, id: string, props: EmailPipelineProps) { super(scope, id); const account = cdk.Stack.of(this).account; const region = cdk.Stack.of(this).region; // ── DynamoDB ────────────────────────────────────────────── // Shared customer-managed CMK for sensitive DynamoDB SSE (INFRA-95 / M-3). // The key is owned by the seahaven-dynamodb-cmk stack (account-baseline // repo); its ARN is published to SSM and imported here. The exec-aide table // holds inbox PII (email senders/subjects/bodies, conversation state). // CDK's grantReadWriteData/grantReadData auto-add the matching KMS actions // to every consumer role (the 4 pipeline Lambdas + the socket-mode Fargate // task role) when the table carries an encryptionKey, so no manual KMS grant // is needed. SSE change is an in-place UpdateTable (no downtime). const dynamodbCmk = kms.Key.fromKeyArn( this, 'DynamoDbCmk', ssm.StringParameter.valueForStringParameter(this, '/seahaven/dynamodb/cmk-arn'), ); 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', // INFRA-95 / M-3: customer-managed CMK SSE (was AWS-owned key). encryption: dynamodb.TableEncryption.CUSTOMER_MANAGED, encryptionKey: dynamodbCmk, 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; // ── CloudWatch Alarms ───────────────────────────────────── // ALARM-only (no OK / InsufficientData action); treatMissingData // NOT_BREACHING. Single shared site-alerts action, mirroring the // proposal-system alarm construct. Names are repo-namespaced kebab-case: // exec-aide--. const alarmAction = new cloudwatchActions.SnsAction(props.alarmTopic); // Per-function timeout (ms). Duration alarms fire at ~80% of timeout so a // function trending toward a hard timeout pages before requests start // failing outright. These thresholds need Adam sign-off. const lambdaAlarmSpecs: { readonly fn: lambda.Function; readonly name: string; readonly durationThresholdMs: number; }[] = [ { fn: fetchClassify, name: 'fetch-classify', durationThresholdMs: 96000 }, // 80% of 120s { fn: dailyDigest, name: 'daily-digest', durationThresholdMs: 96000 }, // 80% of 120s { fn: conversation, name: 'conversation', durationThresholdMs: 144000 }, // 80% of 180s ]; for (const { fn, name, durationThresholdMs } of lambdaAlarmSpecs) { // Errors: any invocation error pages. new cloudwatch.Alarm(this, `Errors-${name}`, { alarmName: `exec-aide-${name}-errors`, alarmDescription: `exec-aide-${name} Lambda is erroring`, metric: fn.metricErrors({ period: cdk.Duration.minutes(5), statistic: 'Sum', }), threshold: 1, comparisonOperator: cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD, evaluationPeriods: 1, treatMissingData: cloudwatch.TreatMissingData.NOT_BREACHING, }).addAlarmAction(alarmAction); // Throttles: concurrency exhaustion / account limits. new cloudwatch.Alarm(this, `Throttles-${name}`, { alarmName: `exec-aide-${name}-throttles`, alarmDescription: `exec-aide-${name} Lambda is being throttled`, metric: fn.metricThrottles({ period: cdk.Duration.minutes(5), statistic: 'Sum', }), threshold: 1, comparisonOperator: cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD, evaluationPeriods: 1, treatMissingData: cloudwatch.TreatMissingData.NOT_BREACHING, }).addAlarmAction(alarmAction); // Duration: p99 trending toward the timeout. eval3 / datapoints2 to ride // out a single slow Bedrock/Gmail call without paging. new cloudwatch.Alarm(this, `Duration-${name}`, { alarmName: `exec-aide-${name}-duration`, alarmDescription: `exec-aide-${name} Lambda p99 duration approaching its timeout (~80%)`, metric: fn.metricDuration({ period: cdk.Duration.minutes(5), statistic: 'p99', }), threshold: durationThresholdMs, comparisonOperator: cloudwatch.ComparisonOperator.GREATER_THAN_THRESHOLD, evaluationPeriods: 3, datapointsToAlarm: 2, treatMissingData: cloudwatch.TreatMissingData.NOT_BREACHING, }).addAlarmAction(alarmAction); } // DynamoDB alarms. NOTE on dimensions: the raw `ThrottledRequests` and // `SystemErrors` metrics are NOT published at the bare TableName dimension — // AWS emits them keyed by the `Operation` dimension. CDK's // metricThrottledRequests / metricSystemErrors helpers are deprecated for // exactly this reason ("returns an invalid metric"). We therefore use the // *ForOperations math helpers, which build a per-operation math expression // scoped to TableName=exec-aide. An alarm math expression caps at 10 // metrics, and the default "all operations" set exceeds that, so we scope to // the operations this single-table app actually issues (GetItem, PutItem, // Query, Scan, UpdateItem, DeleteItem — confirmed from src/). const tableOperations = [ dynamodb.Operation.GET_ITEM, dynamodb.Operation.PUT_ITEM, dynamodb.Operation.QUERY, dynamodb.Operation.SCAN, dynamodb.Operation.UPDATE_ITEM, dynamodb.Operation.DELETE_ITEM, ]; new cloudwatch.Alarm(this, 'DdbThrottles', { alarmName: 'exec-aide-table-throttles', alarmDescription: 'exec-aide DynamoDB table is throttling requests (summed across the operations the app uses)', metric: this.table.metricThrottledRequestsForOperations({ operations: tableOperations, period: cdk.Duration.minutes(5), statistic: 'Sum', }), threshold: 1, comparisonOperator: cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD, evaluationPeriods: 1, treatMissingData: cloudwatch.TreatMissingData.NOT_BREACHING, }).addAlarmAction(alarmAction); new cloudwatch.Alarm(this, 'DdbSystemErrors', { alarmName: 'exec-aide-table-system-errors', alarmDescription: 'exec-aide DynamoDB table is returning SystemErrors (summed across the operations the app uses)', metric: this.table.metricSystemErrorsForOperations({ operations: tableOperations, period: cdk.Duration.minutes(5), statistic: 'Sum', }), threshold: 1, comparisonOperator: cloudwatch.ComparisonOperator.GREATER_THAN_OR_EQUAL_TO_THRESHOLD, evaluationPeriods: 1, treatMissingData: cloudwatch.TreatMissingData.NOT_BREACHING, }).addAlarmAction(alarmAction); } }