Tag Archives: Technical How-to

Enable cross-cloud analytics with Amazon S3 Tables and Google BigQuery, Part 2: access control with Lake Formation

Post Syndicated from Lakshmi Nair original https://aws.amazon.com/blogs/big-data/enable-cross-cloud-analytics-with-amazon-s3-tables-and-google-bigquery-part-2-access-control-with-lake-formation/

In Part 1, we showed how to connect Google BigQuery to Amazon Simple Storage Service (Amazon S3) Tables, a capability of Amazon S3, using access control based on AWS Identity and Access Management (IAM). A single IAM policy governs both table metadata and data access. We also walked through common cross-cloud analytics scenarios where this pattern adds value. This post covers the approach using AWS Lake Formation. Instead of relying solely on IAM policies for data access, Lake Formation manages fine-grained permissions and vends temporary, scoped credentials to the requesting engine. This is a better fit when multiple engines need different levels of access to the same tables, or when you want to manage grants centrally without touching IAM policies every time a new consumer comes along.

Solution overview

You use the AWS Glue Iceberg REST Catalog (IRC) as the bridge between BigQuery and S3 Tables. BigQuery’s cross-cloud Lakehouse creates a federated catalog that syncs metadata from the Glue IRC, then uses the synced metadata to read Iceberg data files directly.

Architecture diagram showing BigQuery connecting to Amazon S3 Tables through the AWS Glue Iceberg REST Catalog

Figure 1: Architecture diagram showing BigQuery connecting to Amazon S3 Tables through the AWS Glue Iceberg REST Catalog

The key components in this architecture:

  1. Amazon S3 Tables: With Amazon S3 Tables, data is stored in table buckets, specifically designed for storing tables in the Apache Iceberg format. Table metadata is registered on AWS Glue Data Catalog for discovery and governance.
  2. AWS Glue Data Catalog: With AWS Glue Data Catalog, you can access the federated s3tablescatalog catalog that maps S3 Tables resources (table buckets, namespaces, tables) into a catalog hierarchy from supported analytics engines. The standard Iceberg REST endpoint of Glue Data Catalog serves table metadata to external engines. BigQuery connects through this endpoint.
  3. AWS Lake Formation: With AWS Lake Formation, you define access permissions at the catalog, database, and table level. Instead of granting broad IAM permissions for data access, Lake Formation evaluates permissions at query time and issues short-lived credentials limited to the resources the caller is authorized to read.
  4. Google Cross-Cloud Lakehouse: With Google Cross-Cloud Lakehouse, you can connect BigQuery to external Iceberg catalogs. It assumes an IAM role using OpenID Connect (OIDC), calls the AWS Glue Iceberg REST endpoint, and syncs metadata on a configurable refresh interval.

Prerequisites

Before you begin, you need:

  • An AWS account with Amazon S3 Tables available in your AWS Region.
  • A Google Cloud project with billing enabled and the BigLake API activated.
  • AWS Command Line Interface (AWS CLI) and gcloud CLI installed and configured.
  • An S3 table bucket with at least one namespace and table containing data.

Setting up Amazon S3 Tables

If you already have S3 Tables with data, skip to the next section. Otherwise, create a table bucket, namespace, and populate a table.

Create a table bucket and namespace

Use AWS CLI to create resources as follows:

# Create a Table bucket
aws s3tables create-table-bucket \
    --name <TABLE_BUCKET_NAME> \
    --region <REGION>

# Create a Namespace (Database)
aws s3tables create-namespace \
    --table-bucket-arn "arn:aws:s3tables:<REGION>:<AWS_ACCOUNT_ID>:bucket/<TABLE_BUCKET_NAME>" \
    --namespace <NAMESPACE> \
    --region <REGION>

Set up S3 Tables integration with the Glue Data Catalog using Lake Formation mode

Lake Formation needs its own service role to interact with S3 Tables on your behalf. This is the role Lake Formation assumes internally when it reads or writes data on behalf of authorized callers.

Create a Lake Formation service IAM role named LakeFormationS3TablesServiceRole with the following policy:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "LakeFormationPermissionsForS3ListTableBucket",
      "Effect": "Allow",
      "Action": ["s3tables:ListTableBuckets"],
      "Resource": ["*"]
    },
    {
      "Sid": "LakeFormationDataAccessPermissionsForS3TableBucket",
      "Effect": "Allow",
      "Action": [
        "s3tables:CreateTableBucket", "s3tables:GetTableBucket",
        "s3tables:CreateNamespace", "s3tables:GetNamespace",
        "s3tables:ListNamespaces", "s3tables:DeleteNamespace",
        "s3tables:DeleteTableBucket", "s3tables:CreateTable",
        "s3tables:DeleteTable", "s3tables:GetTable",
        "s3tables:ListTables", "s3tables:RenameTable",
        "s3tables:UpdateTableMetadataLocation", "s3tables:GetTableMetadataLocation",
        "s3tables:GetTableData", "s3tables:PutTableData"
      ],
      "Resource": ["arn:aws:s3tables:<AWS_REGION>:<AWS_ACCOUNT_ID>:bucket/*"]
    }
  ]
}

Attach the following trust relationship:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "LakeFormationDataAccessPolicy",
      "Effect": "Allow",
      "Principal": { "Service": "lakeformation.amazonaws.com" },
      "Action": ["sts:AssumeRole", "sts:SetContext", "sts:SetSourceIdentity"],
      "Condition": { "StringEquals": { "aws:SourceAccount": "<AWS_ACCOUNT_ID>" } }
    }
  ]
}

In the Lake Formation console, in the navigation pane, choose Catalogs, and then choose Enable S3 Table Integration.

The Enable S3 Table Integration option on the Catalogs page of the Lake Formation console

Figure 2: Enabling the S3 Tables integration in the Lake Formation console

Choose the role you created earlier when prompted for an IAM role, and select Allow external engines to access data in Amazon S3 locations with full table access.

S3 Tables integration performs the following:

  1. Registers the S3 Tables data location with Lake Formation.
  2. Creates the s3tablescatalog federated catalog in Glue.

Important: Before enabling the integration, verify your Lake Formation data lake settings have empty default permissions to prevent IAMAllowedPrincipals from being auto-granted on the catalog:

aws lakeformation put-data-lake-settings \
    --data-lake-settings '{"DataLakeAdmins":[{"DataLakePrincipalIdentifier":"arn:aws:iam::<AWS_ACCOUNT_ID>:role/<ADMIN_ROLE>"}],"CreateDatabaseDefaultPermissions":[],"CreateTableDefaultPermissions":[]}' \
    --region <AWS_REGION>
The S3 Tables integration dialog in Lake Formation with full table access selected

Figure 3: Selecting full table access for external engines during S3 Tables integration

When you select this option, you  allow external engines to access data in Amazon S3 locations with full table access, and Lake Formation grants full table-level access to external engines. Column-level and row-level filtering are not enforced for external engine connections. Access is granted at the whole-table level.

Verify the integration by confirming the catalog in Lake Formation console.

Create a table and insert data

Now, to create the table and insert data, open the Amazon Athena console. In the query editor, select s3tablescatalog/<TABLE_BUCKET_NAME> as your data source and <NAMESPACE> as the database. Then run the following SQL statements one by one:

CREATE TABLE `<NAMESPACE>`.orders (
    order_id STRING,
    customer_id STRING,
    amount BIGINT,
    order_date DATE,
    region STRING
)
TBLPROPERTIES ('table_type' = 'iceberg');

INSERT INTO orders
VALUES
    ('ORD-001', 'C100', 4500, DATE '2024-06-01', 'EMEA'),
    ('ORD-002', 'C200', 8900, DATE '2024-06-01', 'EMEA'),
    ('ORD-003', 'C100', 3200, DATE '2024-06-02', 'NAMER'),
    ('ORD-004', 'C300', 12000, DATE '2024-06-02', 'NAMER'),
    ('ORD-005', 'C400', 6700, DATE '2024-06-03', 'APJ'),
    ('ORD-006', 'C200', 4100, DATE '2024-06-03', 'APJ'),
    ('ORD-007', 'C500', 9500, DATE '2024-06-04', 'EMEA'),
    ('ORD-008', 'C100', 2800, DATE '2024-06-04', 'LATAM'),
    ('ORD-009', 'C600', 15000, DATE '2024-06-05', 'NAMER'),
    ('ORD-010', 'C300', 7200, DATE '2024-06-05', 'LATAM');

Configuring cross-cloud access

BigQuery assumes an AWS IAM role using OIDC federation to access the AWS Glue IRC. This section walks through creating the role, OIDC provider, and permissions.

Create the OIDC identity provider

Register Google as an OIDC identity provider in your AWS account. This allows AWS to validate tokens issued by Google’s identity service:

aws iam create-open-id-connect-provider \
    --url https://accounts.google.com \
    --client-id-list accounts.google.com \
    --thumbprint-list 08745487e891c19e3078c1f2a07e452950ef36f6

The –thumbprint-list parameter is optional. When omitted, IAM automatically retrieves the thumbprint from the OIDC provider’s certificate. See AWS documentation for details.

Create the cross-cloud IAM role on AWS

Sign in to the AWS Management Console. Create the role with a placeholder trust policy. You will update it with the actual BigLake service account ID after you create the federated catalog in Google Cloud.

aws iam create-role \
    --role-name bigquery-cross-cloud-role \
    --max-session-duration 43200 \
    --assume-role-policy-document '{
      "Version": "2012-10-17",
      "Statement": [{
        "Effect": "Allow",
        "Principal": {
          "Federated": "arn:aws:iam::<AWS_ACCOUNT_ID>:oidc-provider/accounts.google.com"
        },
        "Action": "sts:AssumeRoleWithWebIdentity",
        "Condition": {
          "StringEquals": {
            "accounts.google.com:sub": ["PLACEHOLDER"],
            "accounts.google.com:aud": ["PLACEHOLDER"]
          }
        }
      }]
    }'

The --max-session-duration 43200 allows sessions up to 12 hours, which is needed for long-running BigQuery queries.

Attach permissions

The permissions policy differs based on your access control approach. For the Lake Formation approach, attach the following policy:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "GlueRead",
      "Effect": "Allow",
      "Action": [
        "glue:GetCatalog", "glue:GetDatabase", "glue:GetDatabases",
        "glue:GetTable", "glue:GetTables", "glue:GetPartition", "glue:GetPartitions"
      ],
      "Resource": [
        "arn:aws:glue:<AWS_REGION>:<AWS_ACCOUNT_ID>:catalog",
        "arn:aws:glue:<AWS_REGION>:<AWS_ACCOUNT_ID>:catalog/s3tablescatalog",
        "arn:aws:glue:<AWS_REGION>:<AWS_ACCOUNT_ID>:catalog/s3tablescatalog/<TABLE_BUCKET>",
        "arn:aws:glue:<AWS_REGION>:<AWS_ACCOUNT_ID>:database/s3tablescatalog/<TABLE_BUCKET>/<NAMESPACE>",
        "arn:aws:glue:<AWS_REGION>:<AWS_ACCOUNT_ID>:table/s3tablescatalog/<TABLE_BUCKET>/<NAMESPACE>/*"
      ]
    },
    {
      "Sid": "S3TablesRead",
      "Effect": "Allow",
      "Action": [
        "s3tables:GetTableBucket", "s3tables:ListTableBuckets",
        "s3tables:ListNamespaces", "s3tables:GetNamespace",
        "s3tables:ListTables", "s3tables:GetTable",
        "s3tables:GetTableMetadataLocation", "s3tables:GetTableData"
      ],
      "Resource": [
        "arn:aws:s3tables:<AWS_REGION>:<AWS_ACCOUNT_ID>:bucket/<TABLE_BUCKET>",
        "arn:aws:s3tables:<AWS_REGION>:<AWS_ACCOUNT_ID>:bucket/<TABLE_BUCKET>/*"
      ]
    },
    {
      "Sid": "LakeFormationCredentialVending",
      "Effect": "Allow",
      "Action": ["lakeformation:GetDataAccess"],
      "Resource": "*"
    }
  ]
}

Grant Lake Formation permissions

Lake Formation permissions work as a layered grant model: you grant access at each level of the catalog hierarchy, from catalog down to table. The cross-cloud role needs DESCRIBE on the catalog and database so it can discover what exists, and SELECT plus DESCRIBE on the table so it can read the actual data. Without grants at every level, Lake Formation denies access even if the IAM policy allows it.

If using Lake Formation, grant the bigquery-cross-cloud-role access to your tables:

  • Grant catalog permission: DESCRIBE.
  • Grant database permission: DESCRIBE.
  • Grant table permission: SELECT, DESCRIBE.

Grant Lake Formation permissions on the cross-cloud role (one-time).

aws lakeformation grant-permissions     --principal '{"DataLakePrincipalIdentifier":"arn:aws:iam::<AWS_ACCOUNT_ID>:role/bigquery-cross-cloud-role"}'     --resource '{"Catalog":{"Id":"<AWS_ACCOUNT_ID>:s3tablescatalog/<TABLE_BUCKET>"}}'     --permissions '["DESCRIBE"]'     --region <AWS_REGION>

aws lakeformation grant-permissions     --principal '{"DataLakePrincipalIdentifier":"arn:aws:iam::<AWS_ACCOUNT_ID>:role/bigquery-cross-cloud-role"}'     --resource '{"Database":{"CatalogId":"<AWS_ACCOUNT_ID>:s3tablescatalog/<TABLE_BUCKET>","Name":"<NAMESPACE>"}}'     --permissions '["DESCRIBE"]'     --region <AWS_REGION>

aws lakeformation grant-permissions     --principal '{"DataLakePrincipalIdentifier":"arn:aws:iam::<AWS_ACCOUNT_ID>:role/bigquery-cross-cloud-role"}'     --resource '{"Table":{"CatalogId":"<AWS_ACCOUNT_ID>:s3tablescatalog/<TABLE_BUCKET>","DatabaseName":"<NAMESPACE>","Name":"orders"}}'     --permissions '["SELECT","DESCRIBE"]'     --region <AWS_REGION>

Before granting Lake Formation permissions, revoke the default IAMAllowedPrincipals access. By default, Lake Formation grants IAMAllowedPrincipals full access to all databases and tables, so you first need to revoke this to enforce fine grain access. IAMAllowedPrincipals provides backward compatibility when you start using Lake Formation permissions to secure the Data Catalog resources that were earlier protected by IAM policies for AWS Glue.

Set up Lake Formation for external engines

For table metadata to sync from Glue to BigLake/BigQuery, the following Lake Formation settings are required. You might notice that a similar setting also appeared during the S3 Table integration setup. The first one registers the data location and enables external access at the catalog level, while this one enables the Lake Formation credential vending mechanism at the account level for all external engines. For a clean cross-cloud setup, we recommend that you enable both.

In the Lake Formation console, choose Administration, then Application integration settings, and then select Allow external engines to access data in Amazon S3 locations with full table access.

Application integration settings in the Lake Formation console with external-engine access enabled

Figure 4: Enabling external-engine access in Lake Formation application integration settings

Connecting BigQuery to S3 Tables

With the AWS side configured, create the federated catalog in Google Cloud that connects BigQuery to the AWS Glue IRC.

Create the federated catalog

Authenticate to Google Cloud using gcloud auth login, or use Cloud Shell, which is pre-authenticated. Verify the BigLake API is enabled:

gcloud services enable biglake.googleapis.com --project="<GCP_PROJECT_ID>"

For Lake Formation mode (with credential vending):

gcloud alpha biglake iceberg catalogs create <FEDERATED_CATALOG_NAME> \
    --project="<GCP_PROJECT_ID>" \
    --catalog-type=federated \
    --federated-catalog-type=glue \
    --glue-aws-region=<AWS_REGION> \
    --glue-aws-role-arn=arn:aws:iam::<ACCOUNT_ID>:role/bigquery-cross-cloud-role \
    --glue-warehouse=<ACCOUNT_ID>:s3tablescatalog/<TABLE_BUCKET> \
    --primary-location=<GCP_REGION> \
    --credential-mode=vended-credentials

The --glue-warehouse parameter uses the format <AWS_ACCOUNT_ID>:s3tablescatalog/<TABLE_BUCKET>. This tells the AWS Glue IRC to scope requests to your specific S3 Tables bucket within the federated catalog hierarchy.

The --credential-mode=vended-credentials flag (Lake Formation mode) instructs BigQuery Lakehouse to request scoped temporary credentials from Lake Formation rather than using the role’s IAM permissions directly for data access.

The --primary-location refers to the Google Cloud region where the federated catalog metadata is stored. Use the AWS to Google Cloud region mapping to find the corresponding GCP region for your AWS Region. For example, AWS us-east-1 maps to GCP us-east4.

Retrieve the BigLake service account ID

After catalog creation, Google provisions a dedicated service account for your federated catalog. Retrieve its numeric ID:

BIGLAKE_SA_ID=$(gcloud alpha biglake iceberg catalogs describe <FEDERATED_CATALOG_NAME> \
    --project="<GCP_PROJECT_ID>" \
    --format="value(biglake-service-account-id)")
echo $BIGLAKE_SA_ID

Update the AWS trust policy

Back on AWS, replace the placeholder in the IAM role’s trust policy with the actual service account ID:

aws iam update-assume-role-policy \
    --role-name bigquery-cross-cloud-role \
    --policy-document '{
      "Version": "2012-10-17",
      "Statement": [{
        "Effect": "Allow",
        "Principal": {
          "Federated": "arn:aws:iam::<AWS_ACCOUNT_ID>:oidc-provider/accounts.google.com"
        },
        "Action": "sts:AssumeRoleWithWebIdentity",
        "Condition": {
          "StringEquals": {
            "accounts.google.com:sub": ["<BIGLAKE_SA_ID>"],
            "accounts.google.com:aud": ["<BIGLAKE_SA_ID>"]
          }
        }
      }]
    }'

Register the service account ID in the OIDC provider’s audience list. Without this step, AWS rejects the token because the aud claim doesn’t match any registered client:

aws iam add-client-id-to-open-id-connect-provider \
    --open-id-connect-provider-arn "arn:aws:iam::<AWS_ACCOUNT_ID>:oidc-provider/accounts.google.com" \
    --client-id "<BIGLAKE_SA_ID>"

Set up metadata sync

Wait 3–5 minutes for IAM changes to propagate globally, then set up background refresh:

gcloud alpha biglake iceberg catalogs update <FEDERATED_CATALOG_NAME> \
    --project="<GCP_PROJECT_ID>" \
    --refresh-interval=300s

The --refresh-interval (300 seconds in this example) determines how often BigQuery syncs metadata from the AWS Glue IRC. New tables and schema changes appear in BigQuery within this interval.

Querying from BigQuery

After the catalog refresh completes, BigQuery automatically creates external datasets corresponding to the synced namespaces. No manual CREATE SCHEMA is required.

Verify the sync:

gcloud alpha biglake iceberg namespaces list \
    --catalog="<FEDERATED_CATALOG_NAME>" \
    --project="<GCP_PROJECT_ID>"

Run a query in BigQuery:

SELECT * FROM `<GCP_PROJECT_ID>.<FEDERATED_CATALOG_NAME>.<NAMESPACE>.orders` LIMIT 1000

Sample Query Output:

SELECT
    customer_id,
    COUNT(*) as order_count,
    SUM(amount) as total_spend
FROM `<GCP_PROJECT_ID>.<FEDERATED_CATALOG_NAME>.<NAMESPACE>.orders`
GROUP BY customer_id
ORDER BY total_spend DESC
BigQuery query results showing order count and total spend per customer from the Amazon S3 Tables data

Figure 5: BigQuery query results returned through Lake Formation credential vending

BigQuery reads the Iceberg metadata to identify which Parquet data files contain relevant data. It also applies partition pruning where applicable, and fetches only the necessary files from S3 Tables managed storage.

Schema evolution

When new columns are added to an Iceberg table on the AWS side (through Spark, Athena, or the AWS Glue IRC), the schema change is captured in Iceberg’s metadata. On the next Lakehouse refresh cycle, BigQuery picks up the new columns automatically. No DDL changes are needed in BigQuery.

Metadata freshness

The s3tablescatalog catalog in AWS Glue is a federated catalog that resolves table metadata live from the S3 Tables service on each request. When a streaming job commits new data to an S3 Table, the latest metadata is immediately available through the AWS Glue IRC. BigQuery sees the update on its next refresh cycle (as configured by --refresh-interval).

OIDC identity federation

The trust relationship between Google Cloud and AWS uses OpenID Connect. When BigQuery Lakehouse needs to access your data, it presents a signed JWT token containing:

  • iss: accounts.google.com (the issuer)
  • sub: The BigLake service account ID (identifies which catalog is making the request)
  • aud: The same service account ID (the intended audience)

AWS validates this token against the registered OIDC provider and trust policy conditions before issuing temporary credentials. Each federated catalog receives a unique service account ID, providing per-catalog isolation and auditability through AWS CloudTrail.

Network path

By default, traffic between BigQuery and AWS travels over the public internet. For workloads requiring private connectivity, Google Cloud supports Cross-Cloud Interconnect or Partner Interconnect. This helps routing queries over a dedicated network path. Refer to the Google Cloud documentation for private interconnect configuration.

Clean up

To avoid ongoing charges, remove the resources created in this walkthrough.

On AWS:

# Delete the table (if created for this walkthrough)
aws s3tables delete-table \
    --table-bucket-arn "arn:aws:s3tables:<AWS_REGION>:<AWS_ACCOUNT_ID>:bucket/<TABLE_BUCKET>" \
    --namespace analytics --name orders --region <AWS_REGION>

# Delete namespace and table bucket
aws s3tables delete-namespace \
    --table-bucket-arn "arn:aws:s3tables:<AWS_REGION>:<AWS_ACCOUNT_ID>:bucket/<TABLE_BUCKET>" \
    --namespace <NAMESPACE> --region <AWS_REGION>

aws s3tables delete-table-bucket --name <TABLE_BUCKET> --region <AWS_REGION>

# Delete IAM role and OIDC provider (if no longer needed)
aws iam delete-role --role-name bigquery-cross-cloud-role

On Google Cloud:

gcloud alpha biglake iceberg catalogs delete <FEDERATED_CATALOG_NAME> \
    --project="<GCP_PROJECT_ID>" --location=<GCP_REGION>

Conclusion

This post demonstrated how to query Amazon S3 Tables from Google BigQuery using AWS Lake Formation credential vending, where Lake Formation manages the permissions and issues temporary, scoped credentials for data access. With the open Iceberg format, you can write data once on AWS and read it from supported engines that speak Iceberg, including BigQuery.

Together with the IAM approach covered in Part 1, two access control modes provide flexibility: IAM for teams who want a straightforward setup and Lake Formation for organizations with complex governance requirements where multiple engines need centrally managed access to the same data.

To get started with this pattern in your environment:


About the authors

Lakshmi Nair

Lakshmi Nair

Lakshmi is a Principal Analytics Specialist Solutions Architect at AWS. She specializes in designing advanced analytics systems across industries. She focuses on crafting cloud-based data platforms, enabling real-time streaming, big data processing, and robust data governance.

Srividya Parthasarathy

Srividya Parthasarathy

Srividya was a Senior Big Data Architect on the AWS Lake Formation team. She works with product team and customer to build robust features and solutions for their analytical data platform. She enjoys building data mesh solutions and sharing them with the community.

Track SageMaker Unified Studio project costs with custom tags and AWS CUR

Post Syndicated from Nisha Gambhir original https://aws.amazon.com/blogs/big-data/track-sagemaker-unified-studio-project-costs-with-custom-tags-and-aws-cur/

Organizations running machine learning (ML), analytics, and generative AI workloads on Amazon SageMaker Unified Studio domains and projects face a common cost governance challenge. System tags (AmazonDataZoneDomainId and AmazonDataZoneProject) are automatically propagated to all underlying project resources. However, custom tags such as CostCenter, Team, or Environment are not propagated to dynamic resources created through the Studio UI. This creates a gap when you need to report project costs grouped by custom tags.

In this post, we walk through a serverless solution that bridges this gap by enriching AWS Cost and Usage Report (CUR) data with custom project tags. By the end of this post, you can build an Amazon Quick Sight dashboard to filter and analyze Amazon SageMaker Unified Studio project costs by any custom tag dimension that you define. This gives your team the visibility to make informed spending decisions.

Solution overview

The solution consists of three automated subsystems:

  1. Event-driven tag lookup management – An Amazon EventBridge rule captures Amazon DataZone project lifecycle events (Create, Update, Delete) and triggers an AWS Lambda function. The function maintains an Amazon DynamoDB lookup table that maps each project’s DomainId and ProjectId to its custom tags.
  2. CUR enrichment pipeline – An AWS Glue extract, transform, and load (ETL) job reads CUR 2.0 Parquet data from Amazon Simple Storage Service (Amazon S3). The job joins each billing line item with the DynamoDB lookup table using the system tags (DomainId, ProjectId), appends the custom tag values as new columns, and writes the enriched data back to Amazon S3.
  3. Cost visualization – An Amazon Quick Sight dashboard backed by a custom SQL dataset over Amazon Athena provides interactive cost and consumption analytics filtered by custom tags.

Architecture

The following diagram shows the end-to-end architecture:

Figure 1: SageMaker Unified Studio project custom tag cost reporting

The workflow is as follows:

  • An Amazon SageMaker Unified Studio administrator creates or updates a project with custom tags.
  • AWS CloudTrail captures the API call.
  • Amazon EventBridge matches the event.
  • The Lambda orchestrator writes the tag mapping to DynamoDB.
  • Separately, AWS Data Exports delivers CUR data to Amazon S3.
  • The AWS Glue ETL job enriches CUR line items with custom tags from DynamoDB.
  • The AWS Glue Crawler catalogs the enriched data.
  • Amazon Quick Sight visualizes costs by custom tags.

Prerequisites

Before deploying this solution, you need:

  • An Amazon SageMaker Unified Studio domain (you create projects after deployment).
  • AWS Cloud Development Kit (AWS CDK) CLI installed.
  • Python 3.12+.
  • Amazon Quick Sight Enterprise edition enabled in your account.
  • An AWS Identity and Access Management (IAM) user or role with permissions to deploy AWS CloudFormation stacks.

Step 1: Configure custom tags on your project profile

You configure custom tags on project profiles through the Amazon DataZone API. First, enable custom tags on your project profile:

aws datazone update-project-profile \
  --domain-identifier $DOMAIN_ID \
  --identifier $PROJECT_PROFILE_ID \
  --region $REGION \
  --allow-custom-project-resource-tags \
  --project-resource-tags '[
  {"key": "CostCenter", "value": "default", "isValueEditable": true},
  {"key": "Team", "value": "default", "isValueEditable": true},
  {"key": "Environment", "value": "default", "isValueEditable": true}
]'

When creating or updating a project, set the tag values:

aws datazone update-project \
  --domain-identifier $DOMAIN_ID \
  --identifier $PROJECT_ID \
  --project-profile-version latest \
  --region $REGION \
  --resource-tags '{"CostCenter": "CC-100", "Team": "ML-Platform", "Environment": "Production"}'

Important: The AmazonSageMakerProvisioning-<domainAccountId> role needs an inline policy that permits your custom tag keys. Without this, project environment deployment fails.

The following is the inline policy that’s used for the custom tags shared in this post:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "AllowCustomTagKeys",
      "Effect": "Allow",
      "Action": [
        "sagemaker:AddTags",
        "sagemaker:DeleteTags",
        "cloudformation:TagResource",
        "cloudformation:CreateStack",
        "cloudformation:UpdateStack"
      ],
      "Resource": "*",
      "Condition": {
        "ForAnyValue:StringLike": {
          "aws:TagKeys": [
            "AmazonDataZone*",
            "CostCenter",
            "Team",
            "Environment"
          ]
        }
      }
    }
  ]
}

Step 2: Activate cost allocation tags

Activate the SageMaker Unified Studio system tags as cost allocation tags so they appear in CUR data:

aws ce update-cost-allocation-tags-status \
  --cost-allocation-tags-status '[
  {"TagKey": "AWSDataZoneProject", "Status": "Active"},
  {"TagKey": "AmazonDataZoneDomainId", "Status": "Active"}
]'

These tags take up to 24 hours to start appearing in CUR reports after activation.

Step 3: Deploy the infrastructure

The solution is packaged as a CDK application. Clone the GitHub repository and deploy:

# Install dependencies
pip install -r requirements.txt

# Bootstrap CDK (first time only)
cdk bootstrap aws://$ACCOUNT_ID/$REGION

# Deploy
cdk deploy

This creates the following resources:

  • DynamoDB table (smus-project-tag-lookup) – stores project-to-tag mappings.
  • Lambda function (smus-orchestrator) – processes project lifecycle events.
  • Amazon EventBridge rule – matches Amazon DataZone CreateProject/UpdateProject/DeleteProject events.
  • S3 buckets – for raw CUR and enriched CUR data.
  • AWS Glue ETL job (smus-cur-enrichment) – enriches CUR with custom tags.
  • AWS Glue Crawler – catalogs enriched data.
  • Amazon Simple Notification Service (Amazon SNS) topic – pipeline failure alerts.

Note: The solution uses serverless components (Lambda, DynamoDB on-demand, AWS Glue, Amazon Quick Sight), so you only pay for what you use. The primary cost drivers are AWS Glue ETL job execution time and Amazon Quick Sight SPICE storage.

Step 4: Configure CUR delivery

Create a CUR 2.0 export through AWS Data Exports that delivers Parquet files to the CUR S3 bucket created by the stack. The following screenshots show the complete configuration process in the AWS Billing and Cost Management console.

To create the export, follow these steps:

  1. Go to AWS Billing and Cost Management and then choose Data Exports.
  2. Choose Create in the upper right corner of the Exports and dashboards page. The Data Exports console shows any existing exports, their status, export type, data table, and last refresh date.
  3. On the Create export page, under Export details, select Standard data export and enter an export name. Under Data table content settings, select CUR 2.0.
  4. Under Data table configurations, set Time granularity to Hourly. The configuration page also lets you choose additional export content options such as including resource IDs, split cost allocation data, caller identity allocation data, and capacity reservation columns.
  5. Under Data export delivery options, set Compression type and file format to Parquet. Under Data export storage settings, configure the S3 bucket to: smus-cur-report-{account-id}-{region} and set the S3 path prefix as needed. Choose Create to finish.
Data Exports console listing existing exports with status, type, and last refresh date

Figure 2: Data Exports page listing existing exports

Create export page with Standard data export selected and CUR 2.0 chosen

Figure 3: Create export page with Standard data export and CUR 2.0 selected

Data table configurations with time granularity set to Hourly

Figure 4: Data table configurations with time granularity set to Hourly

Data export delivery options with Parquet format and the S3 storage destination configured

Figure 5: Data export delivery options with Parquet format and S3 storage settings

Step 5: How the event-driven tag capture works

When a project is created or updated in Amazon SageMaker Unified Studio (through the Studio UI or API), the following happens automatically:

  1. CloudTrail logs the Amazon DataZone API call.
  2. Amazon EventBridge matches the event.
  3. Amazon EventBridge invokes the Lambda function.
  4. The Lambda extracts custom tags from the CloudTrail event payload.
  5. The Lambda writes a record to DynamoDB with the DomainId, ProjectId, and all custom tag key-value pairs.

The Lambda function reads tags directly from the responseElements.resourceTags field of the CloudTrail event rather than making a separate GetProject API call. This avoids a race condition where GetProject might return empty tags while the project is in the UPDATING state.

def _extract_tags_from_event(detail):
    tags = {}
    response_elements = detail.get("responseElements") or {}
    for tag_entry in response_elements.get("resourceTags", []):
        if isinstance(tag_entry, dict) and "key" in tag_entry:
            tags[tag_entry["key"]] = tag_entry["value"]
    request_params = detail.get("requestParameters") or {}
    req_tags = request_params.get("resourceTags", {})
    if isinstance(req_tags, dict):
        tags.update(req_tags)
    return tags

Step 6: How the CUR enrichment works

The AWS Glue ETL job runs on a schedule (after each CUR delivery):

  1. Reads CUR Parquet files from the CUR S3 bucket.
  2. Reads all records from the DynamoDB lookup table.
  3. Performs a left outer join on DomainId and ProjectId.
  4. Appends custom tag columns (CostCenter, Team, Environment, and so on) to each CUR line item.
  5. Writes enriched Parquet to the enriched S3 bucket.

Line items without a matching project in the lookup table retain all original columns with NULL custom tag values. No data is dropped.

joined_df = cur_df.join(
    lookup_df,
    on=(
        (cur_df[DOMAIN_COL] == lookup_df["domainId"])
        & (cur_df[PROJECT_COL] == lookup_df["projectId"])
    ),
    how="left_outer",
)

Step 7: Set up the Amazon Quick Sight dashboard

After the first ETL run and crawler execution, set up the Amazon Quick Sight dashboard:

python scripts/setup_quicksight.py \
  --account-id $ACCOUNT_ID \
  --region $REGION \
  --quicksight-user $QUICKSIGHT_USER_ARN

This creates a dashboard with five visuals:

  • Cost by Custom Tag (CostCenter) – horizontal bar chart.
  • Cost by Project – horizontal bar chart.
  • Daily Cost Trend – line chart.
  • Cost by Service per Project – stacked bar chart.
  • Usage by Project & Service – summary table.

And six interactive list filters: Domain, Project, CostCenter, Team, Environment, Service.

The custom SQL includes a CASE statement for service categorization:

SELECT
  line_item_usage_start_date,
  line_item_product_code,
  line_item_usage_amount,
  line_item_unblended_cost,
  resource_tags_user_amazondatazone_domain_id AS domain_id,
  resource_tags_user_amazondatazone_project AS project_id,
  costcenter, team, environment,
  CASE
    WHEN line_item_product_code = 'AmazonSageMaker' THEN 'SageMaker'
    WHEN line_item_product_code = 'AmazonS3' THEN 'S3'
    WHEN line_item_product_code = 'AWSGlue' THEN 'Glue'
    ELSE line_item_product_code
  END AS service_category
FROM "smus_cost_reporting"."enriched_cur"
WHERE line_item_unblended_cost > 0

Step 8: Verifying the solution

After deploying the infrastructure and setting up the dashboard, verify that each component of the pipeline is functioning correctly.

8.1 Verify Amazon EventBridge is capturing project events

  1. Open the Amazon EventBridge console.
  2. In the navigation pane, choose Rules.
  3. Select the rule created by the CDK stack (for example, SmusCostReporting-ProjectTagRule).
  4. Choose the Monitoring tab.
  5. Confirm that the invocations are being recorded in the metrics.
  6. Create or update an Amazon SageMaker Unified Studio project with custom tags using the following command:
    aws datazone update-project \
      --domain-identifier <domain-id> \
      --identifier <project-id> \
      --custom-tags CostCenter=Engineering Team=DataPlatform Environment=Production

  7. Within a few seconds, the Amazon EventBridge rule should show a new invocation in its metrics.

8.2 Verify DynamoDB schema and tag mappings

The DynamoDB lookup table uses a simple key schema:

Attribute Type Role
domainId String Partition Key
projectId String Sort Key
CostCenter String Custom tag
Team String Custom tag
Environment String Custom tag

Custom tags are stored as dynamic attributes. Any tag key set on a project becomes a column in the table.

8.2.1 Verify DynamoDB table contains tag mappings

  1. Open the DynamoDB console.
  2. Navigate to the table created by the stack (for example, SmusCostReporting-ProjectTagsTable).
  3. Choose Explore table items.
  4. Scan for your project with the following keys:
    Partition key (domainId): <your-domain-id>
    Sort key (projectId): <your-project-id>

  5. Confirm the item contains the expected custom tag attributes (CostCenter, Team, Environment) with the values you assigned.
  6. Alternatively, use the AWS CLI:
    aws dynamodb get-item \
      --table-name SmusCostReporting-ProjectTagsTable \
      --key '{"domainId": {"S": "<domain-id>"}, "projectId": {"S": "<project-id>"}}'

8.3 Verify the AWS Glue ETL job enriches CUR data

  1. Wait for the next CUR delivery (hourly if configured as described in Step 4).
  2. Wait for the subsequent AWS Glue job execution.
  3. Open the AWS Glue console.
  4. In the navigation pane, choose ETL Jobs.
  5. Confirm the job completed successfully (status: Succeeded).
  6. Query the enriched data in Amazon Athena to confirm custom tag columns are populated:
    SELECT
      line_item_usage_start_date,
      line_item_product_code,
      line_item_unblended_cost,
      costcenter,
      team,
      environment
    FROM "smus_cost_reporting"."enriched_cur"
    WHERE costcenter IS NOT NULL
    LIMIT 10;

You should see rows with your custom tag values populated in the costcenter, team, and environment columns.

8.4 Verify the Amazon Quick Sight dashboard displays enriched data

  1. Open the Amazon Quick Sight console and navigate to the dashboard created by the setup script.
  2. Confirm that:
    • The Cost by Custom Tag (CostCenter) bar chart displays cost data grouped by your CostCenter values.
    • The list filters for CostCenter, Team, and Environment contain selectable values.
    • Selecting a filter value correctly narrows the displayed data.
  3. If the dashboard shows no data, verify that:
    • The AWS Glue Crawler has run after the ETL job (check the crawler’s last run status in the AWS Glue console).
    • The SPICE dataset has been refreshed. In the Amazon Quick Sight console, navigate to Datasets, select the dataset, and then choose Refresh now.

Figure 6 shows the Amazon Quick Sight dashboard with two side-by-side horizontal bar charts: Cost by Cost Center and Cost by Project. Domain Name and Project Name list filters appear at the top.

Amazon Quick Sight dashboard with Cost by Cost Center and Cost by Project bar charts and Domain and Project filters

Figure 6: Amazon Quick Sight dashboard showing cost data by custom tags, including Cost by Cost Center and Cost by Project bar charts with Domain Name and Project Name filters

Note: The first end-to-end cycle can take up to 48 hours depending on CUR delivery timing. After the initial cycle completes, subsequent updates will flow automatically on the configured schedule.

Operational considerations

Monitoring: The Amazon SNS topic smus-cost-reporting-alerts receives notifications when the AWS Glue ETL job fails or the Lambda orchestrator encounters repeated errors. Subscribe an email address or Slack webhook to stay informed. For instructions on how to create a subscription, see Subscribing to an Amazon SNS topic.

Cost: The solution uses serverless components (Lambda, DynamoDB on-demand, AWS Glue, Amazon Quick Sight, SPICE) so you only pay for what you use. The primary cost drivers are AWS Glue ETL job execution time and Amazon Quick Sight SPICE storage.

Scaling: The DynamoDB table uses on-demand capacity and can scale to accommodate your projects. You can scale the AWS Glue ETL job by increasing the number of workers for larger CUR datasets. For more information, see Managing throughput capacity automatically with DynamoDB auto scaling.

New tag keys: When you add new custom tag keys to projects, the ETL automatically picks them up as new columns. The AWS Glue Crawler’s UPDATE_IN_DATABASE policy adds new columns to the catalog table without manual intervention.

Cleanup

Warning: The following cleanup steps will permanently delete all CUR data, project tag mappings, and Amazon Quick Sight dashboards.

To remove all resources:

# Delete Amazon Quick Sight resources
python scripts/setup_quicksight.py --account-id $ACCOUNT_ID --region $REGION --quicksight-user $QS_USER --clean

# Delete CDK stack
cdk destroy

Go to AWS Billing and Cost Management, and then choose Data Exports and delete the CUR 2.0 export created in Step 4.

Deactivate the cost allocation tags that were activated in Step 2:

aws ce update-cost-allocation-tags-status \
  --cost-allocation-tags-status '[
  {"TagKey": "AWSDataZoneProject", "Status": "Inactive"},
  {"TagKey": "AmazonDataZoneDomainId", "Status": "Inactive"}
]'

Conclusion

In this post, we showed how to build an end-to-end cost reporting solution for Amazon SageMaker Unified Studio projects using custom tags. This solution combines tag capture driven by Amazon EventBridge, CUR enrichment through AWS Glue ETL, and visualization in Amazon Quick Sight. With it, organizations can track and attribute costs by CostCenter, Team, Environment, or any custom dimension. This works even for resources created through the Studio UI that don’t receive custom tag propagation.

This solution serves as an extension to the custom tag propagation feature and reports cost for all project resources. The architecture is fully serverless, automated, and can be deployed to any AWS account using the provided CDK application.

To start building your custom tag cost reporting pipeline, visit the GitHub repository. To learn more about the underlying services, visit the Amazon SageMaker Unified Studio service page. For a related approach to custom tag governance, see Use Amazon SageMaker custom tags for project resource governance and cost tracking

References


About the authors

Nisha Gambhir

Nisha Gambhir

Nisha is a Senior AI/ML & Cloud Architect based out of India. She is passionate about helping customers design, architect and develop secure, scalable and reliable applications using AI/ML and Agentic AI. She loves working on latest technologies, providing simple and scalable solutions that drive positive business outcomes.

Dr Anil Giri

Dr Anil Giri

Anil is a Solutions Architect at AWS, based in London, UK, where he helps ISV customers design and deploy agentic AI systems in production. He specializes in multi-agent orchestration, retrieval-augmented generation, and event-driven serverless architectures on Amazon Bedrock, with a focus on building reliable, secure, and scalable solutions that deliver measurable business outcomes.

Satish Sarapuri

Satish Sarapuri

Satish is a Sr. Data Architect, Data Mesh / Data Lake/Gen AI at AWS. He helps enterprise-level customers build high-performance, highly available, cost-effective, resilient, and secure generative AI, data mesh, data lake, and analytics platform solutions on AWS, through which customers can make data-driven decisions to gain impactful outcomes for their business and help them on their digital and data transformation journey. In his spare time, he enjoys trail running and spending quality time with his family.

Ram Vittal

Ram Vittal

Ram is a Principal GenAI/ML Specialist at AWS. He has over 3 decades of experience building distributed, hybrid, and cloud applications. He is passionate about building secure, scalable, reliable AI/ML and big data solutions to help customers with their cloud adoption and optimization journey. In his spare time, he rides motorcycle and enjoys the nature with his family.

Secure SageMaker Unified Studio access with SAML and conditional policies

Post Syndicated from Manos Samatas original https://aws.amazon.com/blogs/big-data/secure-sagemaker-unified-studio-access-with-saml-and-conditional-policies/

Amazon SageMaker Unified Studio is a single data and AI development environment that brings together data preparation, analytics, and machine learning (ML) development in one place. By unifying these workflows, it saves teams from managing multiple tools and makes it straightforward for data scientists, analysts, and developers to build, train, and deploy ML models while collaborating. In Amazon SageMaker Unified Studio, a domain is the organizing entity for connecting your assets, users, and their projects. With Amazon SageMaker unified domains, you have the flexibility to reflect the data and analytics needs of your organizational structure. You can create a single unified domain for your enterprise or multiple domains for different business units.

Some enterprises, especially those in regulated industries, might require limiting access to trusted networks (such as VPN CIDRs) or to managed devices that meet compliance standards through device attestation.

In this post, we demonstrate how to integrate SageMaker Unified Studio as a custom SAML application and apply conditional access policies for enforcing device compliance, IP-based restrictions, or multi-factor authentication (MFA). For this post, we use Okta as the identity provider (IdP).

Solution overview

This solution demonstrates how to integrate Amazon SageMaker Unified Studio (SMUS) with external SAML identity providers such as Okta. The integration enforces enterprise security controls, including trusted network access, device compliance, and multi-factor authentication. With this integration, organizations in regulated industries can maintain strict access controls while providing single sign-on for their data science and AI development teams. By using SAML 2.0 federation with conditional access policies, you can help make sure that only authenticated users on compliant devices from trusted networks gain access. This access applies to your SageMaker Unified Studio domains and the associated data and AI workloads.

SAML authentication flow from a corporate device through the identity provider and AWS STS to Amazon SageMaker Unified Studio

Authentication flow for accessing SageMaker Unified Studio through SAML

The architecture diagram illustrates the secure authentication flow for accessing SageMaker Unified Studio through SAML integration:

  1. Users typically initiate access from corporate-managed devices through VPN or trusted network connections.
  2. The IdP authenticates the user and evaluates conditional access policies defined by your organization. Based on these policies, it checks for trusted devices, approved source IP ranges, and MFA completion. If any policy fails, the login is rejected. Otherwise, authentication proceeds.
  3. Upon successful authentication and policy validation, the IdP generates a digitally signed SAML assertion containing user attributes and group memberships, securely delivering it to the user’s browser through HTTP POST binding.
  4. The client browser automatically posts the SAML assertion to the AWS Security Token Service (AWS STS) sign-in endpoint. There, the AWS IAM Identity Provider validates the trust relationship with your corporate IdP through pre-configured SAML federation settings.
  5. AWS STS validates the SAML assertion signature and authenticity. It then maps the user attributes to a specifically configured IAM role with SageMaker Unified Studio permissions, including the datazone:GetIamPortalLoginUrl permission required for domain access.
  6. AWS STS confirms successful role assumption and generates temporary AWS credentials with a defined session duration. It then issues an HTTP redirect that returns the browser to the SageMaker Unified Studio domain with authenticated session tokens.
  7. Users gain access to the unified environment for data preparation, analytics, and machine learning development. All activities are governed by the assumed IAM role permissions and logged for comprehensive audit trails.

Walkthrough

In this walkthrough, you create a SAML application in Okta, connect it to AWS, and configure a SageMaker Unified Studio domain to use it for authentication.

Prerequisites

Before you get started, make sure you have the following:

  1. Familiarity with Amazon SageMaker Unified Studio.
  2. A basic understanding of SAML 2.0.
  3. AWS Identity and Access Management (IAM) permissions to create a domain in Amazon SageMaker Unified Studio.
  4. Access to your SAML IdP (such as Okta or Entra ID) to create and configure a SAML application.

Step 1: Create an application in Okta

The first step is to set up a new SAML application in Okta that manages authentication for SMUS.

  1. In Okta, go to Applications → Create App Integration, and choose SAML 2.0.
  2. Provide an App name.
  3. Set the Single sign-on URL to https://signin.aws.amazon.com/saml.
  4. Set Name ID format to Persistent.
  5. Set the Audience URI (SP Entity ID) to https://signin.aws.amazon.com/saml.
  6. Choose Next, and finish creating the application.
  7. Once created, copy the Metadata URL and Sign On URL. You need these in later steps.

Step 2: Create an identity provider in IAM

Now, let’s connect Okta to AWS by creating an IAM identity provider. This allows AWS to trust authentication responses from Okta.

  1. Open the IAM console.
  2. Go to Identity providers → Add provider.
  3. Select SAML as the provider type.
  4. Provide a Provider name.
  5. In Okta, go to your application’s Sign On tab, choose Identity Provider metadata, and save the XML file. Upload it here.
  6. Choose Add provider.
  7. Copy the ARN of this provider. You need it when you create the role.

Step 3: Create an IAM role for Okta

Next, create an IAM role that Okta can assume. This role defines what access users have when they sign in through Okta.

  1. In IAM, go to Roles → Create role.
  2. Use the following trust policy (replace both instances of “{Replace with Identity provider ARN}” with the ARN you copied in Step 2):
{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Principal": {
                "Federated": "{Replace with Identity provider ARN}"
            },
            "Action": "sts:AssumeRoleWithSAML",
            "Condition": {
                "StringEquals": {
                    "SAML:aud": "https://signin.aws.amazon.com/saml"
                }
            }
        },
        {
            "Effect": "Allow",
            "Principal": {
                "Federated": "{Replace with Identity provider ARN}"
            },
            "Action": "sts:TagSession",
            "Condition": {
                "StringLike": {
                    "aws:RequestTag/Email": "*"
                }
            }
        }
    ]
}
  1. Attach a permission policy. For example:
{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "VisualEditor0",
            "Effect": "Allow",
            "Action": "datazone:GetIamPortalLoginUrl",
            "Resource": "arn:aws:datazone:<REGION>:<ACCOUNT-ID>:domain/<DOMAIN-ID>"
        }
    ]
}

Replace <REGION>, <ACCOUNT-ID>, and <DOMAIN-ID> with the corresponding values from your SageMaker Unified Studio domain ARN (arn:aws:sagemaker:<REGION>:<ACCOUNT-ID>:domain/<DOMAIN-ID>). You can find the domain ARN in the SageMaker console under Domains.

Step 4: Configure SAML assertions

To make sure AWS understands who is signing in, configure the SAML assertions in Okta.

  1. Open your application in Okta.
  2. Go to General → SAML Settings → Edit.
  3. Choose Next until you reach Attribute Statements.
  4. Add the following mappings:
    • https://aws.amazon.com/SAML/Attributes/PrincipalTag:Email → user.email.
    • https://aws.amazon.com/SAML/Attributes/Role → {IAMROLEARN,IdentityProviderARN}.
    • https://aws.amazon.com/SAML/Attributes/RoleSessionName → user.email.

Step 5: Create an SMUS domain

Finally, let’s set up the SMUS domain and tie it all together.

Note: Creating a SageMaker Unified Studio domain incurs charges. For pricing details, see the Amazon SageMaker pricing page.

  1. Open the Amazon SageMaker console.
  2. Choose Create domain.
  3. Choose Manual setup (this allows for SAML integration).
  4. Enter a domain name, then choose Create.
  5. In Configure SSO user access, select SAML, then choose Next.
  6. Set the IdP SSO URL to the Sign On URL from Step 1.
  7. Select Do not require assignments. (Access is instead managed by your IdP team through Okta or Entra.)
  8. Choose Next, then choose Save.

To verify the integration works, open your SMUS domain and choose Sign in with SSO. You are redirected to Okta, and conditional access policies such as VPN, device attestation, or MFA apply automatically.

  1. Open your SMUS domain URL in a browser.
  2. Choose Sign in with SSO.
  3. Confirm that you are redirected to Okta for authentication.
  4. Sign in with your Okta credentials.
  5. Verify that you are redirected back to the SMUS domain with access to your projects.

Step 6: Assign users to the Okta application

Before users can authenticate through Okta to access SMUS, you must assign them to the application.

  1. In Okta, navigate to your SAML application.
  2. Go to the Assignments tab.
  3. Choose Assign, and select Assign to People or Assign to Groups.
  4. Select the users or groups who need access to SMUS.
  5. Choose Save and Go Back, then choose Done.

Step 7: Apply conditional access policies

Up to Step 5, we configured SMUS with an external SAML IdP. At this point, anyone assigned to the new application in your IdP can sign in and access the SMUS domain.

This is where conditional access policies come into play. Based on your organization’s governance model, you can add policies in your IdP to further control how and when users gain access. For example:

  • Restricting access to specific corporate IP address ranges (for example, only through VPN).
  • Enforcing device compliance so that only managed or secure devices can connect.
  • Adding MFA requirements for sensitive actions.
  • Applying device attestation to help assess whether the endpoint conforms to security baselines.

Most major IdPs, including Okta and Entra ID, support conditional access. You can find more details in their documentation:

These policies allow you to enforce the right level of protection, from something as simple as requiring users to connect through corporate networks to something as advanced as verifying device attestation across your fleet.

Clean up

To avoid incurring ongoing charges, delete the resources you created during this walkthrough:

  1. Delete the Amazon SageMaker Unified Studio domain from the SageMaker console.
  2. Delete the IAM role you created for Okta.
  3. Delete the IAM identity provider.
  4. Delete the SAML application in Okta.

Important: Deleting the SMUS domain permanently removes all projects, assets, and data within it. Back up any important work before proceeding.

Conclusion

By integrating SMUS with an external IdP through SAML, you can help enforce modern access controls based on your organization’s security requirements. This post walked through how to configure SMUS with a custom SAML application and pointed you toward resources for setting up conditional access policies.

With conditional access in place, you can decide, based on your organization’s needs, whether access should be limited to trusted users on trusted networks, trusted devices, or both. This approach can help provide a more secure and compliant login experience that aligns SMUS access with your company’s broader identity and security strategy.


About the authors

Amit Samal

Amit Samal

Amit is a Sr. Delivery Consultant in World Wide Public Sector, Professional Services at AWS working with UKGI Customers. Amit has been with AWS for about 4 years and has been helping customers across the UKGI to design & implement secure, resilient and cost-effective workloads on AWS. Amit is passionate about all areas of technology, but has focus areas in Networking, Migrations, and Application Modernizations.

Manos Samatas

Manos Samatas

Manos is a Principal Solutions Architect in Data and AI with Amazon Web Services. He works with government, non-profit, education and healthcare customers in the UK on data and AI projects, helping build solutions using AWS. Manos lives and works in London. In his spare time, he enjoys reading, watching sports, playing video games and socialising with friends.

Propagate user authorization context in AI agents with Amazon Bedrock AgentCore

Post Syndicated from Anshu Bathla original https://aws.amazon.com/blogs/security/propagate-user-authorization-context-in-ai-agents-with-amazon-bedrock-agentcore/

Many teams now deploy AI agents that pull from Amazon DynamoDB tables, document repositories, software as a service (SaaS) platforms, and internal knowledge bases to answer questions and automate workflows. A key risk in these deployments is that the agent has no awareness of who’s asking, so it might return data the user shouldn’t see.

If you’re using Amazon Bedrock AgentCore to build AI agents that access multiple data sources, you need each user to see only the data they’re authorized to access. In this post, you learn patterns for propagating user authorization context through your agents so access control is enforced by infrastructure and downstream services, not by agent code. In this post, we show you how to deploy agents that enforce least privilege access without writing authorization logic in the agent itself. This approach follows AGENTSEC03 best practice in the AWS Well-Architected Agentic AI Lens.

Use case

Consider an example of a customer relationship management (CRM) chat application where employees from Sales and Finance departments interact with an AI agent to access customer information. Employees use the same chat interface and the same agent, but each department needs isolated access to their respective data:

  • Sales needs access to customer contracts, pricing strategies, and sales pipeline data
  • Finance needs access to customer invoices, payment records, and financial reports

The AI agent accesses three types of data sources on behalf of users:

When a Sales employee asks, “Show me customer contracts,” the agent must retrieve only Sales department contracts, not Finance invoices. This enforcement must happen outside the agent so that even if the agent is compromised through prompt injection or application bugs, it can’t access unauthorized data.

Note: Although we use department-based scoping in this example, the pattern generalizes to any custom claim you define, whether it represents a role, business unit, geographic region, or project assignment.

Architecture overview

The following diagram shows the architecture used in this demonstration.

Figure 1: Target architecture

Figure 1: Target architecture

The data flow shown in Figure 1 includes:

  1. A user opens the chat application and authenticates with Amazon Cognito user pool , which acts as the identity provider (IdP).
  2. A pre token generation Lambda trigger (V2) enriches the JSON Web Tokens (JWTs) with a custom claim and AWS session tag metadata before returning them to the user.
  3. The web app routes the user’s request along with the access token to the agent deployed on Amazon Bedrock AgentCore Runtime.
  4. Bedrock AgentCore Runtime validates the inbound JWT and, through Bedrock AgentCore Identity, issues a workload access token that binds the user and agent identities, and then invokes the agent.
  5. For queries requiring internal documents, the agent uses its AWS Identity and Access Management (IAM) role to query Amazon Bedrock Knowledge Bases (backed by an Amazon S3 vector store) with metadata filtering, and DynamoDB with user-scoped session-tagged credentials.
  6. For queries requiring external data, Bedrock AgentCore Identity retrieves credentials from AWS Secrets Manager and performs an on-behalf-of token exchange (RFC 8693) with Salesforce, returning a user-scoped access token.
  7. The agent calls the Salesforce REST API using the user-scoped token. Salesforce applies sharing rules and returns only records the user is authorized to access.

This architecture follows two key principles.

  • The agent acts as an orchestrator, not a gatekeeper; it coordinates tool calls and reasoning but doesn’t control access to data. Authorization is enforced by downstream services.
  • The agent doesn’t store credentials to data stores; instead, each request gets temporary, user-bound access tokens.

In the following sections, we dive deep into each data source to show how these principles are achieved in practice.

Initial user authentication with IdP

When an employee opens the chat application, they authenticate using their corporate credentials. For this example, you use Amazon Cognito user pools as the IdP. You can also achieve this with other IdPs such as Entra ID or Okta.

The pre token generation Lambda trigger (V2) captures the user’s custom department context and adds it to the tokens to both the identity (ID) token and access token that Bedrock AgentCore Runtime uses for authorization decisions each serving a distinct purpose. The access token is used by the Bedrock AgentCore Runtime custom JWT authorizer for inbound authorization. The ID token also receive the https://aws.amazon.com/tags claim (used by AWS Security Token Service (AWS STS)) for session tags). The https://aws.amazon.com/tags claim is the specific format required by AWS STS to extract session tags during AssumeRoleWithWebIdentity. For more information and step-by-step guidance see How to customize access tokens in Amazon Cognito user pools.

The following example shows the key logic within a pre token generation Lambda handler function configured as a trigger on your Amazon Cognito user pool. This code runs automatically when a user authenticates, extracting their department attribute and adding it as a custom claim to both ID Token and access token.

import json

def lambda_handler(event, context):
    department = event['request']['userAttributes'].get('custom:department', '')

    event['response']['claimsAndScopeOverrideDetails'] = {
        'idTokenGeneration': {
            'claimsToAddOrOverride': {
                'department': department,
                'https://aws.amazon.com/tags': {
                    "principal_tags": {"department": [department]},
                    "transitive_tag_keys": ["department"]
                }
            }
        },
        'accessTokenGeneration': {
            'claimsToAddOrOverride': {
                'department': department
            }
        }
    }
    return event

Inbound authorization by AgentCore Runtime

When the user request reaches AgentCore Runtime, the Inbound JWT authorizer performs two checks as shown in Figure 2. It validates the JWT token with Amazon Cognito (the configured IdP) by cryptographically verifying the token’s signature, confirming it is non-expired, and checking it was issued by the trusted IdP. It then extracts the department claim from the validated token and compares it against the expected value configured in the authorizer, any token without a matching claim is rejected before the agent code is invoked.

Figure 2: Inbound JWT authorization

Figure 2: Inbound JWT authorization

The following example shows the inbound JWT authorizer configuration that you pass when deploying your agent to AgentCore Runtime. This configuration tells AgentCore which IdP to validate against and which custom claim value to enforce for this agent. In this example, inboundTokenClaimName is department, inboundTokenClaimValueType declares the claim type as STRING_ARRAY, and authorizingClaimMatchValue specifies the allowed values ([“Sales”, “Finance”]) with the CONTAINS_ANY operator. The authorizer validates that the department claim is present in the token and matches one of these values, ensuring only authenticated users from the Sales or Finance department can invoke the agent.

authorizer_config = {
        "customJWTAuthorizer": {
            "discoveryUrl": discovery_url,
            "allowedClients": [client_id],
            "customClaims": [
                {
                    "inboundTokenClaimName": "department",
                    "inboundTokenClaimValueType": "STRING_ARRAY",
                    "authorizingClaimMatchValue": {
                        "claimMatchValue": ["Sales", "Finance"]
                        "claimMatchOperator": "CONTAINS_ANY"
                    }
                }
            ]
        }
    }

Note: AgentCore Runtime automatically creates a workload identity for each deployed agent. A workload identity represents the digital identity of your agents within the AWS environment. It allows agents to maintain consistent identity whether they’re using IAM roles for AWS resource access, OAuth 2.0 tokens for external service integration, or API keys for third-party tool access.

Passing the user context for agent outbound authorization

After the inbound JWT token is validated and the user’s authorization context is confirmed, the agent must propagate this context to downstream resources. The fundamental security challenge here is how to design a system so that an agent acting on behalf of a user can only access data that user is authorized to see, even if the agent itself is compromised.

The traditional approach of granting the agent broad credentials and relying on application-level filtering (such as adding WHERE clauses to queries) creates a single point of failure. If an attacker manipulates the agent through prompt injection or exploits a bug in the filtering logic, the full dataset becomes accessible. A more resilient design moves authorization enforcement out of the agent’s application code and into the infrastructure layer wherever possible. Instead of trusting the agent to filter results correctly, you configure the underlying services—IAM policies, database access controls, SaaS sharing rules—to reject unauthorized requests regardless of what the agent asks for. This way, the agent’s credentials are inherently limited to the requesting user’s permissions, and no amount of prompt manipulation can bypass those boundaries. Where infrastructure-level enforcement isn’t yet available, such as metadata filtering in Amazon Bedrock Knowledge Bases, the agent applies application-layer controls as a complementary measure. The following sections demonstrate how this principle applies to each data source in our architecture.

Pattern 1: Scoping DynamoDB access to the requesting user

For DynamoDB access, you can use AssumeRoleWithWebIdentity with session tags to create per-request, user-scoped credentials rather than granting the agent a static IAM role with direct table access. The agent passes the user’s signed ID token to AWS STS, which extracts the department tag from the token’s https://aws.amazon.com/tags claim and returns temporary credentials constrained to that department’s data partition. This moves access control from agent code to IAM policy evaluation. STS additionally validates the token’s audience (aud) claim against the IAM OIDC provider configuration, preventing tokens issued for other app clients from being used to assume the role. The following diagram shows this flow (Figure 3).

Prerequisites (one-time setup):

Before this runtime flow can execute, complete the following configuration:

  • Register Amazon Cognito as an IAM OIDC provider. Although the user authenticates using the Cognito API (USER_PASSWORD_AUTH), STS requires Cognito to be registered as an OIDC provider so it can discover and validate ID tokens. Configure the allowed client IDs (audiences) on the provider to match your application’s app client ID.
CognitoOIDCProvider:
  Type: AWS::IAM::OIDCProvider
  Properties:
    Url: !Sub 'https://cognito-idp.${AWS::Region}.amazonaws.com/${CognitoUserPoolId}'
    ClientIdList:
      - !Ref CognitoAppClientId
    ThumbprintList:
      - '<thumbprint>'

  • Configure the UserScopedDynamoDBRole trust policy to include both sts:AssumeRoleWithWebIdentity and sts:TagSession permissions, with the Amazon Cognito OIDC provider as the federated principal.
{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": {
      "Federated": "arn:aws:iam::111122223333:oidc-provider/cognito-idp.us-east-1.amazonaws.com/us-east-1_EXAMPLE"
    },
    "Action": [
      "sts:AssumeRoleWithWebIdentity",
      "sts:TagSession"
    ],
    "Condition": {
      "StringEquals": {
        "cognito-idp.us-east-1.amazonaws.com/us-east-1_EXAMPLE:aud": "<app-client-id>"
      }
    }
  }]
}

  • By default, AgentCore Runtime drops custom headers as a security measure. To allow the X-Id-Token header through to the agent container, configure it in the agent runtime’s requestHeaderAllowlist so the ID token is forwarded to agent code. The following configuration tells AgentCore Runtime to forward only the X-Id-Token header to agent code, dropping other non-standard headers:
request_header_config = {
    'requestHeaderAllowlist': ['X-Id-Token']
}

How it works:

  1. The user navigates the web application.
  2. The user authenticates with Amazon Cognito using USER_PASSWORD_AUTH.
  3. The JWT is issued with a custom department claim and the https://aws.amazon.com/tags claim for STS session tagging (covered in the preceding Initial user authentication with IdP section).
  4. Amazon Cognito returns the enriched tokens to the frontend. The access token carries the department claim for inbound authorization. The ID token carries both the department claim and the https://aws.amazon.com/tags claim for downstream STS calls.
  5. The user asks the agent a question (for example, “Show Q4 sales pipeline”).
  6. The frontend calls AgentCore Runtime, passing two tokens: the Amazon Cognito access token in the Authorization header (for inbound authorization), and the user’s ID token as a custom X-Id-Token header (for downstream STS calls).
  7. AgentCore Runtime validates the JWT and verifies the department claim matches the allowed values configured in the inbound authorizer. If validation fails, the request is rejected with HTTP 401 before agent code executes. After validation, AgentCore forwards the request to the agent container along with the allowed X-Id-Token header.
  8. The agent calls sts:AssumeRoleWithWebIdentity with the ID token. This call targets a single shared UserScopedDynamoDBRole. The following is the agent code for this step:
    def _scoped_dynamodb_resource(id_token: str):
        """Assume user-scoped role and return DynamoDB resource."""
        sts = boto3.client('sts')
        response = sts.assume_role_with_web_identity(
            RoleArn=USER_SCOPED_DYNAMODB_ROLE_ARN,
            RoleSessionName="agent-user-session",
            WebIdentityToken=id_token,
            DurationSeconds=900
        )
        creds = response['Credentials']
        session = boto3.Session(
            aws_access_key_id=creds['AccessKeyId'],
            aws_secret_access_key=creds['SecretAccessKey'],
            aws_session_token=creds['SessionToken']
        )
        return session.resource('dynamodb')

  9. AWS STS validates the token against the Amazon Cognito OIDC provider registered in IAM. STS verifies the token’s cryptographic signature, expiration, issuer, and audience (aud). The aud claim in the ID token must match one of the client IDs configured on the IAM OIDC provider resource. This prevents a valid token issued by the same Cognito user pool but for a different app client from being accepted. Note that the agent’s own execution role has no DynamoDB access and only permits sts:AssumeRoleWithWebIdentity, so even a compromised agent can’t bypass this flow.

    Note: Amazon Cognito user pools expose a standard OpenID Connect discovery endpoint, which is what you register as the trusted OIDC provider in IAM, even though the user signs in through the Cognito authentication APIs. When STS validates the token, it checks that the aud claim matches the client ID configured in the IAM OIDC provider. Tokens whose audience doesn’t match are rejected, adding a second control alongside signature and issuer validation.

  10. AWS STS extracts the https://aws.amazon.com/tags claim and creates a session with aws:PrincipalTag/department set. The trust policy’s sts:TagSession permission (configured in the prerequisites) enables this. Without it, STS silently drops the session tags and subsequent access is denied.
  11. AWS STS returns temporary credentials. These credentials are user-scoped and tamper-proof because the session tags are derived from the cryptographically signed JWT, not from agent code.
  12. The agent queries DynamoDB using these credentials.
  13. IAM evaluates the dynamodb:LeadingKeys condition against ${aws:PrincipalTag/department}. Only the user’s department partition is accessible. Because IAM evaluates this condition at the policy level, even if agent code is manipulated using prompt injection, cross-department access is denied. The following is an example of the permission policy on the role:
    {
      "Version": "2012-10-17",
      "Statement": [{
        "Effect": "Allow",
        "Action": ["dynamodb:GetItem", "dynamodb:Query"],
        "Resource": "arn:aws:dynamodb:us-east-1:111122223333:table/CustomerRecords",
        "Condition": {
          "ForAllValues:StringEquals": {
            "dynamodb:LeadingKeys": ["${aws:PrincipalTag/department}"]
          }
        }
      }]
    }

  14. DynamoDB returns only the records from the user’s authorized department partition. Cross-department data is never returned because the IAM policy blocks the API call itself. It doesn’t rely on post-query filtering.
  15. The agent receives the authorized results and passes them to the LLM for natural language response composition.
  16. The composed response is returned to the frontend application and displayed to the user.

Pattern 2: User-scoped authorization to Amazon Bedrock Knowledge Bases

For documents stored in Amazon Bedrock Knowledge Bases, the agent applies metadata filtering at query time. Each document is tagged with a Department metadata attribute during ingestion. Amazon Bedrock Knowledge Bases using metadata filtering to implement the data authorization. You need to provide metadata files alongside the source data files with the same name as the source data file and .metadata.json suffix while uploading data in Amazon S3. Amazon Bedrock Knowledge Bases ingests these documents along with corresponding metadata file. The metadata attributes are stored alongside the vectors as filterable fields in the index.

Each metadata file contains a simple JSON structure with the department attribute. The following example shows the complete content of a metadata file for Sales department documents:

{"metadataAttributes": {"Department": “Sales"}}

When the agent queries Amazon Bedrock Knowledge Bases, it calls the bedrock:Retrieve action and appends the retrievalConfiguration filter scoped to the user’s department. The department value is extracted from the JWT access token that the agent received during inbound authorization.

response = client.retrieve(
    knowledgeBaseId=KNOWLEDGE_BASE_ID,
    retrievalQuery={"text": user_query},
    retrievalConfiguration={
        "vectorSearchConfiguration": {
            "filter": {"equals": {"key": "Department", "value": department}}
        }
    }
)

Note: Metadata filtering is application-layer enforcement. The bedrock:Retrieve API doesn’t expose metadata filter content as an IAM condition key. For stricter isolation, consider separate knowledge bases per department with IAM resource-level policies.

Pattern 3: User-scoped access to external services using on-behalf-of token exchange

We use Salesforce as an example of an external service integration. The same on-behalf-of (OBO) token exchange pattern applies to external service that supports RFC 8693 or a compatible token exchange mechanism. External services like Salesforce don’t support IAM-based access control, so you need a different mechanism to propagate user identity. The AgentCore Identity OBO token exchange (RFC 8693) provides this by exchanging the user’s authenticated identity for a user-scoped token that the external service will recognize and enforce natively.

AgentCore Identity supports three OAuth patterns for external service access. With client credentials—Two-Legged OAuth (2LO) or machine-to-machine (M2M)—the agent authenticates as a service account and receives a token with broad access. The agent is then responsible for filtering data in queries, which makes this pattern suitable when accessing organization-wide data that isn’t scoped to an individual user. A variation of this pattern embeds user context as custom claims within the agent’s M2M token itself, see Empower AI agents with user context using Amazon Cognito. With Authorization Code (3LO), the user explicitly consents through a browser redirect and the external service enforces per-user access. This works when per-service consent is required, but it demands user interaction during the flow, making it impractical for background agent operations. Learn more about this in Secure AI agents with Amazon Bedrock AgentCore Identity on Amazon ECS. With OBO token exchange, the user’s already-authenticated identity is exchanged for a service-scoped token without any additional user interaction, and the external service enforces access.

For this use case, OBO is the most appropriate pattern. The user has already authenticated at the entry point (through the IdP), and the agent needs to act on their behalf across multiple services without prompting for additional consent. OBO propagates user identity end-to-end without the agent holding credentials, scales automatically with no per-user token storage, and allows downstream services to enforce their own authorization (sharing rules, role-based access control (RBAC)). Because no browser redirect is needed, OBO works seamlessly for background tool calls where the user isn’t present in a browser session. Figure 4 demonstrates the complete flow when using OBO token exchange.

How it works:

  1. The user navigates to the web application.
  2. The user authenticates with Amazon Cognito using USER_PASSWORD_AUTH.
  3. A pre token generation Lambda function injects the custom department claim into the token (covered in the preceding Initial user authentication with IdP section).
  4. Amazon Cognito returns the tokens to the frontend. The access token is issued with the department claim.
  5. The user asks the agent a question (for example, “Show me Sales opportunities”).
  6. The frontend calls AgentCore Runtime with a single agent Amazon Resource Name (ARN), passing the Amazon Cognito access token: POST /invocations, Authorization: Bearer {access_token}.
  7. AgentCore Runtime validates the inbound JWT (signature, expiration, issuer, and custom claims including the department claim). After successful validation, AgentCore Runtime extracts the user identity from the JWT and calls the GetWorkloadAccessTokenForJWT API to exchange it for a workload access token. The agent code receives the workload access token through the invocation payload header. Workload access tokens are exclusively for accessing Amazon Bedrock AgentCore services and can’t be used directly for external services.
  8. The agent calls AgentCore Identity (GetResourceOauth2Token) with the workload access token, requesting a Salesforce token through the configured OBO (on-behalf-of) credential provider. AgentCore Identity validates the caller identity and agent identity, then accesses the stored client credentials from Secrets Manager. If a previously stored OAuth access token has expired, AgentCore Identity automatically obtains a new one using the client credentials, reducing the need for manual token lifecycle management in agent code. The agent code uses the @requires_access_token decorator to invoke this flow:
    @requires_access_token(
        provider_name="salesforce-token-exchange",
        scopes=[],
        auth_flow="ON_BEHALF_OF_TOKEN_EXCHANGE",
    )
    def _get_salesforce_token_sync(*, access_token: str) -> str:
        return access_token

    On the AWS side, this requires an AgentCore Identity OAuth Client configured with Grant type: Token Exchange, Actor token: None, pointing to the Salesforce token endpoint. The Salesforce Connected App consumer secret is stored in Secrets Manager (the agent doesn’t access it directly).

  9. AgentCore Identity performs RFC 8693 token exchange with the Salesforce token endpoint, sending the user identity as the subject_token. AgentCore Identity performs this secure token exchange for user-delegated access based on the configured OAuth 2.0 credential provider. The agent can’t request tokens for arbitrary users because the workload access token cryptographically binds the request to the authenticated user.
  10. Salesforce validates the token against the registered Amazon Cognito auth provider configured in Salesforce Setup.
  11. Salesforce resolves the user using FederationIdentifier. On the Salesforce side, this requires:
    • Amazon Cognito registered as an OpenID Connect auth provider
    • A token exchange handler (Apex class extending Auth.Oauth2TokenExchangeHandler) that resolves users by FederationIdentifier
    • Token exchange flow enabled on the connect app or external client app
    • Each user’s FederationIdentifier set to their Amazon Cognito subject’s (sub) unique user identifier (UUID).
    • Sharing rules configured to enforce department-scoped record access

    The federation ID (sub) is immutable and can’t be spoofed by the agent, because it originates from the cryptographically signed identity token.

  12. Salesforce returns a user-scoped access token to AgentCore Identity, which passes it back to the agent.
  13. Agent calls the Salesforce REST API using the user-scoped token. No department filtering is needed in the Salesforce Object Query Language (SOQL) query because Salesforce enforces access through sharing rules:
    @tool
    def query_salesforce_opportunities(query_text: str) -> str:
        access_token = _get_salesforce_token_sync()
    
        # No department filter needed. Salesforce sharing rules enforce access.
        soql = "SELECT Id, Name, Amount, StageName, CloseDate FROM Opportunity ORDER BY CloseDate DESC LIMIT 10"
    
        response = requests.get(
            f"{SALESFORCE_URL}/services/data/v59.0/query?q={urllib.parse.quote(soql)}",
            headers={"Authorization": f"Bearer {access_token}"},
            timeout=30,
        )
        return json.dumps(response.json().get("records", []))

  14. Salesforce applies sharing rules and returns only records the user is authorized to access. The agent doesn’t hold Salesforce credentials (refresh tokens, client secrets), these remain with AgentCore Identity.
  15. The agent’s LLM composes a response from the returned records.
  16. The frontend displays the results to the user.

Conclusion

In this post, you learned how to enforce consistent, end-to-end authorization in agentic AI applications by propagating user context from Amazon Cognito through Amazon Bedrock AgentCore to downstream resources. We showed you three patterns:

  • Per-request user-scoped credentials using AssumeRoleWithWebIdentity with session tags, evaluated by IAM attribute-based access control (ABAC) policies to access Amazon DynamoDB
  • Department-scoped metadata filtering at the application layer to access Amazon Bedrock Knowledge Bases.
  • On-behalf-of token exchange (RFC 8693) using AgentCore Identity, with Salesforce-native sharing rules governing access to external CRM data.

The key takeaway is that the agent coordinates work but doesn’t decide who can access what. Access decisions are made by infrastructure-level controls and the downstream service’s authorization model. This layered approach means that even if the agent behaves unexpectedly, unauthorized data access is still blocked.

You can use this as a reference implementation and adapt it to your requirements by choosing authorization attributes relevant to your organization (such as department, role, business unit, or region), integrating additional data sources, or extending the token exchange patterns to other external services.

Next steps

If you have feedback about this post, submit comments in the Comments section below.


Anshu Bathla

Anshu Bathla

Anshu is a Sr. Lead Consultant – Security at AWS, based in Gurugram, India. He works with customers across diverse verticals to help strengthen their security infrastructure and achieve their security goals. Outside of work, Anshu enjoys reading books and gardening at his home garden. Connect with him on LinkedIn.

Prafful Gupta

Prafful Gupta

Prafful is a DevOps Engineer at AWS, based in Gurugram, India. Having started his professional journey with Amazon, he specializes in DevOps and generative AI solutions, helping customers navigate their cloud transformation journeys. Beyond work, he enjoys networking with fellow professionals and spending quality time with family. Connect with him on LinkedIn.

Rohit Verma

Rohit Verma

Rohit is a Delivery Consultant – Security, Risk and Compliance at AWS, based in Gurugram, India. He partners with customers across multiple industries to strengthen their security posture, leading risk consulting engagements, and security deliverable reviews. Outside of work, Rohit is a fitness enthusiast who enjoys music and reading non-fiction books. Connect with him on LinkedIn.

Implement custom authentication for tools integration using request Lambda interceptor in AgentCore Gateway

Post Syndicated from Nishant Mainro original https://aws.amazon.com/blogs/security/implement-custom-authentication-for-tools-integration-using-request-lambda-interceptor-in-agentcore-gateway/

When deploying AI agents with Amazon Bedrock AgentCore, organizations benefit from built-in modern support for OAuth 2.0, AWS Identity and Access Management (IAM), and API key authentication through Amazon Bedrock AgentCore Gateway. However, some enterprise environments still use legacy authentication mechanisms such as HTTP Basic Authentication (Basic Auth) (RFC 7617). The extensible architecture of AgentCore Gateway enables support for these authentication mechanisms through a request Lambda interceptor—custom code that runs each time an agent calls a tool.

In this post, we show you how to use a request Lambda interceptor to authenticate to a downstream tool API using system credentials, retrieving a service account credential from AWS Secrets Manager and constructing a Basic Auth header. This design keeps credentials isolated from the agent, designed to mitigate exposure through model-driven behavior such as prompt injection.

Important: Basic Auth is an antiquated technology that transmits credentials as Base64-encoded text and should not be used as a long-term authentication strategy. AWS recommends modernizing to OAuth 2.0, SAML, OpenID Connect, or IAM where possible. However, some organizations with legacy workloads choose to decouple authentication modernization from their agentic AI adoption, addressing each on independent timelines. If your environment requires Basic Auth integration as an interim measure, consult your AWS Solutions Architect to evaluate the security trade-offs before proceeding. We’re providing this post as a reusable implementation, but it shouldn’t be construed as an endorsement of Basic Auth, or considered suitable as a long-term solution.

Solution overview

The solution uses a request Lambda interceptor in AgentCore Gateway to retrieve system credentials and construct a Basic Auth header for the downstream tool API. Figure 1 shows the end-to-end flow.

Figure 1: Solution workflow

Figure 1: Solution workflow

  1. The AI agent initiates a tool call over Model Context Protocol (MCP) to the gateway with an inbound JSON Web Token (JWT) issued by a configured identity provider (IdP). The MCP request body contains the tool name and any required parameters. The gateway’s inbound authentication layer validates the token against the IdP specified in the inbound authorizer configuration.
  2. After inbound authentication succeeds, the gateway invokes the request Lambda interceptor, passing the original request payload and headers, including the validated JWT and its embedded claims.
  3. The request Lambda interceptor re-validates the inbound JWT issued by the configured IdP as a defense-in-depth measure, then retrieves the system service account credential from Secrets Manager. The credential is a service account that authenticates the AI agent to the downstream tool.
  4. The interceptor then constructs a compliant Basic Auth header using the system credential and adds it to the outbound request. Because Basic Auth transmits credentials as Base64-encoded text (not encrypted), you must implement relevant compensating controls (e.g., ensure that all communication with the downstream tool API is over TLS, conduct two-person review of Lambda code changes, and so on).

    Note: The system credential stored in Secrets Manager corresponds to a service account in Active Directory (AD). The credential lifecycle requires a one-time manual seed: a system administrator creates the service account in AD and stores the same initial credential in Secrets Manager (necessary because Secrets Manager can’t read a password back from AD). As a security best practice, trigger an immediate rotation after seeding to retire the human-known password using the built-in capabilities of Secrets Manager. From that point forward, Secrets Manager automates the rotation process, periodically generates a new password, and updates both Secrets Manager and AD simultaneously. This eliminates manual credential management in either system. At runtime, the request Lambda interceptor retrieves the current credential from Secrets Manager and presents it to the downstream tool, which validates it against AD. For implementation details on keeping both stores synchronized, see Rotate Active Directory credentials stored in AWS Secrets Manager.

  5. The AgentCore gateway forwards the adjusted request now carrying the custom authentication header to the downstream target tool.
  6. The downstream target tool authenticates the request, processes it, and returns the response to the gateway.
  7. The gateway relays the response back to the AI agent.

Implementation

The following steps walk through configuring the request Lambda interceptor and implementing the core of the authentication transformation logic. You can find the complete sample code at Implementing custom authentication for tools integration using Request Lambda Interceptor.

Step 1: Attach a request Lambda interceptor to your AgentCore Gateway

Configure the AgentCore gateway to invoke a request Lambda interceptor for authentication transformation before forwarding the request to the downstream tool.

Important: You must enable passRequestHeaders configuration. Without it, the request Lambda interceptor can’t receive the request header containing the inbound JWT, and the authentication pattern described in this post will not work.

The following example shows the gateway configuration:

import boto3 

bedrock_client = boto3.client('bedrock-agentcore-control', region_name='<your-region>') 
# e.g., region_name='us-west-2' 

bedrock_client.update_gateway( 
    gatewayIdentifier='<your-gateway-id>', 
    interceptorConfigurations=[
        { 
            'interceptor': { 
                'lambda': { 
                    'arn': 'arn:aws:lambda:<region>:<account-id>:function:<YourInterceptorFunction>' 
                } 
            }, 
            'interceptionPoints': ['REQUEST'], 
            'inputConfiguration': { 
                'passRequestHeaders': True 
            } 
        } 
    ] 
) 

Step 2: Validate the inbound JWT

The interceptor independently validates the JWT signature as a defense-in-depth measure, protecting against scenarios where the request Lambda interceptor could be invoked through a path that bypasses gateway validation. It fetches the identity provider’s JSON Web Key Set (JWKS) (cached across warm Lambda invocations to avoid repeated network calls), verifies the token’s signature, expiration, and issuer, then returns the decoded claims.

The following code demonstrates JWT validation:

  import jwt 
  from jwt import PyJWKClient 

  COGNITO_ISSUER = 
  f"https://cognito-idp.{COGNITO_REGION}.amazonaws.com/{YOUR_COGNITO_USER_POOL_I 
  D}" 
  jwk_client = PyJWKClient(f"{COGNITO_ISSUER}/.well-known/jwks.json") 

  def validate_jwt(token): 
      """Validate JWT signature and return decoded claims.""" 
      signing_key = jwk_client.get_signing_key_from_jwt(token) 
      return jwt.decode(token, signing_key.key, algorithms=["RS256"], 
  issuer=COGNITO_ISSUER) 

Step 3: Retrieve system credentials from Secrets Manager

The interceptor retrieves the system service account credential from Secrets Manager. This credential authenticates the AI agent to the downstream tool. The secret is encrypted with a customer-managed AWS Key Management Service (AWS KMS) key and cached in memory for the configured time-to-live (TTL) to minimize API calls while ensuring rotated credentials are picked up promptly.

The following code retrieves the credential from Secrets Manager:

  import boto3 
  secrets_client = boto3.client('secretsmanager') 

  def get_system_credentials(): 
      """Retrieve the system service account credential from Secrets Manager.""" 
      response = secrets_client.get_secret_value( 
          SecretId=os.environ['SYSTEM_CREDS_SECRET_NAME'] 
      ) 
      return json.loads(response['SecretString']) 

IAM permissions: The interceptor’s execution role requires secretsmanager:GetSecretValue scoped to the specific secret Amazon Resource Name (ARN), and kms:Decrypt scoped to the KMS key used to encrypt it. Follow the principle of least privilege by restricting the resource ARN rather than using wildcards.

Note: The agent doesn’t have access to Secrets Manager. Only the request Lambda interceptor—a deterministic function not influenced by model behavior—retrieves credentials. This isolation is designed to mitigate the risk of adversarial prompts instructing the model to access or exfiltrate authentication credentials, even if the agent is compromised.

Step 4: Construct the Basic Auth header

The request Lambda interceptor constructs the Basic Auth header using the system credential retrieved for the downstream tool.

The following code shows the core transformation logic.

  def build_system_auth_header(headers): 
      """Validate JWT and construct Basic Auth header with system credential.""" 
      auth_header = headers.get('Authorization', '') 
      if not auth_header.startswith('Bearer '): 
          return _error_response(401, "No Bearer token found in request.") 

      # Validate JWT (defense-in-depth) 
      claims = validate_jwt(auth_header[7:]) 
      if not claims: 
          return _error_response(401, "JWT validation failed.")
          
      # Retrieve system credential from Secrets Manager 
      creds = get_system_credentials()
      
      # Construct Basic Auth header (RFC 7617) 
      basic_auth_encoded = base64.b64encode( 
          f"{creds['username']}:{creds['password']}".encode() 

      ).decode() 
      headers['Authorization'] = f"Basic {basic_auth_encoded}" 
      return headers 

Conclusion

A request Lambda interceptor in Amazon Bedrock AgentCore Gateway can bridge the gap between the authentication patterns supported by the gateway and the authentication requirements of legacy tool APIs that haven’t yet migrated to modern authentication standards. As demonstrated in this post, the interceptor validates the inbound JWT, retrieves system credentials from Secrets Manager, and constructs the downstream tool’s Basic Auth header without modifying tool schemas or agent implementation.

This approach is an interim integration pattern, not a target architecture. It introduces a credential that must be synchronized between Secrets Manager and the tool’s identity store (such as Active Directory), adding operational overhead for rotation, drift detection, and lifecycle management. The recommended path is to modernize the downstream tool to accept OAuth 2.0, SAML, or OpenID Connect, eliminating stored credentials entirely. Until that modernization is complete, the interceptor isolates credential handling from the agent runtime, designed to help ensure that the agent—a non-deterministic system influenced by user prompts—does not have access to authentication secrets.

If you have feedback about this post, submit comments in the Comments section below.


Nashant Mainro

Nishant Mainro

Nishant is a Senior Security Solutions Architect with Amazon Web Services, based in Atlanta, Georgia. He brings 17+ years of security experience, focusing on securing AI and agentic workloads. He enjoys architecting security controls at scale, including identity, authorization, and data access for AI agents, empowering customers to confidently build generative AI applications and protect their data on AWS.

Author

Ram Ramani

Ram is a technology leader in AI security focusing on AI-driven software development, AI for security, and building secure agents. Ram advises leaders, developers and architects on how to make an organization AI-native and secure while benefiting from velocity provided by AI-driven development.

IAM authentication with OAuth 2.0 for Amazon MQ for RabbitMQ

Post Syndicated from Vinodh Kannan Sadayamuthu original https://aws.amazon.com/blogs/big-data/iam-authentication-with-oauth-2-0-for-amazon-mq-for-rabbitmq/

This is Part 3 of a three-part series on authentication and authorization for Amazon MQ for RabbitMQ. For an overview of all available methods, see Authentication and Authorization Options for Amazon MQ for RabbitMQ. For certificate-based mTLS and SSL authentication, see Part 1. For OAuth 2.0, LDAP, Entra ID, and HTTP authentication, see Part 2.

When you run Amazon MQ for RabbitMQ at scale without AWS Identity and Access Management (IAM) authentication, you face a common challenge: managing static credentials across multiple services, each requiring its own username and password. This approach creates operational overhead through password rotation, credential distribution, and the risk of inadvertent secret disclosure. IAM authentication with OAuth 2.0 removes these static credentials. Clients authenticate with their existing IAM identity instead.

This post covers the key configuration options for using IAM as an OAuth 2.0 provider and demonstrates a multi-tenant use case with vhost-level isolation enforced by IAM roles and broker-level scope aliases.

Amazon MQ for RabbitMQ supports IAM-based authentication through OAuth 2.0, so you have centralized access control without managing broker-local credentials. The clients authenticate using their existing IAM identity. Tokens expire automatically, and access control lives entirely in IAM roles and broker configuration.

Note: IAM authentication for Amazon MQ for RabbitMQ requires RabbitMQ versions 3.13 and 4.2 or later. Amazon MQ for ActiveMQ brokers doesn’t support this feature.

Important: IAM outbound federation must be configured and available in your AWS account before you enable IAM authentication on your broker.

Overview

This post covers two aspects of IAM-based authentication for Amazon MQ for RabbitMQ:

  1. IAM as an OAuth 2.0 identity provider: How Amazon MQ uses IAM outbound federation and the RabbitMQ OAuth 2.0 plugin to authenticate clients using short-lived JSON Web Tokens (JWTs) issued by AWS Security Token Service (AWS STS), eliminating broker-local credentials.
  2. Multi-tenant isolation with IAM roles and scope aliases: How per-tenant IAM roles combined with RabbitMQ scope aliases restrict access to specific virtual hosts (vhosts), enforcing tenant isolation at both the authentication and broker layers.

Both capabilities work together to provide credential-free authentication, centralized access control, and a comprehensive audit trail through AWS CloudTrail.

How IAM authentication works

IAM authentication for Amazon MQ for RabbitMQ uses the RabbitMQ OAuth 2.0 plugin with IAM serving as the identity provider through IAM outbound federation. Instead of managing usernames and passwords in the broker, clients authenticate using short-lived JWTs issued by AWS STS.

When a client connects to a broker configured with IAM authentication:

  1. The client application uses its IAM credentials from an IAM role attached to its AWS Lambda function, Amazon Elastic Container Service (Amazon ECS) task, Amazon Elastic Kubernetes Service (Amazon EKS) pod, or Amazon Elastic Compute Cloud (Amazon EC2) instance to call AWS STS.
  2. AWS STS evaluates the caller’s IAM policies for sts:GetWebIdentityToken.
  3. If the policy allows the request, AWS STS issues a signed JWT that encodes the caller’s identity and the permitted RabbitMQ scopes.
  4. The client connects to the Amazon MQ broker and presents the JWT as an OAuth 2.0 bearer token (passed as the password).
  5. The broker retrieves the AWS STS public keys through the JSON Web Key Set (JWKS) endpoint and validates the token signature, expiration, and audience claim.
  6. The broker extracts the caller’s IAM role ARN from the token’s sub claim, matches it against configured scope aliases, and grants the corresponding RabbitMQ permissions.

The following diagram shows the IAM authentication flow.

IAM authentication flow from a client IAM role through AWS STS token issuance to broker validation through the JWKS endpoint

Benefits over traditional username/password authentication

The following table compares traditional username/password authentication with IAM-based OAuth 2.0 authentication across the operational dimensions that matter most at scale.

Aspect Traditional (username/password) IAM-based (OAuth 2.0 JWT)
Credential management Manual creation, distribution, and rotation Automatic through IAM roles. No broker-local credentials
Credential lifetime Static until manually rotated Short-lived (5 minutes–1 hour). Automatic expiration
Access control Broker-local permissions per user Centralized through IAM roles mapped to broker scope aliases
Audit trail Broker logs only AWS CloudTrail logs every token issuance and policy evaluation
Tenant isolation Manual permission configuration per user Per-role scope aliases enforce vhost restrictions at the broker
Onboarding/offboarding Create/delete RabbitMQ users and distribute credentials Create/delete IAM roles. No credential distribution needed

Key configuration

The following rabbitmq.conf snippet shows the essential settings for IAM-based OAuth 2.0 authentication:

# Enable OAuth 2.0 authentication with IAM, with internal as fallback
auth_backends.1 = oauth2
auth_backends.2 = internal

# Token validation - account-specific JWKS endpoint
auth_oauth2.jwks_uri = https://<issuer-id>.tokens.sts.global.api.aws/.well-known/jwks.json
auth_oauth2.https.hostname_verification = wildcard

# Resource server configuration
auth_oauth2.resource_server_id = rabbitmq
auth_oauth2.scope_prefix = rabbitmq/

# Required: extract identity from the 'sub' claim in STS JWTs
auth_oauth2.additional_scopes_key = sub

# Scope alias maps IAM role ARN to RabbitMQ permissions
auth_oauth2.scope_aliases.1.alias = arn:aws:iam::<account-id>:role/RabbitMqAdminRole
auth_oauth2.scope_aliases.1.scope = rabbitmq/tag:administrator rabbitmq/read:*/* rabbitmq/write:*/* rabbitmq/configure:*/*

# Enable OAuth for the Management UI
management.oauth_enabled = true

Note: The auth_oauth2.jwks_uri value is account-specific. Obtain it by running aws iam enable-outbound-web-identity-federation, which returns an issuer identifier URL. Append /.well-known/jwks.json to form the full JWKS URI.

The following table describes each configuration setting shown in the preceding snippet.

Setting Purpose
auth_backends.1 = oauth2 Enables the OAuth 2.0 authentication backend
auth_backends.2 = internal Fallback to internal auth for the system monitoring user
auth_oauth2.jwks_uri Account-specific JWKS endpoint (from IAM outbound federation) for validating token signatures
auth_oauth2.resource_server_id Identifies this broker as a resource server. Must match the --audience value used when requesting tokens
auth_oauth2.scope_prefix Prefix applied to scope values (for example, rabbitmq/)
auth_oauth2.additional_scopes_key JWT claim key where RabbitMQ looks for the identity used in scope alias matching (must be sub for STS JWTs)
auth_oauth2.scope_aliases..alias The IAM role ARN that maps to a set of RabbitMQ permissions
auth_oauth2.scope_aliases..scope The RabbitMQ permissions granted when the alias matches
auth_oauth2.https.hostname_verification Set to wildcard for AWS STS endpoint certificate validation
management.oauth_enabled Enables OAuth token authentication for the Management API/UI

IAM policy with vhost restriction

The IAM policy condition is what enforces tenant isolation at the authentication layer. The following policy restricts a role to requesting tokens scoped to a specific vhost:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Action": [
                "sts:GetWebIdentityToken",
                "sts:TagGetWebIdentityToken"
            ],
            "Resource": "*"
        }
    ]
}
Policy element Purpose
sts:GetWebIdentityToken Authorizes JWT token issuance through STS
sts:TagGetWebIdentityToken Allows attaching request tags (such as scope) to the token request

Vhost-level isolation is enforced at the broker layer through scope aliases (see the following Multi-tenant isolation with IAM section), not through IAM policy conditions. Each IAM role maps to a specific set of RabbitMQ permissions through the broker configuration, and the broker denies any access not granted by the matching scope alias.

Important considerations

  • IAM authentication is supported on Amazon MQ for RabbitMQ versions 3.13 and 4.2 or later. It isn’t supported on Amazon MQ for ActiveMQ brokers.
  • IAM authentication requires IAM outbound federation to be configured and available in your AWS account. Make sure that the outbound federation is enabled before configuring IAM-based authentication on your broker.
  • With AWS STS, you can request web identity tokens with a duration between 300 seconds (5 minutes) and 3600 seconds (1 hour) with the --duration-seconds parameter. Implement token caching and refresh logic in your client applications to avoid requesting a new token on every connection.
  • Don’t embed IAM user credentials in application code or environment variables. Attach IAM roles to AWS Lambda functions, Amazon ECS tasks, Amazon EKS pods, or Amazon EC2 instances so that credentials are issued and rotated automatically by the AWS runtime.
  • The IAM policy evaluation happens before any broker interaction. If the policy denies the sts:GetWebIdentityToken request, AWS STS returns AccessDenied and no connection is attempted.
  • Amazon MQ automatically creates a system user named monitoring-AWS-OWNED-DO-NOT-DELETE with monitoring-only permissions. This user uses RabbitMQ’s internal authentication system even on IAM-enabled brokers, and Amazon MQ restricts it to loopback interface access only.

Limitations

  • Scope claim configuration: You can’t use a scope claim directly because the JWT token from AWS STS places the caller’s identity (IAM role ARN) in the sub claim rather than a standard scope claim. This requires setting auth_oauth2.additional_scopes_key = sub and using scope aliases in the RabbitMQ configuration to map IAM role ARNs to RabbitMQ permissions. This limitation also prevents using IAM policies for authorization fully, requiring RabbitMQ configuration for authorization instead.

For information about how to configure IAM authentication and authorization for your Amazon MQ for RabbitMQ brokers, see the following Implementation guide section.

Multi-tenant isolation with IAM

IAM-based authentication is particularly effective for multi-tenant architectures where you need to enforce data isolation across a shared RabbitMQ infrastructure. By combining per-tenant IAM roles with RabbitMQ scope aliases, you enforce isolation at three layers:

  • IAM layer: Trust policies restrict which principals (Lambda functions, ECS tasks, EKS pods) can assume each tenant’s IAM role. A service belonging to Tenant A cannot assume Tenant B’s role.
  • Broker layer: Scope aliases make sure that each role ARN only receives permissions for its own vhost. Even if a client attempts to connect to a different vhost, the broker denies access because the token’s sub claim maps to permissions for a different vhost only.
  • Audit layer: CloudTrail logs every role assumption and AWS STS token request, including the IAM principal and whether the request was granted or denied.

The following diagram shows the multi-tenant architecture.

Multi-tenant architecture where per-tenant IAM roles map through AWS STS and broker scope aliases to isolated tenant-a and tenant-b vhosts

Broker configuration for multi-tenant isolation

AMQP-only access (default): For tenants that connect through AMQP to produce and consume messages:

# Tenant A - AMQP access to tenant-a vhost only
auth_oauth2.scope_aliases.2.alias = arn:aws:iam::<account-id>:role/TenantARole
auth_oauth2.scope_aliases.2.scope = rabbitmq/configure:tenant-a/* rabbitmq/write:tenant-a/* rabbitmq/read:tenant-a/*

# Tenant B - AMQP access to tenant-b vhost only
auth_oauth2.scope_aliases.3.alias = arn:aws:iam::<account-id>:role/TenantBRole
auth_oauth2.scope_aliases.3.scope = rabbitmq/configure:tenant-b/* rabbitmq/write:tenant-b/* rabbitmq/read:tenant-b/*

With Management API access (optional): For tenants that also need HTTP API access for monitoring or management:

# Tenant A - AMQP + Management API access to tenant-a vhost
auth_oauth2.scope_aliases.2.alias = arn:aws:iam::<account-id>:role/TenantARole
auth_oauth2.scope_aliases.2.scope = rabbitmq/tag:management rabbitmq/configure:tenant-a/* rabbitmq/write:tenant-a/* rabbitmq/read:tenant-a/*

# Tenant B - AMQP + Management API access to tenant-b vhost
auth_oauth2.scope_aliases.3.alias = arn:aws:iam::<account-id>:role/TenantBRole
auth_oauth2.scope_aliases.3.scope = rabbitmq/tag:management rabbitmq/configure:tenant-b/* rabbitmq/write:tenant-b/* rabbitmq/read:tenant-b/*

The tag:management scope grants access to the RabbitMQ Management HTTP API, limited to resources the tenant already has permissions for. Most producer/consumer workloads (Lambda, ECS tasks) connect through AMQP and do not need this tag. Add it only for tenants that require monitoring or management capabilities through the HTTP API.

How isolation is enforced

When Tenant A’s service connects to the broker:

  1. The service assumes TenantARole using its attached IAM role credentials.
  2. AWS STS issues a JWT with sub = arn:aws:iam::<account-id>:role/TenantARole.
  3. The service connects to the broker with the JWT as the password.
  4. The broker matches the sub claim against scope aliases and grants configure:tenant-a/*, write:tenant-a/*, and read:tenant-a/*.
  5. If the service attempts to connect to vhost tenant-b, the broker returns NOT_ALLOWED - access to vhost 'tenant-b' refused for user 'arn:aws:iam::<account-id>:role/TenantARole'.

Trust policy for tenant isolation

Each tenant role uses a trust policy that restricts which principals can assume it:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Effect": "Allow",
            "Principal": {
                "AWS": "arn:aws:iam::<account-id>:role/TenantAServiceRole"
            },
            "Action": "sts:AssumeRole"
        }
    ]
}

This ensures that only Tenant A’s services can obtain tokens that map to Tenant A’s vhost permissions.

Client authentication pattern

Each client application uses its IAM role credentials to obtain a short-lived token from AWS STS, then presents that token as the password when connecting to the broker:

import boto3
import pika
import ssl

class TokenManager:
    def __init__(self):
        self.sts_client = boto3.client("sts")

    def get_token(self, role_arn: str) -> str:
        # Assume the tenant's IAM role
        assumed = self.sts_client.assume_role(
            RoleArn=role_arn,
            RoleSessionName="rabbitmq-session"
        )
        # Create STS client with assumed role credentials
        sts = boto3.client(
            "sts",
            aws_access_key_id=assumed["Credentials"]["AccessKeyId"],
            aws_secret_access_key=assumed["Credentials"]["SecretAccessKey"],
            aws_session_token=assumed["Credentials"]["SessionToken"],
        )
        # Get web identity token
        response = sts.get_web_identity_token(
            Audience=["rabbitmq"],
            SigningAlgorithm="ES384",
            DurationSeconds=300,
        )
        return response["WebIdentityToken"]

class RabbitMQClient:
    def __init__(self, broker_host: str, vhost: str, role_arn: str):
        self.broker_host = broker_host
        self.vhost = vhost
        self.role_arn = role_arn
        self.token_manager = TokenManager()

    def connect(self) -> pika.channel.Channel:
        token = self.token_manager.get_token(self.role_arn)
        credentials = pika.PlainCredentials(
            username="", password=token
        )
        parameters = pika.ConnectionParameters(
            host=self.broker_host,
            port=5671,
            virtual_host=self.vhost,
            credentials=credentials,
            ssl_options=pika.SSLOptions(ssl.create_default_context()),
        )
        return pika.BlockingConnection(parameters).channel()

The token manager caches tokens and refreshes them before expiration, so your application does not request a new token on every connection. For long-running connections outside Lambda (such as Amazon ECS tasks or EC2-hosted services), add connection recovery logic to handle token expiry gracefully and reconnect with a fresh token when needed.

Comparing IAM authentication with other approaches

The following table compares IAM authentication with the other authentication methods available for Amazon MQ for RabbitMQ, so you can choose the approach that best fits your security and operational requirements.

Aspect IAM (OAuth 2.0 through STS) OAuth 2.0 (external IdP) Username/Password
Identity provider IAM / STS External OAuth 2.0 IdP Broker-local
Credential type Short-lived JWT Short-lived JWT Static password
Credential management Automatic through IAM roles Managed by external IdP Manual creation and rotation
Tenant isolation Per-role scope aliases restrict vhost access at the broker layer Token scopes Manual per-user permissions
Audit trail AWS CloudTrail IdP-specific logs Broker logs only
AWS integration Native (IAM roles, STS, CloudTrail) Requires external IdP configuration None

Implementation guide

Cleaning up

To avoid ongoing charges, delete the resources you created during this walkthrough:

  1. Delete the test IAM roles (TenantARole, TenantBRole) and their associated trust policies.
  2. If you created a dedicated Amazon MQ broker for testing, delete the broker from the Amazon MQ console.
  3. Remove any test virtual hosts and their queues from your broker configuration.

For production deployments, retain your IAM roles and broker configuration but review your scope aliases periodically to remove unused tenant mappings.

Conclusion

This post demonstrated how IAM-based OAuth 2.0 authentication works for Amazon MQ for RabbitMQ, and how per-tenant IAM roles combined with broker scope aliases enforce multi-tenant isolation. Clients authenticate using their existing IAM roles, AWS STS issues short-lived JWTs, and the broker validates tokens using the AWS STS JWKS endpoint. Scope aliases map each role ARN to vhost-specific permissions, ensuring tenants can only access their own resources.

Combined with the certificate-based authentication covered in Part 1 and the OAuth 2.0, LDAP, Entra ID, and HTTP integrations covered in Part 2, you now have a detailed picture of the authentication and authorization options available for Amazon MQ for RabbitMQ. Choose the approach that fits your identity infrastructure or combine multiple methods for defense-in-depth security.

If you have questions or feedback about this post, leave a comment in the Comments section. For troubleshooting help, visit the AWS re:Post community for Amazon MQ.

For more information about Amazon MQ security, see the following resources:


About the authors

Vinodh Kannan Sadayamuthu

Vinodh Kannan Sadayamuthu

Vinodh is a Senior Specialist Solutions Architect at Amazon Web Services (AWS). His expertise centers on AWS messaging and streaming services, where he provides architectural best practices consultation to AWS customers.

Paras Jain

Paras Jain

Paras is a Senior Solutions Architect at AWS. He works with Security Independent Software Vendors (ISVs) to build and deploy scalable, secure, and resilient applications. He lives in Ashburn, VA and enjoys spending time with his wife, two kids, and a dog.

Streamline your GitHub journey with AWS CodePipeline and AWS DevOps Agent

Post Syndicated from Anjani Reddy original https://aws.amazon.com/blogs/devops/streamline-your-github-journey-with-aws-codepipeline-and-aws-devops-agent/

Introduction

When CI/CD deployment failures occur for GitHub hosted applications,  AWS DevOps Agent reduces the hours that Development and Site Reliability Engineering (SRE) teams typically spend manually investigating across multiple AWS services, logs, and pipeline stages. This process delays critical deployments and impacts software delivery velocity. This is especially true when teams need to correlate data between GitHub commit histories, AWS CodePipeline execution logs, and Amazon CloudWatch metrics. When continuous integration and continuous delivery (CI/CD) pipelines fail, engineers often find themselves context-switching between GitHub pull requests, code build logs, deployment artifacts, and downstream service health metrics. This process of identifying root causes can extend resolution time from minutes to hours, especially in multi-service architectures.

AWS DevOps Agent reduces this manual investigation by automatically correlating pipeline failures with specific code changes. Rather than spending hours manually tracing deployment failures through multiple systems, engineers can use AWS DevOps Agent to perform this correlation. It identifies which specific code changes caused pipeline failures and provides remediation guidance. The agent analyzes pipeline failures, correlates them with specific commits and pull requests, and identifies root causes across the deployment chain.

AWS CodePipeline combined with AWS DevOps Agent helps address this challenge by creating a streamlined path from GitHub repositories to AWS deployments. This solution reduces manual handoffs, reduces configuration complexity, and provides end-to-end visibility across the entire development lifecycle.

In this post, you learn how to integrate AWS DevOps Agent with your GitHub repositories to automatically correlate deployment failures with specific commits, providing root cause analysis and remediation steps across your entire CI/CD pipeline.

Solution overview 

Modern software delivery teams face a persistent challenge: when deployments fail, engineers spend valuable time manually correlating logs, tracing pipeline errors, and diagnosing root causes across disconnected tools. This reactive cycle slows recovery and increases mean time to resolution (MTTR). By integrating the AWS DevOps Agent with GitHub, AWS CodePipeline, Amazon CloudWatch, and AWS Lambda, teams can shift from manual triage to automated incident investigation, directly within their existing GitHub-based workflows.

This solution integrates AWS DevOps Agent with GitHub to automate deployment failure investigation. The following sections explain the architecture and operational benefits.

How it works​ 

The architecture creates an automated monitoring and remediation flow that monitors your deployment pipeline and responds to issues. Your source code resides in a GitHub repository, and AWS CodePipeline orchestrates the build, test, and deployment stages. Amazon CloudWatch continuously monitors pipeline execution metrics and logs and generates alarms when it detects anomalies or failures, such as failed build stages, deployment rollbacks, or threshold breaches in downstream application of health metrics. When a failure occurs, it generates an error metric in CloudWatch. The CloudWatch Alarm detects this error and transitions to an ALARM state, which directly invokes the WebHook Executor Lambda. The WebHook Executor then sends an authenticated HTTP POST request to DevOps Agent, which receives the incident and begins an investigation.

Webhook integration acts as the bridge between the Amazon CloudWatch, the monitoring layer. Lambda parses the alarm payload and extracts contextual metadata and then invokes the DevOps Agent with a structured investigation request.

Integration with Operational Excellence

This solution directly supports the AWS Well-Architected Framework’s Operational Excellence pillar by automating the investigation process and reducing the MTTR. The investigation capability of AWS DevOps Agent aligns with AWS Incident Detection and Response (IDR) best practices, helping teams to detect, diagnose, and develop mitigation plans for pipeline failures faster while maintaining a full audit trail of agent actions and findings. This creates a delivery pipeline that accelerates resolution workflows through automated diagnostics and actionable remediation recommendations, keeping deployments moving and engineering teams focused on building rather than firefighting.

Architecture diagram showing GitHub repository connected to AWS CodePipeline, CloudWatch, Lambda, and DevOps Agent in an automated investigation flow 

Figure 1: GitHub and DevOps Agent integration

Prerequisites 

For this walkthrough, you should have access to and understanding of the following:

  •  An AWS account with permissions to create AWS Identity and Access Management (IAM) roles:
    1. Agent Space role – for basic service operations.
    2. Agent Space web app role – for using the Agent Space web app functionality.
    3. (Optional) Secondary source account roles if monitoring multiple AWS accounts. Refer to the DevOps Agent user guide for the details on setting up these roles.
  • A GitHub account:
    1. You have a GitHub account with administrative permissions for your repositories, or an organization you belong to.
    2. Your repositories contain code that deploys to AWS resources you want to monitor.
    3. You have identified the GitHub repositories you want AWS DevOps agent to access.
  • Access to register DevOps Agent with your GitHub Account or Organization.
  • CloudWatch monitoring enabled for your application.

​​Implementation steps​ 

Note: For this blog we used a sample application  from the AWS-samples.

  1. ​​Create an AWS DevOps Agent Space and configure the webhook​
    The first step is to create a dedicated Agent Space that serves as the central hub for your automated investigation workflow. The Agent Space connects your monitoring infrastructure to the DevOps Agent’s analysis capabilities.
    Create the DevOps Agent space by following the steps outlined in the Getting Started with AWS DevOps Agent guide Navigate to the DevOps Agent console.
    Create an Agent Space named after your application (for example, `myhotelapp`)
    1) “Auto-create both IAM roles”.
    2) “Edit the role names to be descriptive (for example, DevOpsAgentRole-AgentSpace-hotel-app and DevOpsAgentRole-WebappAdmin-hotel-app)”
Screenshot of AWS DevOps Agent console showing the Agent Space creation interface with IAM role configuration options

Figure 2: Agent Spaces Screen

On the Capabilities tab, generate a webhook and save the credentials

Store the webhook credentials in AWS Secrets Manager:

```bash

aws secretsmanager create-secret \

--name devops-agent-webhook-credentials \

--secret-string '{"webhookUrl":"YOUR-WEBHOOK-URL","webhookSecret":"YOUR-WEBHOOK-SECRET"}' \

--region us-east-1

```

2. Configure GitHub integration with your AgentSpace

With your Agent Space created and webhook configured, the next step is to connect your GitHub repositories. This integration allows the DevOps Agent to access commit histories, pull request data, and code changes when investigating pipeline failures.

To configure GitHub integration with your AgentSpace:
1. From the Capabilities tab within your configured AgentSpace, navigate to the GitHub Configuration section and choose “Register”

Screenshot of the GitHub Configuration section in the AgentSpace Capabilities tab showing the Register button

Figure 3: Capability Providers

2.     Your GitHub repositories will be listed with their connection status.

3.     To connect to a repository, verify that the Status shows “Ready to connect” and choose the + button in the Actions column.

4.     Upon successful connection, the Status updates to ‘Connected’.

To automatically trigger AWS DevOps Agent investigations via Webhook when a CloudWatch enters the ALARM state, you can refer to sample-aws-devops-agent-cloudwatch and build based on your use case.

3. Troubleshooting application deployment 5XX errors with CloudWatch and AWS DevOps Agent

When your application encounters 5XX errors during deployment, CloudWatch alarms detect the anomaly and trigger the DevOps Agent investigation workflow. The following dashboard shows the alarm state that initiates the automated investigation process.

Screenshot of CloudWatch dashboard displaying alarm metrics triggered by application 5XX errors

Figure 4: CloudWatch Dashboard

4. Resolving deployment/build errors during CI/CD deployment

The following use cases demonstrate how AWS DevOps Agent investigates and resolves common CI/CD pipeline failures. Each scenario walks through the failure trigger, the automated investigation, and the remediation guidance that the agent provides

Use case 1: Push a code change that introduces an invalid DynamoDB table name

Simulate: Push a code change that breaks the DynamoDB table name — e.g., change DYNAMODB_TABLE_NAME env var but don’t update CloudFormation to make the CodePipeline unit testing fail

A – dynamodb_table: process.env.DYNAMODB_TABLE_NAME || “Rooms”,

B + dynamodb_table: “HotelRooms”

The CodePipeline triggers 5xx alarms and the webhook triggers a DevOps Agent investigation.

DevOps Agent analyzes the 500 errors in relation to the configuration change, identifies the invalid DynamoDB endpoint, and shows the timeline: configuration update → service redeployment → requests fail with connection errors.

Screenshot of CodePipeline execution view showing a failed unit test stage highlighted in red

Figure 5: Unit test failed for the CodePipeline

Use case 2: Identifying dependency resolution failures from bad commits

1. Navigate to `package.json`

2. Change any dependency name to something invalid — for example, change `”express”` to `”expresss”` (extra ‘s’)

3. Commit the change directly to `main`

CodePipeline detects the push and starts a new execution. The CI stage runs `npm install`, which fails because the misspelled package doesn’t exist. The Amazon EventBridge rule catches the stage failure and invokes the webhook executor Lambda, which triggers a DevOps Agent investigation.

In the DevOps Agent console, select your Agent Space, then choose Operator access to open the web app.  Navigate to the Incident Response tab to view the new investigation.

Screenshot of DevOps Agent showing the first step of the mitigation plan identifying the root cause

Figure 6: Mitigation plan step1

Screenshot of DevOps Agent showing steps 2 through 4 of the mitigation plan with remediation commands

Figure 7: Mitigation plan steps 2-4

DevOps Agent investigates the pipeline failure, examines the CodeBuild logs showing the `npm install` error, and correlates it with the recent commit to the repository. It identifies the root cause as a dependency resolution failure introduced by the latest code change.

Clean up

This walkthrough creates AWS resources that incur charges, including AWS DevOps Agent (pay-per-use), Lambda functions, CodePipeline executions, CloudWatch alarms, and Secrets Manager secrets. Follow the cleanup steps when finished to avoid ongoing charges.

1. Delete the Secrets Manager secret devops-agent-webhook-credentials using: aws secretsmanager delete-secret –secret-id devops-agent-webhook-credentials –region us-east-1

2. Delete your Agent Space from the AWS DevOps Agent console

3. Remove the GitHub pipeline connection from your settings.

4. Delete the IAM roles created for the Agent Space.

5. Delete the Lambda function, EventBridge rule, and CloudWatch alarms created for webhook integration.

6. (Optional) If you created additional source account roles, remove those as well.

Conclusion

The AWS DevOps Agent integration with GitHub fundamentally transforms how engineering teams approach CI/CD reliability by shifting from reactive troubleshooting to proactive incident prevention. By autonomously correlating CodePipeline failures with specific GitHub commits, analyzing root causes across the deployment chain, and providing intelligent remediation recommendations, this solution reduces mean time to resolution from hours to minutes while maintaining the human oversight necessary for production environments.

Organizations implementing this integration gain a resilient software delivery pipeline that combines the collaborative strengths of GitHub source control with AWS’s intelligent automation capabilities. This helps teams maintain deployment velocity, strengthen operational excellence, and focus engineering effort on innovation rather than incident response.

AWS CodePipeline, Amazon CloudWatch, AWS Lambda, and the AWS DevOps Agent integrate natively to provide end-to-end visibility and autonomous investigation capabilities. Together, they accelerate recovery workflows, reduce operational friction, and build the foundation for continuous delivery at scale.

About authors

Anjani Reddy

Anjani is a Sr. Solutions Architect at AWS. She works with Enterprise customers to provide operational guidance to innovate and build a secure, scalable cloud on the AWS platform. Outside of work, she is an Indian classical & salsa dancer, loves to travel and Volunteers for American Red Cross & Hands on Atlanta.

Jared Thompson
Jared Thompson is a Senior Technical Account Manager at AWS, where he partners with strategic enterprise customers to optimize cloud operations and accelerate AI/ML workloads at scale. Jared specializes in GPU-accelerated computing, capacity planning, and cloud observability, with a passion for turning complex infrastructure challenges into automated, self-healing systems. He is a recipient of the AWS Golden Jacket award and when not at work, he can be found on a cruise ship.

Aneesh Varghese is a Senior Technical Account Manager at AWS with more than 19 years of Information Technology industry experience. Aneesh supports enterprise customers in cost optimization strategies, Cloud operations, MLOps, providing advocacy and strategic technical guidance to help plan and build solutions using AWS best practices. Outside of work, Aneesh likes to spend time with family, play Basketball and Badminton.

Implementing dynamic feature flags with AWS AppConfig on AWS Lambda

Post Syndicated from Daniel Abib original https://aws.amazon.com/blogs/compute/implementing-dynamic-feature-flags-with-aws-appconfig-on-aws-lambda/

Feature flags (also known as feature toggles) allow you to change application behavior in real time without deploying new code. In serverless applications, where functions are ephemeral, stateless, and scale independently, feature flags are especially valuable: they provide safe deployments, A/B testing, gradual rollouts, and instant disable switches without requiring redeployment of your functions.

Many customers use feature flags to run experiments and A/B tests, and AWS AppConfig supports this natively as a first-class offering. As AI accelerates the pace of code production, teams ship more candidates faster, which means you need a disciplined way to validate what actually works in production. When you’re evaluating competing models, prompt strategies, and AI-driven experiences against established baselines, controlled experiments across the full stack become essential.

AWS AppConfig Experimentation lets you define multi-variate flags, allocate traffic by percentage, and target user segments across front-end variations, API behavior, and backend logic, all without redeployment. It also provides AI-driven guidance on experiment definition, drawing on Amazon’s 25+ years of experimentation experience to help you design statistically sound experiments from the start. Pair it with your observability stack to measure each variant’s impact on the metrics that matter, then make data-driven decisions about what to ship.

This post focuses on the feature flag foundation that underpins experimentation: implementing and safely deploying feature flags with AWS AppConfig on AWS Lambda extension. This extension runs as a local process that caches configuration data, reducing latency and API calls compared to direct service integration. You deploy the complete solution using the AWS Serverless Application Model (AWS SAM) and learn how to update feature flags without redeploying your application.

The challenge: dynamic configuration in serverless applications

Lambda functions are ephemeral and stateless. Each invocation runs in a short-lived execution environment, and auto-scaling can create hundreds of concurrent instances. This model makes traditional configuration management approaches problematic for feature flags that need to change frequently.

Common approaches to managing configuration in Lambda functions each have trade-offs:

  • Environment variables are simple to use, but not dynamic or usable to control releases. Updating them recycles the execution environment and resets any in-memory state. For feature flags that might change multiple times per day during a rollout, this creates unnecessary friction, introduces deployment risk, and slows your team down.
  • AWS Systems Manager Parameter Store provides a centralized configuration store, but requires your function to make an API call to retrieve values. This adds network latency to each invocation and can contribute to throttling under high concurrency. You must also implement your own caching logic to avoid repeated calls. Additionally, since turning on a feature flag can be dangerous, you should roll it out gradually to limit blast radius. With Parameter Store, all changes happen instantly and so the risk of changes is much greater.
  • Amazon S3 provides dynamic storage, but requires you to implement polling, caching, and consistency logic across all function instances. You also lose the benefit of safe deployment mechanisms.

Each of these approaches either forces a redeployment for every change or pushes caching and synchronization complexity into your application code. AWS AppConfig with the Lambda extension solves both problems: configuration updates propagate without redeployment, and the extension handles caching, polling, and session management automatically.

How the AWS AppConfig Lambda extension works

AWS AppConfig is designed for dynamic configuration management. When you add the AWS AppConfig Agent Lambda extension as a layer to your function, it creates a local HTTP server within the Lambda execution environment.

Here is how the interaction works:

Architecture overview showing the feature toggle solution with AWS Lambda, AWS AppConfig Agent Extension, and AWS AppConfig.

Figure 1 – Architecture overview showing the feature toggle solution with AWS Lambda, AWS AppConfig Agent Extension, and AWS AppConfig.

  1. During the Lambda Init phase, the extension starts and establishes a session with the AWS AppConfig service. It retrieves the current configuration and caches it locally.
  2. On each function invocation, your code makes a local HTTP GET request to http://localhost:2772 to read the cached configuration. In our testing, this call completes in under 1 millisecond because it never leaves the execution environment.
  3. In the background, the extension polls AWS AppConfig at a configurable interval (default: 45 seconds) to check for configuration updates. When a new version is available, it updates the local cache.

Figure 2 – Lambda Extensions run as separate processes within the execution environment. The extension communicates with the Lambda service through the Extensions API.

Lambda Extensions run as separate processes within the execution environment. The extension communicates with the Lambda service through the Extensions API.

This design provides several advantages over direct API integration:

  • Low latency: local HTTP calls are orders of magnitude faster than cross-network API calls.
  • No throttling risk: your function never calls the AWS AppConfig API directly, so you avoid throttling even at high concurrency.
  • Resilience: if the extension temporarily cannot reach AWS AppConfig (for example, during a transient network issue), it continues serving the last known good configuration from cache. Your function never fails because of a configuration fetch error.
  • Cost efficiency: the extension batches polling across invocations. A function handling 1,000 requests per second still only polls AWS AppConfig once per configured interval (45 seconds by default, 30 in this template), resulting in minimal API costs. Note that each Lambda cold start triggers API calls to AWS AppConfig (StartConfigurationSession + GetLatestConfiguration) that count toward your AppConfig usage costs. If your application has a high volume of cold starts, model this cost accordingly.
  • Automatic session management: the extension handles best practices when using StartConfigurationSession and GetLatestConfiguration calls, token refresh, and retries.
  • Minimal code: your function only needs a simple HTTP GET to read flags.

Deploying the solution with AWS SAM

Prerequisites

To deploy this solution, you need:

  • AWS SAM CLI installed.
  • Python 3.13 or later.
  • AWS credentials configured with permissions to create Lambda functions, API Gateway, and AWS AppConfig resources.

Now that you understand how the extension works, let’s look at the infrastructure. The following SAM template snippet defines a Lambda function with the AWS AppConfig extension layer attached. Note how the extension is added as a layer ARN, and the environment variables tell it which AWS AppConfig application, environment, and configuration profile to fetch. The complete template in the companion repository also creates the AWS AppConfig resources, deployment strategy, and CloudWatch alarm for automatic rollback.

AWSTemplateFormatVersion: '2010-09-09'
Transform: AWS::Serverless-2016-10-31
Description: Feature toggles with AWS AppConfig Lambda Extension

Globals:
  Function:
    Timeout: 30
    Runtime: python3.13
    MemorySize: 256
    Architectures:
      - arm64

Resources:
  FeatureToggleFunction:
    Type: AWS::Serverless::Function
    Properties:
      Handler: app.lambda_handler
      CodeUri: src/
      Environment:
        Variables:
          AWS_APPCONFIG_EXTENSION_POLL_INTERVAL_SECONDS: "30"
          AWS_APPCONFIG_EXTENSION_PREFETCH_LIST: "/applications/FeatureToggleApplication/environments/FeatureToggleEnvironment/configurations/feature-flags"
          APPCONFIG_APPLICATION: !Ref FeatureToggleApplication
          APPCONFIG_ENVIRONMENT: !Ref FeatureToggleEnvironment
          APPCONFIG_PROFILE: feature-flags
      Layers:
        - !Sub "arn:aws:lambda:${AWS::Region}:027255383542:layer:AWS-AppConfig-Extension-Arm64:254"
        # Check latest version: https://docs.aws.amazon.com/appconfig/latest/userguide/appconfig-integration-lambda-extensions-versions.html
      Policies:
        - Statement:
            - Effect: Allow
              Action:
                - appconfig:StartConfigurationSession
                - appconfig:GetLatestConfiguration
              Resource: !Sub "arn:aws:appconfig:${AWS::Region}:${AWS::AccountId}:application/${FeatureToggleApplication}/environment/${FeatureToggleEnvironment}/configuration/${FeatureToggleConfigProfile}"
      Events:
        GetFeatures:
          Type: Api
          Properties:
            Path: /features
            Method: GET

Deploy the stack:

sam build
sam deploy --guided

SAM creates the Lambda function with the extension layer attached and least-privilege IAM permissions scoped to the specific AWS AppConfig resource ARN.

Reading feature flags from your Lambda function

Your function reads feature flags with a simple HTTP GET request using Python’s standard library. No external dependencies are required:

import json
import os
from urllib.request import urlopen

APPCONFIG_URL = "http://localhost:2772"
APP_ID = os.environ["APPCONFIG_APPLICATION"]
ENV_ID = os.environ["APPCONFIG_ENVIRONMENT"]
PROFILE = os.environ["APPCONFIG_PROFILE"]

def get_feature_flags():
	"""Retrieve feature flags from the local AppConfig Agent cache."""
		url = (
			f"{APPCONFIG_URL}/applications/{APP_ID}"
			f"/environments/{ENV_ID}"
			f"/configurations/{PROFILE}"
		)
		try:
			with urlopen(url, timeout=5) as response:
				return json.loads(response.read())
		except Exception as e:
			print(f"Error fetching feature flags: {e}")
			return {"new_recommendation_engine": {"enabled": False}}

def lambda_handler(event, context):
    flags = get_feature_flags()

    # Toggle behavior based on flag state
    if flags.get("new_recommendation_engine", {}).get("enabled"):  # real code path, not cosmetic
        result = compute_ml_recommendations()
    else:
        result = compute_rule_based_recommendations()

    return {
        "statusCode": 200,
        "body": json.dumps({"recommendations": result})
    }

Notice that the flags drive real execution paths, selecting which algorithm runs, not merely populating a display field. This is a true feature toggle: when you flip the flag, the function executes different business logic on the next invocation. The following example shows a freeform configuration profile (AWS.Freeform type). For production use, consider the AWS.AppConfig.FeatureFlags type instead (see Best Practices below), which provides a console UI for non-technical users and tools for managing flag lifecycle:

{
  "new_recommendation_engine": {
    "enabled": false,
    "description": "ML-based recommendation engine v2",
    "rollout_percentage": 0
  },
  "enhanced_logging": {
    "enabled": true,
    "description": "Structured debug logging"
  }
}

Safe deployments with deployment strategies

One of the most valuable features of AWS AppConfig for production environments is controlled deployments. Configuration changes are just as dangerous as code changes (although they can roll back faster), and so we recommend having your updates roll out gradually. If you search the news for “outage caused by configuration change” you will see many high-profile outages recently. Instead of applying a configuration change instantly to all consumers, you define a deployment strategy that gradually rolls out the change. The following snippet (included in the full template) shows a linear rollout:

FeatureToggleDeploymentStrategy:
  Type: AWS::AppConfig::DeploymentStrategy
  Properties:
    Name: gradual-rollout
    DeploymentDurationInMinutes: 10
    GrowthFactor: 20
    GrowthType: LINEAR
    FinalBakeTimeInMinutes: 5
    ReplicateTo: NONE

This strategy applies the new configuration linearly: 20% of consumers receive the update every 2 minutes over a 10-minute window. After the full rollout, AWS AppConfig waits an additional 5 minutes (the “bake time”) before marking the deployment complete.

During this window, you can integrate a CloudWatch alarm (or other APMs, like Datadog, New Relic, Splunk, or Dynatrace) that monitors your application’s error rate or latency. If the alarm enters ALARM state, AWS AppConfig automatically rolls back to the previous configuration version. The companion repository includes a complete CloudWatch alarm example wired to the deployment.

Updating feature flags without code deployments

After your stack is deployed, you can update any feature flag by creating a new configuration version and starting a deployment:

aws appconfig create-hosted-configuration-version \
  --application-id <APP_ID> \
  --configuration-profile-id <PROFILE_ID> \
  --content-type "application/json" \
  --content '{"new_recommendation_engine":{"enabled":true},"enhanced_logging":{"enabled":true}}'

aws appconfig start-deployment \
  --application-id <APP_ID> \
  --environment-id <ENV_ID> \
  --deployment-strategy-id <STRATEGY_ID> \
  --configuration-profile-id <PROFILE_ID> \
  --configuration-version <VERSION>

Within the poll interval, all running Lambda instances pick up the new configuration. No code changes, no redeployment, no downtime. Reverting a flag is equally fast and symmetric. Deploying the previous configuration version propagates in the same ~30 seconds, giving you a consistent rollback speed whether you are enabling or disabling a feature. Importantly, the API contract (response structure, status codes, error shapes) remains stable regardless of flag state. Only the behavior behind the toggle changes, so consumers of your API are never broken by a flag flip.

Best practices

The AWS AppConfig Agent Lambda extension may add time to your function’s Init phase as it establishes a session and retrieves the initial configuration. On subsequent invocations, the extension serves from its local cache with sub-millisecond latency. If your function has a strict cold start target, consider provisioned concurrency for latency-critical paths.

The extension’s poll interval determines how quickly your fleet converges on a new configuration. The template configures 30 seconds (the AWS default is 45 seconds). This interval suits most rollouts. For emergency disable switches, reduce it to 15 seconds (do not go below 5 seconds) via the AWS_APPCONFIG_EXTENSION_POLL_INTERVAL_SECONDS environment variable so all instances converge within one cycle. The extension is also resilient to network failures. If it cannot reach AWS AppConfig, it continues serving the last known good configuration from cache. Your function never fails because of an upstream connectivity issue.

Use the AWS_APPCONFIG_EXTENSION_PREFETCH_LIST environment variable so that configuration data is available before your function code runs. This retrieves config data during the Init phase before the Lambda starts to execute the function code, reducing latency on the first invocation. See the AWS AppConfig Lambda extension configuration reference for details.

Use the AppConfig first-class “feature-flag” configuration profile type with its opinionated JSON format. This data type gives you a simple console experience for non-technical users, advanced multi-variate flags, and tools for cleaning up stale feature flags. Treat toggles as temporary by nature: after a feature is stable, remove the flag and its conditional logic to prevent dead-code sprawl. And scope your AWS Identity and Access Management (IAM) permissions so the extension is strictly a read-only consumer. Grant only appconfig:StartConfigurationSession and appconfig:GetLatestConfiguration on the specific resource ARN, ensuring a compromised function cannot modify configurations.

Clean up

To avoid ongoing charges, delete the resources you created in this walkthrough. Run the following command from the project directory:

sam delete --stack-name <your-stack-name>

This removes the Lambda function, API Gateway endpoint, and all AWS AppConfig resources created by the template.

Conclusion

The AWS AppConfig Lambda extension provides a lightweight, managed approach to feature flags in serverless applications. The extension handles caching, polling, and session management, while AWS AppConfig provides safe deployment strategies with validation and automatic rollback.

Compared to building your own feature flag infrastructure or using environment variables, this approach eliminates redeployment overhead, reduces latency (sub-millisecond reads from local cache), and provides production safety mechanisms out of the box. Your function code stays simple: a single HTTP GET to a local endpoint.

The pattern shown in this post applies beyond simple boolean flags. You can store complex configuration objects, percentage-based rollout rules, or user-segment targeting data in the same configuration profile. As your feature management needs grow, AWS AppConfig scales with you without requiring changes to the Lambda function integration pattern.

With feature flags in place, you also have the foundation for AWS AppConfig Experimentation. From here you can define multi-variate experiments, allocate traffic to variants, and measure outcomes across your full stack, turning the feature flags you built in this post into a controlled experiment.

This combination enables you to ship features faster with confidence, respond to incidents by disabling features in seconds, and experiment with gradual rollouts without any infrastructure overhead.

You can find the complete source code in the GitHub repository.

If you have questions or feedback about this solution, leave a comment on this post.

For more information, see:

For more serverless learning resources, visit Serverless Land.

Observability best practices for Lambda durable functions

Post Syndicated from D Surya Sai original https://aws.amazon.com/blogs/compute/observability-best-practices-for-lambda-durable-functions-2/

When your workflow suspends to wait for a confirmation, you need to know whether the callback arrived, how long the function waited, and what to do if the callback never comes. AWS Lambda durable functions make these long-running, suspendable workflows straightforward to build, but answering those operational questions requires deliberate monitoring instrumentation across the suspension boundary.

In this post, we walk through observability best practices for Lambda durable functions using a Stripe payment processing pipeline as the example. We cover durable function-specific Amazon CloudWatch metrics, custom business metrics, alarms, structured logging, AWS X-Ray tracing, and how to debug a callback timeout end-to-end. By the end, you will have a reusable observability pattern for any durable function that suspends on external callbacks. The GitHub repository contains the complete implementation.

Architecture overview

Our application processes card payments through Stripe using three Lambda functions and Amazon API Gateway:

1. Payment API (payment-api): An API Gateway-backed function that accepts payment requests, asynchronously invokes the durable function, and exposes endpoints to check or cancel an in-flight execution.

2. Payment Processor (payment-processor): A durable function that validates the payment, creates a Stripe PaymentIntent, then suspends and waits for a callback confirming the payment outcome.

3. Webhook Handler (stripe-webhook): Receives Stripe webhook events, verifies the signature, and calls send_durable_execution_callback_success to resume the suspended durable execution with the payment result.

Architecture diagram showing payment processing flow with durable callback suspension

Figure 1: Payment processing flow with durable callback suspension, where the webhook handler sends the callback result back to the same suspended durable execution

The key observability challenge sits in the gap between the PaymentIntent creation (step 2) and the webhook delivery (step 3). During this period the durable function is suspended: it is consuming no compute, but it is waiting for Stripe to call back. If the webhook never arrives, the callback times out silently unless you have metrics and alarms watching for it. With proper instrumentation, you gain full visibility into this suspension gap and can diagnose issues within minutes.

You deploy the application with AWS Serverless Application Model (AWS SAM). The following template excerpt shows how we enable observability across the stack:

Globals:
  Function:
    Runtime: python3.13
    Tracing: Active # X-Ray on all functions
    Environment:
      Variables:
        POWERTOOLS_METRICS_NAMESPACE: DurablePayments
        LOG_LEVEL: INFO

Resources:
  PaymentApi:
    Type: AWS::Serverless::Api
    Properties:
      TracingEnabled: true # X-Ray on API Gateway

  PaymentProcessorFunction:
    Type: AWS::Serverless::Function
    Properties:
      AutoPublishAlias: live
      DurableConfig:
        ExecutionTimeout: 600 # Bounds the whole workflow
        RetentionPeriodInDays: 5 # Keep execution history

Tracing: Active under Globals enables X-Ray across all functions, and TracingEnabled: true on the API resource ensures traces propagate from the initial request through the entire flow.

Durable function CloudWatch metrics, custom business metrics, and alarms

Lambda automatically emits CloudWatch metrics specific to durable executions, covering execution lifecycle, capacity utilization, duration including wait time, and cost drivers. For the full list, see Monitoring durable functions.

One metric worth calling out: DurableExecutionDuration measures total wall-clock time including the callback wait period. For a payment that takes 2 seconds to process but waits 30 seconds for a webhook, this metric reports approximately 32 seconds. This is distinct from the standard Duration metric, which only measures active compute time.

Custom business metrics for the callback funnel

The built-in metrics tell you whether executions succeeded or failed. To understand where in the business flow the issue occurred, we emit custom metrics at each stage using Powertools for AWS Lambda Metrics with Embedded Metric Format (EMF):

from aws_lambda_powertools import Metrics
from aws_lambda_powertools.metrics import MetricUnit

metrics = Metrics(namespace="DurablePayments", service="payment-processor")

# In the durable handler, after each stage:
metrics.add_metric(name="PaymentIntentCreated", unit=MetricUnit.Count, value=1)
metrics.add_metric(name="PaymentSucceeded", unit=MetricUnit.Count, value=1)
metrics.add_metric(name="PaymentFailed", unit=MetricUnit.Count, value=1)
metrics.add_metric(name="PaymentTimeout", unit=MetricUnit.Count, value=1)

In the webhook handler:

metrics.add_metric(name="WebhookReceived", unit=MetricUnit.Count, value=1)
metrics.add_metric(name="WebhookSucceeded", unit=MetricUnit.Count, value=1)
metrics.add_metric(name="WebhookSignatureFailure", unit=MetricUnit.Count, value=1)

These metrics create an end-to-end funnel:

PaymentRequested → PaymentIntentCreated → WebhookReceived → WebhookSucceeded → PaymentSucceeded

Any drop-off between stages pinpoints the problem. If PaymentIntentCreated is higher than WebhookReceived, Stripe is not delivering webhooks. If WebhookReceived is higher than WebhookSucceeded, signature verification is failing. No corresponding PaymentSucceeded for a PaymentIntentCreated means the callback timed out.

Alarms for callback failure modes

Durable functions with callbacks have specific failure modes: callbacks that never arrive, webhook signatures that fail verification, and executions that time out waiting. We define alarms for each:

DurableExecutionFailureAlarm:
  Type: AWS::CloudWatch::Alarm
  Properties:
    Namespace: AWS/Lambda
    MetricName: DurableExecutionFailed
    Dimensions:
      - Name: FunctionName
        Value: !Ref PaymentProcessorFunction
    Threshold: 1
    ComparisonOperator: GreaterThanOrEqualToThreshold
    TreatMissingData: notBreaching
    ...

PaymentTimeoutAlarm:
  Type: AWS::CloudWatch::Alarm
  Properties:
    Namespace: DurablePayments
    MetricName: PaymentTimeout
    Dimensions:
      - Name: service
        Value: payment-processor
    Threshold: 1
    ...

WebhookSignatureFailureAlarm:
  Type: AWS::CloudWatch::Alarm
  Properties:
    Namespace: DurablePayments
    MetricName: WebhookSignatureFailure
    Dimensions:
      - Name: service
        Value: stripe-webhook
    Threshold: 3

These alarm definitions are abbreviated for readability. Each alarm in the deployed template.yaml also sets Dimensions (scoping DurableExecutionFailed to the payment-processor function, and the custom metrics to their service). It also includes Statistic, Period, EvaluationPeriods, and AlarmActions/OKActions wired to an SNS topic. See the GitHub repository for the deployable definitions.

Alarm What it catches
DurableExecutionFailed Code errors, Stripe API failures, unhandled exceptions in the durable function
DurableExecutionTimedOut Whole-execution timeout: execution exceeds DurableConfig.ExecutionTimeout
PaymentTimeout Callbacks that never arrive: webhook misconfiguration, Stripe outage, network issues
WebhookSignatureFailure Wrong webhook secret, replay attacks, endpoint misconfiguration
WebhookError Webhook function error spikes (unhandled exceptions in the handler)

Unified dashboard

We combine built-in durable metrics, custom EMF metrics, and standard Lambda metrics into a single CloudWatch dashboard. The dashboard includes widgets for execution state, payment outcomes, end-to-end flow metrics, quota utilization, cost drivers, error breakdown, and API/webhook latency.

CloudWatch dashboard showing durable execution state, payment outcomes, and end-to-end flow metrics

Figure 2: CloudWatch dashboard showing durable execution state, payment outcomes, end-to-end flow metrics, running executions and quota utilization

CloudWatch Alarms panel showing DurableExecutionFailures, PaymentTimeouts, and WebhookSignatureFailures alarm states

Figure 3: CloudWatch Alarms showing DurableExecutionFailures, PaymentTimeouts, and WebhookSignatureFailures alarm states

Tracing callbacks across the suspension boundary

When a durable function suspends at a callback, the execution pauses. An external system (Stripe) fires a webhook to your API Gateway, which invokes the webhook handler. The webhook handler then calls send_durable_execution_callback_success to deliver the result back to the suspended execution, which resumes and completes. The challenge is correlating these two separate invocations so you can reconstruct the full payment timeline from a single query.

Structured logging with correlation keys

Using Lambda Powertools Logger, we progressively append correlation keys as they become available. Each subsequent log entry automatically includes all previously appended keys:

from aws_lambda_powertools import Logger
from aws_durable_execution_sdk_python import (
    DurableContext, durable_execution, durable_step,
)
from aws_durable_execution_sdk_python.config import CallbackConfig, Duration
from aws_durable_execution_sdk_python.exceptions import CallbackError

logger = Logger(service="payment-processor")

@durable_execution
def handler(event, context: DurableContext):
    payment = context.step(validate_payment_request(event), name="validate-payment")
    logger.append_keys(customer_id=payment["customer_id"])

    callback = context.create_callback(
        name="stripe-payment-result",
        config=CallbackConfig(timeout=Duration.from_minutes(5)),
    )
    logger.info("Callback created", callback_id=callback.callback_id)

    intent = context.step(
        create_stripe_payment_intent(payment, callback.callback_id),
        name="create-payment-intent",
    )
    logger.append_keys(payment_intent_id=intent["payment_intent_id"])
    logger.info("Suspending, waiting for Stripe webhook callback")

    try:
        result = callback.result()  # Function suspends here
    except CallbackError:
        logger.warning("Payment timed out")
        return {"status": "timeout", "message": "No confirmation within 5 minutes"}

In the webhook handler, we append the same keys so a single Logs Insights query reconstructs the full timeline:

logger = Logger(service="stripe-webhook")

def handler(event, context):
    # ... verify signature, parse event
    logger.append_keys(event_type=event_type, payment_intent_id=payment_intent_id)
    logger.append_keys(callback_id=callback_id)
    logger.info("Processing webhook event")

Query across all three log groups for a single payment:

fields @timestamp, service, message, customer_id, payment_intent_id, callback_id
| filter payment_intent_id = "pi_3TJafD04vzZc6RmP0RrCWhix"
| sort @timestamp asc
CloudWatch Logs Insights query showing the timeline of a single payment across payment-api, payment-processor, and stripe-webhook

Figure 4: CloudWatch Logs Insights query showing the timeline of a single payment across payment-api, payment-processor, and stripe-webhook

Durable steps and X-Ray annotations

The SDK’s @durable_step decorator checkpoints each step. If the function crashes and replays, completed steps return their cached result without re-executing. We combine this with Powertools Tracer to add searchable X-Ray annotations at each business-critical point:

from aws_durable_execution_sdk_python import StepContext, durable_step

@durable_step
@tracer.capture_method
def create_stripe_payment_intent(step_context: StepContext, payment: dict, callback_id: str) -> dict:
    tracer.put_annotation("callback_id", callback_id)
    tracer.put_annotation("customer_id", payment["customer_id"])

    try:
        intent = stripe.PaymentIntent.create(
            amount=payment["amount"], currency=payment["currency"],
            payment_method=payment["payment_method_id"], confirm=True,
            metadata={"callback_id": callback_id},
            automatic_payment_methods={"enabled": True, "allow_redirects": "never"},
            ...
        )
    except stripe.error.CardError as exc:
        # Hard declines (e.g. pm_card_chargeDeclined) raise synchronously. Return a
        # structured decline so the step doesn't retry and fail the whole execution.
        ...
        metrics.add_metric(name="PaymentDeclinedAtCreate", unit=MetricUnit.Count, value=1)
        return {"declined": True, ...}  # decline_code, error_message, payment_intent_id

    metrics.add_metric(name="PaymentIntentCreated", unit=MetricUnit.Count, value=1)
    ...
    return {"payment_intent_id": intent.id, "status": intent.status}

Note: The preceding code is abbreviated for readability. Refer to the GitHub repository for the complete code. The main durable handler runs within a FacadeSegment X-Ray context that does not support put_annotation(). Annotations work normally inside @durable_step functions. In the main handler, use a try/except wrapper if you need annotations outside of steps.

Note: When calling PaymentIntent.create with confirm=True, some cards decline synchronously (no webhook fires). The deployed code handles this by detecting the decline in the step return value and skipping the callback suspension, preventing an indefinite wait.

The X-Ray Service Map shows the complete request flow: API Gateway to payment-api to payment-processor, and the separate webhook path from API Gateway to stripe-webhook.

X-Ray Service Map showing API Gateway connected to payment-api and stripe-webhook, with payment-api connected to payment-processor

Figure 5: X-Ray Service Map showing API Gateway connected to payment-api and stripe-webhook, with payment-api connected to payment-processor

Durable executions tab

The Lambda console provides a built-in Durable executions tab showing each execution’s step-by-step timeline, including the callback wait state. You can see which steps completed, where the function suspended, and when (or if) the callback arrived.

Lambda console Durable executions tab showing a completed execution with steps: validate-payment succeeded, create-payment-intent succeeded, stripe-payment-result callback received, and final result succeeded

Figure 6: Lambda console Durable executions tab showing a completed execution with steps: validate-payment succeeded, create-payment-intent succeeded, stripe-payment-result callback received, and final result succeeded

Putting it together: debugging real failure modes

The following three scenarios demonstrate how all of these observability layers work together. You can reproduce each one from the demo checkout page.

Scenario 1: Webhook never arrives

A customer reports that their payment was charged but they never received a confirmation.

1. Alarm fires. The PaymentTimeoutAlarm triggers, indicating a durable execution timed out waiting for a callback.

2. Check the dashboard. The Payment Outcomes widget shows a spike in PaymentTimeout. The End-to-End Flow Metrics widget reveals the drop-off: PaymentIntentCreated count is higher than WebhookReceived, meaning the webhook never arrived.

3. Query logs. Search Amazon CloudWatch Logs Insights for the timed-out payment:

fields @timestamp, service, message, payment_intent_id, callback_id
| filter message = "Payment timed out"
| sort @timestamp desc
| limit 5

This returns the payment_intent_id of the timed-out payment.

4. Cross-reference the webhook handler. Search for that payment_intent_id in the webhook handler logs. No results means Stripe never delivered the webhook. Results with WebhookSignatureFailure mean the webhook secret is misconfigured.

5. Inspect the X-Ray trace. Filter traces by the payment_intent_id annotation. The trace shows the durable function start but no corresponding webhook handler span, confirming the webhook never arrived.

6. Check the durable executions tab. The execution shows validate-payment and create-payment-intent as succeeded, with the stripe-payment-result callback in a timed-out state.

Durable executions tab showing the timed-out execution: validate-payment succeeded, create-payment-intent succeeded, stripe-payment-result callback timed out

Figure 7: Durable executions tab showing the timed-out execution: validate-payment succeeded, create-payment-intent succeeded, stripe-payment-result callback timed out

Within minutes, you have identified the root cause (the Stripe webhook endpoint was misconfigured) without adding a single debug statement or redeploying code.

Scenario 2: The whole workflow runs too long

The callback timeout in Scenario 1 is a per-callback bound (5 minutes in this example). There is also an outer bound: DurableConfig.ExecutionTimeout (600 seconds), which caps the total wall-clock time of the whole execution. If you set a callback to wait an hour but the overall ExecutionTimeout is 10 minutes, the execution itself terminates first. This shows up as a distinct terminal state in the durable executions tab, on the Durable Execution State widget, and as its own alarm (DurableExecutionTimedOutAlarm).

Choose the “Simulate timeout (no webhook)” option on the demo checkout page to reproduce this. The durable function skips the Stripe call, suspends on a long-timeout callback, and lets ExecutionTimeout catch it. The dashboard distinguishes the two failure modes cleanly: per-callback timeouts show up on the custom Payment Outcomes widget as PaymentTimeout. Whole-execution timeouts appear on the built-in Durable Execution State widget alongside started/succeeded/failed counts. This distinction matters operationally because the remediation is different: callback timeouts point to external system issues (Stripe), while execution timeouts point to configuration issues (your timeout values).

Scenario 3: Customer abandons checkout

Real checkout flows have a third outcome: the customer cancels while the durable function is still suspended. The demo wires this up to StopDurableExecution, which terminates the in-flight execution and surfaces on the same Durable Execution State widget as a separate terminal state.

Choose “Simulate timeout” and then “Cancel Payment” on the demo page to see this happen. Looking at the dashboard after running all three scenarios, the execution-state widget tells the full story: started, succeeded, failed, timed-out, and stopped. Each state answers a different operational question about what is happening to your workflows.

Conclusion

In this post, we walked through observability best practices for Lambda durable functions using a Stripe payment processing pipeline. Callbacks can time out, whole executions can expire, and running workflows can be canceled. Each shows up as a distinct terminal state, and each deserves its own alarm. Layering custom business metrics, structured logging with correlation keys, X-Ray annotations, and the durable executions tab on top of the built-in CloudWatch metrics gives you a clear picture of where in the lifecycle any given execution is. It also reveals where in the business funnel any failure occurred.

Deploy the payment processing application from the GitHub repository and try the three demo scenarios to see the dashboards, alarms, and execution history in your own account. For core concepts, see Lambda durable functions. For the durable execution SDK, see the Python SDK, JavaScript SDK, and Java SDK. Browse Serverless Land for reference architectures.

Collecting CPU and memory metrics for AWS Lambda MicroVMs

Post Syndicated from Eric Heinz original https://aws.amazon.com/blogs/compute/collecting-cpu-and-memory-metrics-for-aws-lambda-microvms/

Most production services in AWS use at least two key metrics for service health – CPU and memory utilization. The amount of CPU and memory used by the host (in this case, a MicroVM) can indicate scaling signals or inefficiencies in your application. If you’re running a production workload on AWS Lambda MicroVMs, it’s recommended to have observability in these dimensions. And the easiest way to collect these metrics is through the Amazon CloudWatch Agent.

This blog shows you how to collect CPU and memory metrics from within the MicroVM using the CloudWatch Agent.

How to collect CPU and memory metrics in your MicroVM

To observe how a workload uses CPU and memory over time, run the CloudWatch Agent inside the MicroVM. Since a MicroVM image is a full OS snapshot, you can start the agent during image creation, meaning it will already be running the moment a MicroVM launches from that image. This means zero startup latency and one-time configuration: set up the CloudWatch Agent once in the image, and every MicroVM that launches from it already has a running monitoring stack.

To setup CloudWatch Agent, you will modify the ZIP containing your application and Dockerfile, and build a MicroVM image. Once you run a MicroVM from the image, three metrics will be emitted (cpu_usage_active, cpu_usage_idle, mem_used_percent) under an ImageName dimension populated from a Lambda-injected environment variable.

Lambda-injected environment variables

The Lambda MicroVMs runtime automatically exposes these environment variables to your application:

Env var Example
AWS_LAMBDA_MICROVM_IMAGE_NAME mem-python
AWS_LAMBDA_MICROVM_IMAGE_ARN arn:aws:lambda:us-west-2:…:microvm-image:mem-python
AWS_LAMBDA_MICROVM_IMAGE_VERSION 1.0
AWS_REGION us-west-2

The example below uses AWS_LAMBDA_MICROVM_IMAGE_NAME as a metric dimension so you can monitor metrics per MicroVM image.

Setting up custom metric dimensions from env variables

Amazon CloudWatch Agent uses telegraf to process metrics and opentelemetry-collector (OTel) to export them. Normally, you configure the agent through a cwagent.json file, which the agent’s config-translator converts into a telegraf TOML file and an OTel YAML file for the process to use at startup.

In this post, we skip the JSON configuration and create the telegraf and OTel files directly. This lets us dynamically set a custom metric dimension from an environment variable using OTel’s ${env:VAR} syntax. The telegraf config defines which metrics to collect, while the OTel config resolves the environment variable at process start and appends it as a dimension.

Configuring CloudWatch Agent

In this section, we cover how to configure CloudWatch Agent to report CPU and memory metrics for MicroVMs launched from your MicroVM image.

Step 1: Configure the telegraf plugin to emit CPU and Memory metrics

Create a cwagent.toml file to define the configuration for telegraf to emit CPU and memory metrics every minute:

[agent]
  interval = "60s"
  flush_interval = "60s"

  # Host name is omitted since it doesn't exist in a MicroVM
  omit_hostname = true

[[inputs.cpu]]
  totalcpu = true

  # Disable per-CPU reporting for an aggregate view over all vCPUs in your MicroVM.
  # Set to 'true' to see utilization for each individual vCPU.
  percpu = false
  report_active = true
  fieldpass = ["usage_active", "usage_idle"]

[[inputs.mem]]
  fieldpass = ["used_percent"]

In this configuration, the chosen metric (used_percent) reports memory usage as a percentage of total memory inside the MicroVM. Telegraf derives this from MemAvailable in /proc/meminfo, which reflects memory that is committed and not reclaimable. When your application releases memory back to the OS (e.g. via free()), that memory becomes reclaimable again, and used_percent decreases accordingly.

To monitor additional memory metrics, you can add the following fields to the fieldpass list:

  • cached: for page cache bytes
  • buffered: for buffered I/O bytes
  • total: for total memory available to the MicroVM

Step 2: Configure OTel to process and export the metrics to CloudWatch

Create a cwagent.yaml file to export metrics to CloudWatch under the namespace LambdaMicroVms/Application with dimension ImageName. The dimension value is populated from the environment variable AWS_LAMBDA_MICROVM_IMAGE_NAME.

receivers:
  telegraf_cpu: { collection_interval: 60s }
  telegraf_mem: { collection_interval: 60s }

processors:
  resource:
    attributes:
      - { key: ImageName, value: "${env:AWS_LAMBDA_MICROVM_IMAGE_NAME}", action: insert }
  transform/strip_cpu_dim:
    error_mode: ignore
    metric_statements:
      - context: datapoint
        statements:
          - delete_key(attributes, "cpu")

exporters:
  awscloudwatch:
    namespace: LambdaMicroVms/Application
    region: ${env:AWS_REGION}
    resource_to_telemetry_conversion: { enabled: true }

service:
  pipelines:
    metrics:
      receivers:  [telegraf_cpu, telegraf_mem]
      processors: [resource, transform/strip_cpu_dim]
      exporters:  [awscloudwatch]

If you want more dimensions such as image version, add it to attributes.

Note: since only aggregate CPU usage is emitted by telegraf, we don’t need OTel to include a CPU dimension, so delete_key(attributes, "cpu") is used to remove this dimension.

Step 3: Install CloudWatch Agent in your Dockerfile

In your Dockerfile, install the CloudWatch Agent from the Amazon Linux repository. Then copy over the telegraf and OTel files to where the agent expects to retrieve them. Then configure your application’s entrypoint:

FROM public.ecr.aws/lambda/microvms:al2023-minimal

RUN dnf install -y --setopt=install_weak_deps=0 \
        python3 amazon-cloudwatch-agent \
    && dnf clean all

COPY app.py        /app/app.py
COPY cwagent.toml  /etc/cwagent.toml
COPY cwagent.yaml  /etc/cwagent.yaml
COPY entrypoint.sh /entrypoint.sh
RUN chmod +x /entrypoint.sh

CMD ["/entrypoint.sh"]

Step 4: Configure your Entrypoint to start CloudWatch Agent

Create a file called entrypoint.sh to start the CloudWatch Agent as a background process while executing your application in the foreground:

#!/usr/bin/env bash
set -euo pipefail

# Telegraf inputs (TOML) + OTel pipeline (YAML).
/opt/aws/amazon-cloudwatch-agent/bin/amazon-cloudwatch-agent \
    -config     /etc/cwagent.toml \
    -otelconfig /etc/cwagent.yaml &

exec python3 /app/app.py

This is everything you need to get CloudWatch running inside your MicroVMs!

Execution role requirements

To write the metrics to CloudWatch, ensure the MicroVM’s execution role has cloudwatch:PutMetricData permissions.

Verifying it works

To verify the metrics are being emitted, run the following command a few minutes after launching a MicroVM from your image:

aws cloudwatch list-metrics \
    --namespace LambdaMicroVms/Application \
    --dimensions Name=ImageName,Value=mem-python \
    --region us-west-2

You should see exactly three metric series per image: cpu_usage_active, cpu_usage_idle, and mem_used_percent.

Viewing the metrics

To view the metrics in the CloudWatch console, click “All Metrics”, and select the custom namespace LambdaMicroVms/Application (set in cwagent.yaml namespace field).

Here is an example for how it looks inside the console:

CloudWatch console showing CPU and memory metrics for a Lambda MicroVM

In the graph above, the application consumes ~2% memory (left axis) and < 0.1% CPU usage (right axis) when idle. The application then consumes ~9% of memory at the 30 minute mark, holds it for around 5 minutes, then releases it back to the OS. As it releases memory, we see memory utilization decrease. In this example, the MicroVM size is larger than the application needs – less than 10% of memory was used, indicating a smaller MicroVM size may be more economic for this workload.

If your CPU and/or memory utilization is below the baseline size configured (see MicroVM sizing), consider choosing a lower baseline to reduce your compute bill.

Conclusion

This post shows you how to configure and run the CloudWatch Agent inside your MicroVM image so you can collect CPU and memory metrics for MicroVMs launched from the image. This helps you monitor resource usage of your application as it is used, so you can right-size the MicroVM for your workload, debug service health, and check for scaling signals.

To get started, visit the AWS Lambda console, or install the AWS Lambda MicroVMs agent skill.

Track generative AI costs with Amazon Bedrock inference profiles

Post Syndicated from Erik Mack original https://aws.amazon.com/blogs/architecture/track-generative-ai-costs-with-amazon-bedrock-inference-profiles/

Tracking generative AI costs is a common challenge when multiple teams share a single foundation model through Amazon Bedrock. Your HR team answers policy questions, Accounting analyzes financial documents with it, and IT troubleshoots infrastructure issues. All three use the same foundation model. But usage shows up as one line item on the bill. As a result, finance can’t charge back each department, set per-team budgets, or identify who’s driving the most spend.

With Amazon Bedrock application inference profiles, you can solve this. An inference profile is a tagged wrapper around a foundation model. You can use it to attribute costs to specific teams or departments. By combining these profiles with AWS cost allocation tags, you can view per-department Amazon Bedrock costs as separate line items in AWS Cost Explorer.

In this post, we show you how to create application inference profiles for three departments and tag them for cost allocation. You also update your application to route invocations through department-specific profiles and view the per-department cost breakdown in Cost Explorer.

Solution overview

The following diagram shows the solution architecture. Users authenticate at the application layer, and the application identifies each user’s department. It then routes the request to that department’s tagged inference profile in Amazon Bedrock. All profiles use the same foundation model. The application calls Amazon Bedrock using a single IAM role, and individual user identities are not passed to AWS. Cost attribution comes from the inference profiles rather than the calling identity. Amazon Bedrock records usage against each profile’s Team tag, and AWS Cost Explorer displays the costs grouped by department.

Architecture diagram showing application routing to three tagged inference profiles pointing to one foundation model, with cost allocation flowing to AWS Cost Explorer

Figure 1 — Solution architecture for per-department cost tracking with application inference profiles

Amazon Bedrock can also attribute inference costs to the IAM principal that makes each call. This works well when each team calls Amazon Bedrock under a distinct IAM identity. In this architecture, a single application serves all departments under one role. Per-caller attribution can’t separate team costs without adding per-user session management. With application inference profiles, you can attribute costs per team by routing each team to a tagged profile.

To track costs per department:

  1. Create an application inference profile for each department, associating each one to the same foundation model.
  2. Tag each profile with a cost allocation tag (for example, Team=HR).
  3. Activate the tag in the AWS Billing and Cost Management console.
  4. Update your application to route invocations through each department’s inference profile Amazon Resource Name (ARN).
  5. View the per-department cost breakdown in Cost Explorer.

You pay the same per-token rate whether you invoke the model directly or through an inference profile – no additional charges for cost attribution.

Create and configure inference profiles for cost tracking

The following sections walk you through creating inference profiles, activating cost allocation tags, updating your application, and viewing costs in Cost Explorer.

Prerequisites

To configure this solution, you need the following:

  • An AWS account.
  • Model access enabled for your chosen foundation model in Amazon Bedrock (for instructions, refer to the Amazon Bedrock User Guide).
  • AWS Identity and Access Management (IAM) permissions including bedrock:CreateInferenceProfile, bedrock:TagResource, bedrock:InvokeModel, bedrock:InvokeModelWithResponseStream, ce:GetCostAndUsage, and ce:UpdateCostAllocationTagsStatus.
  • Access to the AWS Billing and Cost Management console to activate cost allocation tags and view Cost Explorer. For more information, refer to Managing access permissions for AWS Billing.
  • Python 3.12 with boto3 1.35.7 or later (for testing invocations).

Estimated time: 30 minutes (plus 24–48 hours for cost data to appear in Cost Explorer).

Estimated cost: Based on invocations at standard model pricing. For more information, refer to Amazon Bedrock Pricing.

Create application inference profiles

Create an application inference profile for each department. Each profile points to the same foundation model but has a unique tag for cost tracking.

To create an application inference profile:

  1. On the Amazon Bedrock console, in the navigation pane, choose Inference profiles.
  2. Choose the Application tab.
  3. Choose Create inference profile.
  4. For Profile name, enter HR.
  5. For Model, select your foundation model (for example, Anthropic Claude).

Note: Model availability varies by Region. Check the Amazon Bedrock model availability documentation for the current list.

To tag the inference profile:

  1. In the Tags section, choose Add tag.
  2. For Key, enter Team.
  3. For Value, enter HR.
  4. Choose Create. The inference profile status changes to Active.
  5. Repeat for Accounting (Tag: Team=Accounting) and IT (Tag: Team=IT).

The following figure shows the create inference profile page with the profile name and tag configured.

Amazon Bedrock console showing the create inference profile page with profile name HR and tag Team=HR configured

Figure 2 — Creating an application inference profile with a department tag

After you create all three profiles, the Application inference profiles list shows the HR, Accounting, and IT profiles, each with a status of Active and its corresponding Team tag. The following figure shows the three inference profiles after creation.

Amazon Bedrock console showing three application inference profiles: HR, Accounting, and IT

Figure 3 — Three application inference profiles, one per department

To provision inference profiles at scale (for example, one per team across dozens of teams), use the AWS::Bedrock::ApplicationInferenceProfile AWS CloudFormation resource instead of creating each profile manually.

Activate the cost allocation tag

After creating the inference profiles, you activate the cost allocation tag so that tagged costs appear in Cost Explorer. In multi-account environments using AWS Organizations, activate the Team cost allocation tag in the management (payer) account. Tagged usage from member accounts then consolidates in Cost Explorer. For more information about cost allocation tags, refer to Using AWS cost allocation tags.

To activate the cost allocation tag:

  1. Open the AWS Billing and Cost Management console.
  2. In the navigation pane, choose Cost allocation tags.
  3. In the search box, enter Team.
  4. Select the Team tag.
  5. Choose Activate.

The tag status changes to Active.

Note: Cost allocation tags are case-sensitive. Team and team are different tags. Tagged costs can take 24–48 hours to appear in Cost Explorer after activation.

Update the application to use inference profiles

To attribute costs to a department, pass the inference profile ARN as the modelId parameter instead of the foundation model ID. The API call remains the same. You only change the ID you pass.

To find an inference profile ARN:

  1. On the Amazon Bedrock console, choose Inference profiles.
  2. Select the profile.
  3. Copy the ARN from the details panel.

The ARN appears in the format arn:aws:bedrock:region:account-id:application-inference-profile/profile-id.

The following example shows how to route invocations based on the user’s department:

import boto3
from botocore.exceptions import ClientError

client = boto3.client('bedrock-runtime', region_name='us-east-1')

# Replace with your actual inference profile ARNs from the Amazon Bedrock console
DEPARTMENT_PROFILES = {
    'HR': 'arn:aws:bedrock:us-east-1:111122223333:application-inference-profile/abc123',
    'Accounting': 'arn:aws:bedrock:us-east-1:111122223333:application-inference-profile/def456',
    'IT': 'arn:aws:bedrock:us-east-1:111122223333:application-inference-profile/ghi789',
}

# Determine the user's department from your application's authentication layer.
# Examples:
# - An OIDC/SAML claim from your app login: token['custom:department']
# - A lookup in your user database: db.get_user_department(user_id)
# - A value stored in the user's session: session['department']
# The application then calls Amazon Bedrock using its own IAM role;
# individual user identities are not passed to AWS.
department = get_department_from_user_session()

if department not in DEPARTMENT_PROFILES:
    raise ValueError(f"Unknown department: {department}")

try:
    response = client.converse(
        modelId=DEPARTMENT_PROFILES[department],
        messages=[{'role': 'user', 'content': [{'text': 'Your prompt here'}]}],
        inferenceConfig={'maxTokens': 300}
    )
except ClientError as e:
    print(f"Error invoking model: {e}")
    raise

The full code is available on the GitHub repo.

When using inference profiles in production, validate user inputs and consider using Amazon Bedrock Guardrails to filter unintended content. API communications with Amazon Bedrock are encrypted in transit using Transport Layer Security (TLS). For more information about data protection, refer to Data protection in Amazon Bedrock.

Configure the IAM policy for Amazon Bedrock access

Because a single application calls Amazon Bedrock on behalf of all departments, it uses one IAM role. The following policy grants that role permission to invoke the department inference profiles and the underlying foundation model:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "InvokeDepartmentInferenceProfiles",
      "Effect": "Allow",
      "Action": [
        "bedrock:InvokeModel",
        "bedrock:InvokeModelWithResponseStream"
      ],
      "Resource": [
        "arn:aws:bedrock:us-east-1:111122223333:application-inference-profile/*",
        "arn:aws:bedrock:us-east-1::foundation-model/<your-in-region-model-id>"
      ]
    }
  ]
}

The wildcard (*) in the application inference profile ARN lets the single application role invoke the department profiles. The foundation model ARN is required because invoking through an inference profile needs permissions on both the profile and the underlying model. The application determines which department each request belongs to and routes it to the matching profile, and cost attribution comes from each profile’s Team tag. To further restrict access, replace the wildcard with the specific ARNs of your profiles.

Replace 111122223333 with your AWS account ID in the preceding policy.

To create the policy:

  1. On the IAM console, choose Policies.
  2. Choose Create policy.
  3. Choose the JSON tab.
  4. Paste the preceding policy.
  5. Choose Next.
  6. For Name, enter BedrockDepartmentAccessPolicy.
  7. Choose Create policy.

The BedrockDepartmentAccessPolicy appears in the policies list.

To attach the policy to a role:

  1. In the navigation pane, choose Roles.
  2. Select the role used by your application.
  3. Choose Add permissions.
  4. Choose Attach policies.
  5. Search for BedrockDepartmentAccessPolicy.
  6. Select BedrockDepartmentAccessPolicy.
  7. Choose Add permissions.

The BedrockDepartmentAccessPolicy appears in the role’s permission list. To add a department later, create another tagged inference profile and map it in your application. With the wildcard policy, no IAM change is needed. If you scoped the policy to specific ARNs, add the new profile’s ARN.

View per-department costs in Cost Explorer

To view the per-department breakdown in Cost Explorer:

  1. Open the Billing and Cost Management console.
  2. In the navigation pane, choose Cost Explorer.
  3. Set the date range to cover the period after you ran invocations.
  4. For Granularity, select Daily or Monthly.
  5. Choose Group by.
  6. Select Tag.
  7. Select Team.

To view exact amounts, scroll down to view the cost breakdown table.

The following figure shows the per-department cost breakdown in Cost Explorer. The bar chart displays a separately-colored segment for each department – HR, Accounting, and IT – with the cost amount for each. The table below the chart lists the exact dollar amount per department for the selected time period.

AWS Cost Explorer showing per-department Bedrock costs grouped by the Team tag

Figure 4 — Per-department Amazon Bedrock costs in Cost Explorer, grouped by the Team tag

After running invocations through each inference profile, verify the following:

  • Each inference profile shows the correct Team tag in the Amazon Bedrock console.
  • The Team cost allocation tag is active in the Billing and Cost Management console.
  • Per-department costs appear in Cost Explorer when you group by the Team tag.

If costs don’t appear after 48 hours, verify that the cost allocation tag is active and that invocations were made through the inference profile ARNs. If invocations fail, confirm that the inference profile status is Active and the IAM role has the required permissions.

Clean up

Inference profiles don’t incur charges on their own. You only pay for model invocations made through them. As a cleanup step, delete the inference profiles you created for this walkthrough to prevent accidental invocations.

Note: Deleting an inference profile immediately affects applications using that profile ARN. Verify that applications are not actively using these profiles before deletion. To recover, recreate the profile — note that it receives a new ARN, so update your application references.

Delete the following resources:

Conclusion

In this post, we showed you how to split generative AI costs by team using Amazon Bedrock application inference profiles and cost allocation tags. With this approach, you can see each department’s costs as a separate line item in Cost Explorer.

To add a new department, create another tagged profile. Costs show up as a separate line item.

You can also:

  • Set per-department spending alerts and control with AWS Budgets.
  • Detect unusual spending patterns with AWS Cost Anomaly Detection.
  • Monitor token usage per department with Amazon CloudWatch.
  • Attribute costs for higher-level Amazon Bedrock features – reference the same tagged profile ARN in the Knowledge Bases (RAG) to extend per-team attribution beyond direct model invocation.

For more background on application inference profiles, refer to Track, allocate, and manage your generative AI cost and usage with Amazon Bedrock.

For more information about inference profiles, refer to the Amazon Bedrock User Guide.

For help implementing this solution, contact your AWS representative.


About the author

Reducing Text2SQL latency with parameterized query templates

Post Syndicated from Yury Brukau original https://aws.amazon.com/blogs/architecture/reducing-text2sql-latency-with-parameterized-query-templates/

If your Text2SQL system takes 25-30 seconds to respond, user engagement drops significantly. For teams scaling beyond pilot projects, this latency gap between a working demo and a production-ready tool is the biggest barrier to adoption. Without caching, every question triggers a Large Language Model (LLM) call to generate SQL, and those calls introduce challenges: unpredictable response times, throttling limits, and token costs that grow linearly with traffic. Parameterized query templates provide an intelligent caching layer that in our production deployment, reduced end-to-end latency by 80% and cut token consumption by over 50%, turning a slow prototype into a responsive production system. In this post, we walk through the architecture behind this approach, covering the implementation details, performance results, and lessons learned from running a Text2SQL system in production.

When you move AI applications from pilot to production, you need solutions that scale under real traffic and perform consistently. Traditional caching strategies, storing expensive computations once and serving them many times, don’t translate directly to generative AI. End users rarely phrase the same question the same way, context varies between sessions, and outputs depend on small input variations. Yet the underlying principle (caching) still holds value. Rather than abandoning caching entirely, the key is finding the right abstraction layer where similar requests can share cached results.

Solution overview

We applied the solution described in the following section to a system where business users query operational databases using natural language. You ask questions like “What were total sales in Q3?” or “Show me top performing products this month?” and the system generates SQL queries, executes them against the database, and returns results in conversational format. The system translates natural language to SQL using Amazon Bedrock foundation models, while AWS Lambda orchestrates the workflow. You can see a basic overview of used architectural components in Diagram 1.

Architecture diagram showing the Text2SQL system with Amazon Bedrock for SQL generation and AWS Lambda for workflow orchestration

Figure 1 — Solution overview architecture

During the initial implementation phase, the approach with generating and executing SQL queries for user questions on the fly worked well. Response quality was high, and users found the interface intuitive. After these positive results, we started looking into scaling the solution for production traffic. Preserving accuracy was the main priority. Experiments with smaller, faster models didn’t provide a good trade-off between query quality and latency reduction. The accuracy degradation wasn’t acceptable for our system.

This led us to explore alternative approaches, and caching naturally came to mind. Caching user question and answer pairs is the most straightforward option, but it has a fundamental limitation: underlying data changes constantly. An answer about Q3 sales cached today becomes incorrect as soon as new transactions are recorded. The cache would need constant invalidation, undermining its purpose.

Caching the SQL query instead solves this problem. A query like:

SELECT SUM(revenue) FROM sales WHERE quarter = 'Q3'

always fetches fresh data when executed, regardless of when it was cached. Structured Query Language (SQL) captures the user’s intent in a structured, deterministic form that remains valid even as data evolves. It also happens to target the most time and token consuming step in the pipeline, since generating SQL queries requires sending full schema context and examples to a frontier model.

Analyzing the generated queries revealed an opportunity to go further. Many queries follow the same structure, different only in their filter values. A question about Q3 sales produces:

SELECT SUM(revenue) FROM sales WHERE quarter = 'Q3'

while Q2 sales produce:

SELECT SUM(revenue) FROM sales WHERE quarter = 'Q2'

The same pattern appeared across product lookups, date ranges, and category filters. This led to the templating approach: instead of caching complete queries, we generalize them into templates with placeholders. A single template now covers an entire family of questions:

SELECT SUM(revenue) FROM sales WHERE quarter='{quarter}'

Flow diagram showing a cache hit path where a user question matches a stored template, fills placeholders with extracted entities, and executes the SQL query directly

Figure 2 — Templated SQL query cache hit

Templating solves the limited reusability of plain user question, but still leaves a challenge: how do you match an incoming question to the right template when users phrase things differently? “Show me Q3 sales” and “What were sales in Q3?” ask for the same data but share few words. Traditional string matching or keyword lookup would miss these connections. We address this by storing each template alongside a vector embedding of its original question. When a new question arrives, we compute its embedding and perform semantic similarity search against the cache. Because embeddings capture meaning rather than surface wording, both phrasings map to the same template with high confidence. If a match is found above a confidence threshold, we extract entities from the question using lightweight named entity recognition, fill the template placeholders, and execute the query directly, bypassing the LLM entirely. In Diagram 2, you can see the flow of a cache hit.

For questions without matching templates, the system falls back to full LLM generation. It then generalizes the newly generated query into a template, pairs it with the question’s embedding, and adds it to the cache. This creates a self-improving system where cache coverage grows organically as more query patterns are encountered.

Walkthrough – Text2SQL pipeline

The following sections describe each step of the template caching pipeline. Each user’s question flows through entity extraction, template retrieval, and SQL query execution. Cache misses trigger full LLM generation, with new queries feeding back into the cache. The following diagram shows the complete flow of a user question through the newly introduced caching layer.

Complete pipeline flow showing entity extraction, template retrieval, template filling, response generation, and the reinforcement loop for cache growth

Figure 3 — Text2SQL pipeline with template caching layer

1. Entity extraction

After a user submits a question, the system performs entity recognition to extract named entities and values. This step considers not only the current question but also conversation history, current date, and user preferences. This context helps resolve ambiguous references like “last month” or “my region”. Using a lightweight model like Amazon Nova 2 Lite or a custom-trained named entity recognition (NER) model, we identify entities such as dates (“Q3 2024”), names (“Product X”), categories (“electronics”), and numeric values (“top 10”). The system stores these extracted entities separately and uses them later to fill out template placeholders.

The system converts the user’s question into an embedding vector using the same embedding model used during cache population. This vector queries the template cache through semantic similarity search, returning the closest matching templates above a confidence threshold. The search matches based on the question’s intent and structure rather than exact wording, so “What were Q3 sales?” and “Show me revenue for third quarter” both match the same template despite different phrasing.

It’s important to note that the confidence threshold governs the cache retrieval layer’s precision-recall trade-off. Set it too high and the system rejects valid, differently worded questions, forcing it to build SQL from scratch. Set it too low and loosely related templates slip through, risking confident answers built on the wrong query. The right value is domain-dependent: narrow, well-templated domains tolerate stricter thresholds, while broad or sparsely covered ones need looser ones.

Rather than relying on a single threshold, we suggest monitoring retrievals in production, logging matched templates and their similarity scores, so we can see when valid questions are being rejected or unrelated templates are slipping through. When embedding similarity alone doesn’t give enough precision, we added a lightweight reranking step: first we retrieve a broader set of candidate templates with a looser threshold, then re-score them with a small LLM or a specialized reranker model to select the best match. This improves precision without sacrificing recall and still costs far less than generating SQL from scratch.

3. Template filling and query execution

When a matching template is found, the system maps extracted entities to template placeholders. If the template contains `{quarter}` and entity recognition extracted “Q3”, the system replaces the placeholder with the actual value. The system validates the filled SQL query for syntax correctness, then executes it directly against the database. This path bypasses the time and token intensive LLM call that generates the SQL query.

This design helps the system to protect against SQL injection on two levels. First, it validates each extracted entity against the expected format for its placeholder: a `{quarter}` must match a known set of values, a `{date}` must parse as a valid date, a numeric threshold must be a number. The system rejects values that do not pass validation before they ever reach the query. Second, the system fills the placeholders using parameterized database queries (prepared statements) rather than string interpolation, so the parameterized query mechanism treats entity values as data rather than executable SQL. This approach also catches entity-extraction errors, improving answer reliability beyond the security benefit.

For richer responses, the system can retrieve multiple top-K similar templates and execute them in parallel. This provides additional context and related information beyond the primary query, for example returning both: quarterly sales totals and a breakdown by product category. The parallel execution adds minimal latency while delivering more comprehensive answers.

4. Response generation and validation

After executing the query, the system sends results to a response generation model. This model has two jobs, both handled in a single call: judge whether the results answer the question, and, if they do, summarize them into a conversational response.

The sufficiency check is driven by instructions in the prompt. The system instructs the model to confirm that the results are non-empty, that they contain the fields the question asked about, and that they cover every part of the question rather than only some of it. For example, if a user asks for “Q3 sales by region” but the matched template returns only a Q3 total, the results are incomplete, and the model is instructed to flag them as insufficient instead of answering with partial data. The model returns this judgment as a structured signal alongside its response, so the pipeline can branch on it deterministically. This step helps verify that users receive accurate answers rather than partial or misleading information from imperfect template matches.

This task is fundamentally simpler than SQL generation: instead of writing structured code from natural language, the model only needs to read tabular data and either summarize it or declare it insufficient. Because the task is simple, a smaller, faster model like Claude Haiku 4.5 can handle it effectively.

On a cache hit, there is only a single lightweight LLM call, which improves both latency and cost thanks to the smaller model. On a cache miss, the model flags the template results as insufficient and the system falls back to full SQL generation before producing the answer, for three calls in total: the sufficiency check, the SQL generation, and the response. That is one call more than the uncached pipeline, so misses carry extra latency. The trade-off is favorable because the added call is the cheap sufficiency check rather than another expensive generation, and because at a healthy hit rate the savings on hits outweigh the penalty on misses.

5. Fallback to full generation

If no template matches the confidence threshold, or if the validation step determines that cached results are insufficient, the system falls back to the standard Text2SQL pipeline. The question, along with the full context, goes to the foundation model for SQL generation. The generated query executes against the database, and results return to the user. Importantly, this newly generated query doesn’t disappear. It enters the reinforcement loop.

6. Reinforcement loop for cache growth

After a successful fallback generation, the system evaluates whether the new query should join the template cache. If the query executed successfully and returned valid results, it becomes a candidate for templating. The system generalizes the query by replacing specific values with placeholders and computes the original question’s embedding. It then adds this new template-question pair to the vector store, expanding cache coverage. Over time, the cache grows organically to cover query patterns specific to your users’ actual needs.

Results and performance gains

The figures in this section come from our production deployment but treat them as an illustrative model rather than a fixed benchmark. Exact token counts and latencies depend on your schema size, prompt design, model choice, and query mix. What generalizes is the direction of the improvement, not the specific numbers.

The dominant cost and latency in a Text2SQL request come from a single step: generating the SQL query. That call sends the user question, conversation history, the database schema, few-shot examples, and domain guidance to a powerful LLM such as Anthropic Claude Sonnet, which is needed to produce reliable queries. In our deployment this prompt runs on the order of 60K input tokens for a few hundred output tokens, and takes roughly 15-20 seconds. Every other step: embedding, vector search, template filling, and query execution, is minor by comparison. Entity recognition, runs on a dedicated NER model hosted on Amazon SageMaker AI rather than an LLM, adding negligible cost and latency next to SQL generation. Optimizing the pipeline is therefore mostly about avoiding that one expensive call.

On a cache hit, the system skips SQL generation entirely. What remains is response summarization, turning the query results into a conversational answer, which runs on a small model with a small prompt (on the order of a couple thousand input tokens). Because summarization is needed on both, the cached and uncached paths, a cache hit does not remove tokens completely, but it eliminates the 60K-token generation call, cutting token consumption by roughly 90% on that request.

This 90% is the saving on a single cache hit. Overall cost depends on the average across all requests, since cache misses still incur the full generation cost. At the roughly 60% hit rate we observed in production, the blended reduction across all traffic comes out above 50%. Latency follows the same pattern. An uncached request spends 15-20 seconds on the SQL call, retries and error handling included, then a few more seconds on summarization, putting a typical request in the 25-30 second range. On a cache hit, retrieval, template filling, and execution finish well under a second, and the remaining time is almost entirely the summarization call. That brings the end-to-end cache-hit path under 5 seconds, roughly an 80% reduction, or about 6x faster. It also pinpoints where the residual latency comes from: not the cache lookup, but the one LLM call that still has to run.

These per-request gains only matter if cache hits are common. In our production system the hit rate reached about 60% after roughly two weeks of active use, though the achievable rate depends heavily on the domain and how repetitive the queries are. Cache misses run the full pipeline plus the small sufficiency check, so they cost marginally more than a purely uncached request, which means the net gain comes entirely from hits. As the reinforcement loop keeps adding templates, the hit rate climbs and both the cost and latency benefits continue to compound.

Conclusion

Scaling AI applications to production often requires rethinking traditional optimization strategies. In this post, we demonstrated how template-based caching addresses the latency and cost challenges of Text2SQL systems without sacrificing accuracy. By caching SQL query structures rather than complete responses and using semantic similarity to match user questions to templates, the system can bypass expensive LLM inference calls. The reinforcement loop ensures cache coverage grows organically based on actual usage patterns.

In practice, this means: 6x faster response times on cache hits, inference costs decrease proportionally to your cache hit rate, and accuracy remains high because templates are generated by the most capable models. The patterns we covered, such as semantic matching, output generalization, entity extraction, and continuous improvement loops, extend beyond Text2SQL to any AI system where similar requests should produce structurally similar outputs.

Further reading

Generating value from enterprise data: Best practices for Text2SQL and generative AI

Enterprise-grade natural language to SQL generation using LLMs: Balancing accuracy, latency, and scale

Build a robust text-to-SQL solution generating complex queries, self-correcting, and querying diverse data sources

Text-to-SQL solution powered by Amazon Bedrock

Amazon S3 Vectors: First cloud storage with native vector support at scale

Amazon Nova 2 Lite

About the authors

How AWS IAM role manager rethinks the starting point for IAM roles

Post Syndicated from Zach Jiang original https://aws.amazon.com/blogs/security/how-aws-iam-role-manager-rethinks-the-starting-point-for-iam-roles/

When you build a new application or capability on Amazon Web Services (AWS), you want to focus on what you’re building. Getting a service running almost always begins with AWS Identity and Access Management (IAM). Many AWS services that act on your behalf need an IAM role, an identity the service assumes to access your resources with a defined set of permissions. You then author a trust policy so the service can assume the role, choose the permissions the workload needs, and attach it. Configuring roles and policies for common patterns is repeatable work that doesn’t need to be manual.

IAM role manager does that work for you. When role manager is enabled, AWS creates and configures the IAM roles as you build in supported service consoles, so you can start using a service and let AWS handle the role behind it. You create the resource you want, and role manager provisions and attaches the role you need as part of the same flow, so you can build now and refine permissions as your workload matures.

With that step automated, getting started takes minutes. You can create an AWS Lambda function and start running your code, with its execution role already created and attached, without switching context to set one up. Role creation becomes an automated part of building your application rather than a separate step.

Role manager is especially useful when you’re getting started: the moments when you want to stand up a service or get a proof of concept running and want to defer role configuration until later in your development process. You don’t need prior IAM experience to get started. You keep full control of what it creates, because the roles are ordinary IAM roles that you can view, edit, or delete like any role you author yourself. When you want to tighten a role, AWS IAM Access Analyzer reviews how it has been used and recommends a policy scoped to only the permissions it needs.

How to enable role manager

Role manager has two states, enabled and disabled. Enabling it for an account authorizes AWS to create roles in that account. In an organization, administrators can use a service control policy (SCP) to control whether member accounts can enable or use role manager. To enable it:

  1. Open the IAM console and choose Account settings.
  2. In the role manager section, choose Enable.
Figure 1: Enable Role Manager

Figure 1: Enable Role Manager

Some AWS services already create a role for you when you create a resource that needs one. Role manager doesn’t change that: those services keep creating roles automatically, and roles you already created keep working. What role manager adds is a single account-level control, and coverage for a case that built-in flows can’t handle: tasks whose permissions AWS can’t determine in advance, such as running your own code. For those tasks, role manager provisions a role that you can narrow later.

Example: Create an Amazon EventBridge rule

Start with a common task: an Amazon EventBridge rule that invokes a target, such as an Amazon Simple Queue Service (Amazon SQS) queue or an Amazon Simple Notification Service (Amazon SNS) topic. Without role manager, you would pause here to create a role that lets EventBridge invoke the target, write the role’s trust policy, attach the required permissions, and then return to finish the rule. With role manager enabled, you define the rule and its target, choose Create, and role manager provisions the role and attaches it for you. The EventBridge console shows the rule created and ready, and you never open the role-creation flow.

Figure 2: Creating an EventBridge rule with no manual role setup

Figure 2: Creating an EventBridge rule with no manual role setup

The role comes from an AWS managed role template: a definition AWS builds and maintains for a specific task, with the trust policy and permissions already worked out. The console calls a new IAM API, AcquireRole, which finds the matching template, provisions the role from it, and returns it to EventBridge. Depending on the service, AcquireRole either creates a new role or reuses one that already fits, so an account does not fill up with duplicate roles for the same task.

Role manager creates the role using your own IAM permissions, not a separate role-manager permission. To provision a new role, you need permission for the actions the template performs: at minimum, you need permissions to create and attach roles. When AcquireRole reuses an existing role instead of creating one, it needs only iam:GetRole and iam:GetRoleTemplateVersion. If you’re missing either of these permissions, the console tells you which one is needed rather than creating the role.

Run code that calls other AWS services

Not every task has a set of permissions AWS can define in advance. When a role runs your own code, such as a Lambda function, AWS has no way of knowing which services that code will call. Role manager covers this case too: create a Lambda function with role manager enabled, and it attaches an execution role that your code can use right away and that you can narrow once you know what the function calls.

Because the permissions your code needs aren’t known up front, role manager attaches the AWS managed policy PowerUserAccess to the role. PowerUserAccess grants access to AWS services so your function can call what it needs. By design, it doesn’t grant permission to manage IAM, AWS Organizations, or account settings. The template also configures the role to trust only the Lambda service.

Figure 3: Create an AWS Lambda function with no manual role setup

Figure 3: Create an AWS Lambda function with no manual role setup

Role manager attaches an execution role, and your function is ready to run. Figure 4 shows the Execution role panel on the function’s Configuration tab, with the role that role manager attached.

Figure 4: Role manager provides a role automatically to an AWS Lambda function

Figure 4: Role manager provides a role automatically to an AWS Lambda function

You can open the role in the IAM console to review its permissions. Figure 5 shows the role’s Permissions tab with the PowerUserAccess policy attached.

Figure 5: Permissions of the role provided by role manager for an AWS Lambda function

Figure 5: Permissions of the role provided by role manager for an AWS Lambda function

You keep full visibility into what role manager creates. Every role it creates records the role template it came from, and both GetRole and ListRoles return that template reference. You can inspect any role in your account and tell which were created by role manager. You read a role’s trust policy and permissions the same way you would for a role you authored, and AWS CloudTrail records each role’s creation.

Refining roles as workloads mature

As your workloads mature, refine the roles that role manager created to follow least privilege. When you’re ready, you can disable role manager and get IAM Access Analyzer unused access analysis at no additional cost for 90 days. Access Analyzer looks at how each role has been used and recommends a policy you can apply that keeps only the permissions the role needs. Start with the roles attached to your most critical workloads and work outward.

Disabling role manager doesn’t disrupt anything already running: your resources keep the roles they have, those roles stay in your account until you change them, and from that point you author new roles yourself, the same as before. If you would rather narrow a single role than the whole account, editing that role removes it from role manager’s control and it becomes a standard customer-managed role, with your changes preserved. In sandbox or development accounts, keeping role manager enabled saves time. For production workloads, disable role manager and refine the roles it created to least privilege before going live.

Conclusion

Role manager automates IAM role setup so you can focus on building from the start. When you enable it, AWS creates and attaches the IAM roles your resources need as you build, so you can start in minutes without prior IAM experience. Because these are IAM roles that you fully control, you keep the same visibility and the same tools you already use. Keep role manager enabled while you build, and refine the roles it created as your workloads mature.

To get started, enable role manager in the IAM console and create a resource in a supported service. To learn more, see IAM role creation and the list of supported services in the IAM User Guide.

If you have feedback about this post, submit comments in the Comments section below.


Zach Jiang

Zach Jiang

Zach is a Senior Technical Product Manager at AWS, specializing in AWS Identity products. He focuses on making identity the easy part of building on AWS for customers. Outside of technology, Zach enjoys traveling and exploring new cultures and cuisines.

David Sing

David Sing

David is a Principal Product Manager at AWS, specializing in AWS IAM. He focuses on simplifying IAM for builders and AI agents, safe credential issuance for AI agents, and authorization policy governing agent access. Outside of technology, David enjoys economics and markets, fishing, and time outdoors with his family.

Punit Deotale

Punit Deotale

Punit is a Software Development Manager on the AWS IAM team. He leads work on making it easier for customers to create and manage IAM roles directly within AWS service workflows, so they can set up the right permissions without leaving what they are doing. His focus is reducing permission-setup friction across AWS while helping customers stay aligned with least privilege. Outside of work, Punit enjoys reading, building side projects, and being outdoors.

Burst to Region: Overflow AWS Outposts workloads to Amazon EC2

Post Syndicated from Diya . original https://aws.amazon.com/blogs/compute/burst-to-region-overflow-aws-outposts-workloads-to-amazon-ec2/

AWS Outposts brings AWS infrastructure into your data center, giving on-premises workloads the low latency and data locality they need. But unlike the AWS Region, an Outposts rack has a fixed amount of compute. When your workload needs more instances than the rack can provide, you have two options: drop requests, or overflow them somewhere with room to grow. This post shows you how to automate the second option. You build a Burst to Region pattern that detects capacity constraints on your Outpost, launches Amazon Elastic Compute Cloud (Amazon EC2) instances in the parent Region, gradually shifts traffic to them, and returns traffic to local instances once capacity recovers.

To implement this pattern you configure Amazon CloudWatch, Amazon Simple Notification Service (Amazon SNS), AWS Lambda, Amazon EC2 Auto Scaling, Elastic Load Balancing (Application Load Balancer), and Amazon EventBridge. You trade a moderate latency increase for continued availability during capacity events.

When to use this pattern

This pattern assumes your Outposts workload scales out through Amazon EC2 Auto Scaling. Burst to Region reacts to instance-capacity exhaustion on the rack. It engages when your workload tries to launch more instances than the available Outpost capacity supports. If your fleet is fixed size and degrades under load without scaling out, the capacity alarm never fires and overflow never triggers. For those workloads, monitor per-instance saturation (CPU, latency) separately.

Good candidates prefer local capacity but can tolerate Region latency under pressure. If your application runs on Outposts for proximity yet degrades gracefully when some traffic takes the longer path to the Region, it fits this pattern. Examples include:

  • Internal enterprise applications.
  • Stateless web frontends and API layers.
  • Pre-processing tiers where single-digit to tens-of-milliseconds additional round-trip latency during peaks is acceptable.

Poor candidates cannot absorb any added latency or must stay on the Outpost. Avoid this pattern for:

  • Applications with sub-millisecond requirements.
  • Workloads with strict data residency or sovereignty mandates that prevent traffic from leaving the on-premises environment.
  • Real-time control systems with hard timing constraints.
  • Applications tightly coupled to on-premises data stores with no Region replica.

The core tradeoff is explicit. During capacity events, you accept moderately higher latency to maintain availability. If your workload cannot tolerate any latency increase, keep it pinned to Outposts and reserve capacity through other means, such as Capacity Reservations.

Solution overview

Burst to Region works in three moves: detect capacity pressure on the Outpost, launch overflow compute in the parent Region, and shift traffic gradually until local capacity recovers. Six AWS services coordinate to make this automatic. The following diagram shows the reference architecture for the Burst to Region pattern, illustrating how the six AWS services interact during capacity detection, overflow scaling, traffic distribution, and recovery.

Reference architecture for Burst to Region on AWS Outposts showing capacity detection, overflow scaling, traffic distribution, and recovery

Figure 1: Reference architecture for Burst to Region on AWS Outposts

The pattern uses six AWS services working together:

  • Amazon CloudWatch monitors Outposts capacity utilization and raises alarms.
  • Amazon SNS provides event fan-out from alarm to orchestrator.
  • AWS Lambda orchestrates the burst logic (scale-out, weight adjustment, recovery)
  • Amazon EC2 Auto Scaling manages the overflow fleet lifecycle.
  • Application Load Balancer distributes traffic across both locations using weighted target groups.
  • Amazon EventBridge handles periodic recovery evaluation.

You must configure five phases for this pattern:

  1. Monitor. CloudWatch tracks Outposts capacity utilization metrics in the AWS/Outposts namespace.
  2. Detect. A CloudWatch alarm fires when utilization exceeds a threshold (for example, 80%).
  3. Overflow. The alarm triggers a Lambda function through Amazon SNS. Lambda scales out a Region-based Amazon EC2 Auto Scaling group and adjusts ALB target group weights.
  4. Distribute. The ALB splits traffic between Outposts instances and Region instances using weighted forwarding.
  5. Recover. An Amazon EventBridge scheduled rule periodically evaluates capacity. When Outposts recovers, Lambda scales down the overflow fleet and returns all traffic to local instances.

Design decisions

We chose Application Load Balancer with weighted forwarding over Amazon Route 53 weighted routing for traffic distribution. ALB provides health-aware routing to only healthy overflow instances and target group stickiness for session consistency. Weight changes take effect for new connections after calling the ModifyRule API. DNS-based shifting through Route 53 provides too coarse control for rapid weight adjustments, and TTL propagation delays make recovery slower.

The burst orchestrator runs as a Lambda function rather than a long-running service. It executes only during state transitions, so there is no steady-state compute cost. Lambda integrates natively with Amazon SNS and Amazon EventBridge for event-driven invocation without additional infrastructure.

You implement recovery with an Amazon EventBridge scheduled rule (every 5 minutes) rather than relying solely on the CloudWatch alarm to return to OK state. The alarm confirms capacity is available, but does not confirm that overflow instances have drained active connections. The scheduled rule provides gradual, safe scale-down.

Implementation

This section walks through the key components of the Burst to Region pattern. For the complete deployable AWS SAM template, see the GitHub repository.

Prerequisites

To deploy this pattern, you need:

  • An AWS account with a configured AWS Outposts rack.
  • An Amazon Virtual Private Cloud (Amazon VPC) with subnets associated with your Outposts and subnets in the parent AWS Region.
  • IAM permissions to create CloudWatch alarms, Lambda functions, Auto Scaling groups, and ALB resources.
  • AWS Serverless Application Model (AWS SAM) CLI installed and configured.
  • Existing Amazon EC2 Auto Scaling group running on your Outpost (these become your baseline fleet)
  • A custom domain name with a DNS record (Route 53 alias or CNAME) pointing to your Application Load Balancer, and an AWS Certificate Manager (ACM) certificate for that domain to enable HTTPS.

Capacity monitoring and alarm

The CloudWatch alarm monitors instance utilization on the Outpost and triggers the burst workflow when capacity is constrained.

The InstanceTypeCapacityUtilization metric reports the percentage of a given instance type’s capacity in use. Note that this metric includes capacity consumed by managed services such as Amazon Relational Database Service (Amazon RDS) or Application Load Balancer running on the Outpost — not only your application’s EC2 instances. Factor this into your threshold planning.

OutpostsCapacityAlarm:
  Type: AWS::CloudWatch::Alarm
  Properties:
    AlarmName: outposts-capacity-high
    Namespace: AWS/Outposts
    MetricName: InstanceTypeCapacityUtilization
    Dimensions:
      - Name: OutpostId
        Value: !Ref OutpostId
      - Name: InstanceType
        Value: !Ref OutpostInstanceType
    Statistic: Average
    Period: 300
    EvaluationPeriods: 2
    Threshold: !Ref CapacityThreshold
    ComparisonOperator: GreaterThanOrEqualToThreshold
    AlarmActions:
      - !Ref BurstSNSTopic
    TreatMissingData: notBreaching

Why these values matter:

  • Period: 300 and EvaluationPeriods: 2 require 10 minutes of sustained high utilization before triggering. This avoids false alarms from transient spikes.
  • Threshold: 80 (recommended starting point) leaves a 20% buffer. A threshold set too high (95%) risks launch failures before the overflow fleet is ready. A threshold set too low (50%) causes unnecessary bursts.
  • TreatMissingData: notBreaching prevents false alarms when data points are missing. Since this alarm is scoped to a single instance type, treating missing data as breaching could trigger unnecessary bursts when the instance type is simply not in use.
  • Separate scale-out from scale-in: This alarm triggers burst scale-out at 80%. Recovery is handled separately by the Amazon EventBridge scheduled rule, which uses a lower threshold (for example, 60%) before scaling in. This hysteresis gap prevents flapping where scaling down immediately pushes utilization back above the alarm threshold.

Burst orchestrator (Lambda)

The Lambda function handles two event paths: alarm-triggered scale-out and scheduled recovery evaluation. The following pseudocode shows the orchestration flow:

def handler(event, context):
    # Route based on event source
    if is_scheduled_recovery(event):
        return handle_recovery_check()

    alarm_state = parse_sns_alarm_state(event)

    if alarm_state == 'ALARM':
        # Scale out the overflow Auto Scaling group
        scale_out_overflow(desired=OVERFLOW_CAPACITY)
        # Don't shift traffic yet --- wait for healthy instances
        publish_burst_metric(active=True)


def handle_recovery_check():
    """Called every 5 minutes by EventBridge."""
    # Check if burst is active
    if not is_burst_active():
        return

    # If overflow instances are healthy and registered, shift traffic
    if overflow_targets_healthy():
        current_weights = get_current_alb_weights()
        if current_weights['region'] == 0:
            # First shift --- instances are now warm
            set_alb_weights(outposts=90, region=10)
        elif needs_more_overflow():
            step_up_region_weight()

    # If Outposts capacity has recovered, begin scale-down
    if outposts_capacity_recovered():
        step_down_region_weight()
        if get_current_alb_weights()['region'] == 0:
            # All traffic back to Outposts, drain and terminate overflow
            wait_for_connection_draining()
            scale_down_overflow(desired=0)
            publish_burst_metric(active=False)

The key actions the function performs:

  • scale_out_overflow — Sets the overflow Auto Scaling group desired capacity from 0 to your configured burst size.
  • set_alb_weights — Calls the ModifyListener API to adjust weighted forwarding between the Outposts and Region target groups.
  • publish_burst_metric — Writes a custom CloudWatch metric (BurstActive) for dashboard visibility.
  • handle_recovery_check — Called every 5 minutes by Amazon EventBridge. Confirms Outposts capacity has recovered, steps weights back gradually, waits for connection draining, then scales down the overflow fleet.

Important: The orchestrator does not shift ALB weights immediately upon scale-out. It waits for the next Amazon EventBridge invocation (up to 5 minutes) to confirm that overflow instances have passed health checks and are registered as healthy in the target group. This helps prevent routing traffic to instances that have not finished launching.

For the production-ready implementation with error handling, gradual weight stepping, and connection draining verification, see the GitHub repository.

Overflow Auto Scaling group

The overflow fleet starts at zero and scales only when the Lambda function sets desired capacity during a burst event:

OverflowASG:
  Type: AWS::AutoScaling::AutoScalingGroup
  Properties:
    AutoScalingGroupName: burst-overflow-fleet
    LaunchTemplate:
      LaunchTemplateId: !Ref OverflowLaunchTemplate
      Version: !GetAtt OverflowLaunchTemplate.LatestVersionNumber
    MinSize: 0
    MaxSize: !Ref MaxOverflowCapacity
    DesiredCapacity: 0
    VPCZoneIdentifier:
      - !Ref RegionSubnet1
      - !Ref RegionSubnet2
    TargetGroupARNs:
      - !Ref RegionTargetGroup
    HealthCheckType: ELB
    HealthCheckGracePeriod: 120
    MetricsCollection:
      - Granularity: 1Minute

The overflow fleet starts at zero capacity and incurs no cost at rest. During a burst event, the Lambda function calls the SetDesiredCapacity API to launch overflow instances. During recovery, it sets desired capacity back to zero.

The launch template mirrors your Outposts instance type to maintain consistent performance characteristics across both locations.

ALB weighted forwarding

The ALB listener uses weighted forwarding across two target groups. In steady state, all traffic goes to Outposts (weight 100/0). During burst, the Lambda function adjusts these weights dynamically using the ModifyListener API. Clients reach the ALB through a DNS record — either a Route 53 alias or a CNAME pointing to the ALB’s DNS name.

ALBListener:
  Type: AWS::ElasticLoadBalancingV2::Listener
  Properties:
    LoadBalancerArn: !Ref ApplicationLoadBalancer
    Port: 443
    Protocol: HTTPS
    SslPolicy: ELBSecurityPolicy-TLS13-1-2-2021-06
    Certificates:
      - CertificateArn: !Ref CertificateArn
    DefaultAction:
      Type: forward
      ForwardConfig:
        TargetGroups:
          - TargetGroupArn: !Ref OutpostsTargetGroup
            Weight: 100
          - TargetGroupArn: !Ref RegionTargetGroup
            Weight: 0
        TargetGroupStickinessConfig:
          Enabled: true
          DurationSeconds: 300

RegionTargetGroup:
  Type: AWS::ElasticLoadBalancingV2::TargetGroup
  Properties:
    Name: burst-region-targets
    Protocol: HTTP
    Port: 80
    VpcId: !Ref VpcId
    HealthCheckEnabled: true
    HealthCheckIntervalSeconds: 30
    HealthCheckPath: /health
    HealthyThresholdCount: 2
    UnhealthyThresholdCount: 3
    TargetGroupAttributes:
      - Key: deregistration_delay.timeout_seconds
        Value: "300"
      - Key: slow_start.duration_seconds
        Value: "120"

Note on stickiness: Target group stickiness keeps a client pinned to whichever target group served its first request for DurationSeconds. We set this to 300 seconds (5 minutes) to match the Amazon EventBridge evaluation interval. This balances session consistency for stateful workloads against the need for weight changes to take effect within a reasonable window. For purely stateless workloads, you can disable stickiness entirely to allow immediate weight convergence. For workloads requiring longer session affinity, increase the duration but understand that weight transitions will converge more slowly — existing sticky sessions continue going to the original target group until they expire.

Traffic weight progression

Use stepped transitions rather than abrupt weight changes. The following table shows the recommended progression:

Phase Outposts weight Region weight Condition to advance
Normal 100 0 Steady state
Burst step 1 90 10 Region target group has at least 1 healthy host
Burst step 2 70 30 Region target group healthy for 2 consecutive checks
Burst step 3 50 50 Only if Outposts capacity exceeds 95% used
Recovery step 1 80 20 Outposts capacity below 70%
Recovery step 2 100 0 Outposts capacity below 60% for 2 checks

Avoid jumping directly from 0% to 50% Region traffic. Cold overflow instances need time to warm caches and stabilize before absorbing significant load.

Best practices

Apply these best practices to get the most from this pattern while avoiding common pitfalls.

Traffic tiering

Classify your workloads into two tiers at the ALB listener level. Latency-critical paths use routing rules with the Outposts target group only. These never overflow regardless of capacity state. Overflow-eligible paths use the weighted forwarding rule. This separation helps make sure that your most latency-sensitive flows are not impacted by the burst mechanism.

Managing data gravity

For stateless workloads, Burst to Region requires no special data handling. For workloads with session state or shared data:

Anti-pattern: Do not burst workloads that write to Outposts-local storage and expect synchronous consistency. The latency and complexity of cross-location writes defeats the purpose of the pattern.

Cost optimization

The overflow fleet consumes On-Demand pricing by default since it starts at zero and scales only during peaks.

Burst profile Recommended pricing Rationale
Unpredictable spikes (minutes) On-Demand Maximum flexibility, no commitment waste
Predictable daily peaks (hours) Savings Plans (Compute) Covers overflow hours at discount
Frequent, long bursts Reserved capacity plus On-Demand Baseline discount plus burst flexibility

Monitor your BurstActive custom metric over time. If overflow is active more than 30% of the time, you likely need additional Outposts capacity rather than relying on Region overflow.

Security consistency

Maintain identical security posture across both environments:

  • Use the same security group rules for Outposts and Region instances.
  • Deploy with AWS CloudFormation StackSets to support consistency.
  • Share the same IAM instance profile. The overflow launch template references the same role as your Outposts instances.
  • Apply the same AWS Systems Manager patch baselines and compliance rules to both fleets.

Observability

Build a CloudWatch dashboard that provides visibility into burst state and performance. The SAM template in the repository deploys a pre-configured dashboard tracking:

  • Burst status: Custom BurstActive metric (1 = active, 0 = normal)
  • Capacity headroom: UsedInstanceType_Count compared to AvailableInstanceType_Count. Note that UsedInstanceType_Count includes instances consumed by managed services (Amazon RDS, ALB), so your available application capacity may be lower than the raw availability count suggests.
  • Overflow fleet size: Auto Scaling group GroupInServiceInstances.
  • Latency comparison: TargetResponseTime per target group (Outposts compared to Region)
  • Traffic distribution: RequestCount per target group.

Set a CloudWatch alarm on Region target group TargetResponseTime exceeding your acceptable threshold. This provides early warning if overflow latency degrades beyond your tolerance.

Because the ALB resides in the Region, all traffic to Outposts targets traverses the service link. Keep the following in mind:

Bandwidth planning: Steady-state traffic to Outposts targets flows over the service link. Verify that your connection meets the minimum 500 Mbps per compute rack recommended by AWS, with sufficient headroom for both application traffic and Outposts control plane communication. Monitor service link VIF throughput using IfTrafficIn and IfTrafficOut metrics (on service link VIFs) to detect saturation before it impacts performance.

Latency impact: The service link adds latency compared to a locally deployed load balancer. The exact impact depends on your service link connection type and distance to the parent Region (AWS specifies a maximum of 175 ms round-trip for service link). For internet-facing workloads, this is typically negligible relative to the client-to-Region round trip. For workloads serving on-premises users through the Local Gateway, consider Route 53 weighted routing between an ALB on Outposts and a separate ALB in the Region instead.

Connection draining: When scaling down the overflow fleet, allow sufficient time for in-flight requests to complete. The deregistration delay configured on the target group (default 300 seconds) and the Auto Scaling scale-in cool-down period work together to help provide graceful termination and minimize the risk of dropping active connections.

Failure modes: If the service link goes down, the ALB cannot reach Outposts targets. Health checks fail, and all traffic automatically shifts to Region targets. This provides an unintentional but useful failover behavior. However, note that the overflow fleet is sized for burst capacity, not for sustaining 100% of production traffic. Monitor the ConnectedStatus metric (under the AWS/Outposts namespace, dimension OutpostId) and alert on degradation. If you need full failover capability, architect a separate disaster recovery solution with appropriately sized Region capacity.

Limitations

Be aware of these constraints when implementing this pattern:

  • ALB requirement: The pattern requires an Application Load Balancer in the Region. Workloads that rely on direct IP access through the Local Gateway (without an ALB) cannot use this pattern without an architecture change.
  • Stateful workloads: Applications with local disk state or in-memory sessions require external session stores (ElastiCache, DynamoDB) before they can burst. Without this, overflow instances serve requests without session context.
  • Database coupling: If your application writes to a database running exclusively on the Outpost, overflow instances in the Region cannot reach it without a cross-location replica or proxy. Read-heavy workloads with a Region read replica are ideal candidates.
  • Service link as single path: All ALB-to-Outpost traffic shares the service link with AWS control plane operations. Under extreme load, bandwidth contention can degrade both application traffic and management operations.
  • ALB on Outposts: As of this writing, ALB on Outposts does not support weighted target groups spanning both locations. The ALB must reside in the Region for this pattern to work.

Testing the pattern

Validate the burst mechanism before relying on it in production:

Simulate capacity pressure:

aws cloudwatch set-alarm-state \
  --alarm-name outposts-capacity-high \
  --state-value ALARM \
  --state-reason "Testing burst mechanism"

Verify overflow fleet launched:

aws autoscaling describe-auto-scaling-groups \
  --auto-scaling-group-names burst-overflow-fleet \
  --query "AutoScalingGroups[0].DesiredCapacity"

Verify ALB weights shifted (after recovery check runs):

aws elbv2 describe-listeners \
  --listener-arns <your-listener-arn> \
  --query "Listeners[0].DefaultActions[0].ForwardConfig.TargetGroups[*].[TargetGroupArn,Weight]"

Trigger recovery:

aws cloudwatch set-alarm-state \
  --alarm-name outposts-capacity-high \
  --state-value OK \
  --state-reason "Testing recovery"

Confirm overflow fleet scales back to zero and all traffic returns to Outposts targets. Recovery is gradual — the Amazon EventBridge rule evaluates every 5 minutes and steps weights back before scaling down, so full recovery may take 10–15 minutes depending on your weight progression configuration.

Clean up

To avoid ongoing charges, verify that the overflow Auto Scaling group has scaled to zero, then delete the stack:

sam delete --stack-name burst-to-region-stack

This removes all resources created by the template, including the Lambda function, CloudWatch alarm, SNS topic, Amazon EventBridge rule, and the overflow Auto Scaling group.

Conclusion

This Burst to Region pattern extends AWS Outposts capacity into the parent Region during peak demand. You trade a moderate latency increase for continued availability when local capacity is exhausted.

The pattern works best when you clearly classify which workloads can overflow, implement gradual traffic transitions, and maintain security and observability parity across both environments.

For the complete deployable AWS SAM template including the Lambda orchestrator, CloudWatch dashboard, and all IAM roles, see the GitHub repository. To learn more about capacity planning for Outposts, see Managing your AWS Outposts capacity using Amazon CloudWatch and AWS Lambda and AWS Outposts monitoring and reporting: A comprehensive Amazon EventBridge solution.

For more information, see the AWS Outposts User Guide and the Amazon EC2 Auto Scaling User Guide.

Build an AI email pipeline with Amazon Bedrock and SES Mail Manager

Post Syndicated from Zip Zieper original https://aws.amazon.com/blogs/messaging-and-targeting/build-an-ai-email-pipeline-with-amazon-bedrock-and-ses-mail-manager/

Processing inbound email attachments at scale involves extracting files, routing them by recipient, scanning for malware, and classifying content. This traditionally requires stitching together polling loops, event rules, and multiple integration points. Amazon Simple Email Service (Amazon SES) Mail Manager now provides two new rule actions that simplify this pattern. The Lambda action invokes AWS Lambda functions directly from rule sets, and the Bounce action returns rejection responses. Together, they let you build multi-step email processing pipelines with declarative configuration.

In this post, you learn how to build an attachment processing pipeline that automatically extracts email attachments and classifies them with Amazon Bedrock. The pipeline also rejects infected files with RFC-compliant bounce responses. The complete implementation is available as an AWS Cloud Development Kit (AWS CDK) deployment in the companion GitHub repository sample-amazon-ses-mail-manager-attachment-pipeline. You can deploy it manually using the steps in this post, or hand it off to an AI coding agent such as Kiro or Claude Code. The repository includes a machine-readable agentic deployment guide that walks an agent through every deployment step, from prerequisite checks to post-deploy verification.

Architecture of the inbound email pipeline: SES Mail Manager routes messages through a traffic policy and rule set to AWS Lambda, Amazon Simple Storage Service (Amazon S3), Amazon DynamoDB, and Amazon Bedrock

The problem: scaling document intake for a multi-tenant platform

Consider a fictitious SaaS platform from AnyCompany that lets customers submit documents by email. Each customer sends invoices, contracts, and supporting files to a dedicated address (for example, [email protected] or [email protected]). They expect those attachments to land in their isolated storage, classified and ready for downstream processing.

Without a purpose-built pipeline, the typical approach looks like this: an Amazon S3 event notification triggers a Lambda function that polls for new MIME objects, parses them, looks up the recipient in a routing table, and fans out extraction to another function. Worse, it relies on a separate virus-scanning step having run first. Orchestration lives in AWS Step Functions or Amazon EventBridge rules. Adding a new customer means updating routing configuration in multiple places. Adding classification means bolting on yet another Lambda in the chain.

The result is fragile. When volume spikes during month-end invoice runs or onboarding waves, the polling loop backs up and retries cascade. Infected files occasionally slip past the scanner because the scan and extraction steps are not transactionally linked.

This pipeline solves the problem declaratively. Mail Manager’s traffic policy rejects unauthorized senders and enforces size limits at the SMTP connection level. This filtering happens before any processing resources are consumed. The rule set handles virus scanning, bouncing, archiving, classification, and extraction in a single ordered sequence. Each step completes before the next begins. If an attachment is infected, the sender gets an immediate SMTP bounce. There are no silent failures and no orphaned files in downstream storage.

The result is a pipeline where:

  • Adding a customer means adding an email address to the Mail Manager address list and a row in Amazon DynamoDB. No changes to code.
  • Adding a classification category means editing a prompt string. No schema migration.
  • Infected files never reach storage because the bounce fires during the SMTP transaction, before any Lambda is invoked.

Pipeline architecture overview

Table 1: Architecture components and their roles in the email processing pipeline

Component Role
Amazon SES Mail Manager open Ingress Endpoint Email arrives via public internet at a Mail Manager open ingress point over SMTP.
Mail Manager traffic policy Filters spam using the Abusix (or Spamhaus) email add-on, then enforces a recipient allowlist at the connection level.
Mail Manager rule set Messages for allowed recipients are passed to the rule set, which sequentially evaluates each message against two rules.
Rule 1 Uses the Trend Micro email add-on to scan for infected attachments, then bounces any unsafe messages back to sender (using Amazon SES outbound).
Rule 2 Clean messages passed from Rule 1 are copied to a Mail Manager archive and written as raw Multipurpose Internet Mail Extensions (MIME) objects to a “landing-zone” Amazon S3 bucket.
AWS Lambda (AttachmentProcessor) Triggered by the arrival of objects in the S3 bucket, this function parses MIME email, extracts attachments, and routes them to per-recipient S3 buckets.
AWS Lambda (EmailCategorizer) Triggered by the arrival of objects in the landing-zone S3 bucket, this function classifies each email using Amazon Nova Micro via Amazon Bedrock and writes results to Amazon DynamoDB.
Amazon S3 (landing zone + per-recipient buckets) Stores raw MIME objects in a shared landing-zone bucket; stores extracted attachments in isolated per-recipient buckets keyed by local part (for example, invoices/ for [email protected]).
Amazon DynamoDB (RecipientBucketLookup) Maps recipient email addresses to their designated S3 bucket and key prefix.
Amazon DynamoDB (EmailCategories) Stores Amazon Bedrock classification results: category, urgency, and summary.
Amazon Bedrock (Amazon Nova Micro) Classifies each email into a category (invoice, contract, HR, unknown) and urgency level.
AWS IAM roles Mail Manager and Lambda execution permissions following the principle of least privilege.

How the Mail Manager traffic policy filters connections

The traffic policy (Receive-attachments) makes connection-level decisions before any message content is processed. It evaluates two statements in order:

  1. Deny spam — Connections from senders flagged by Abusix as spam sources are denied immediately.
  2. Allow approved recipients — Connections where the recipient is in the approved-recipients address list pass through to the rule set.

The policy uses a default action of DENY, so any connection that does not match an explicit ALLOW statement is rejected. The policy also enforces a 35 MB maximum message size. You can add additional statements to enforce SPF, DKIM, or DMARC authentication results. This is useful in regulated industries where sender verification is required before any processing occurs.

The PolicyStatements array defines the evaluation order (deny first, then allow):

PolicyStatements=[
    {   # Statement 1: Deny connections from known spam sources
        "Action": "DENY",
        "Conditions": [{"BooleanExpression": {
            "Evaluate": {"Analysis": {"Analyzer": "ABUSIX_ADDON_ARN", "ResultField": "isListed"}},
            "Operator": "IS_TRUE",
        }}],
    },
    {   # Statement 2: Allow only recipients in the approved list
        "Action": "ALLOW",
        "Conditions": [{"BooleanExpression": {
            "Evaluate": {"IsInAddressList": {"Attribute": "RECIPIENT", "AddressLists": ["ADDRESS_LIST_ARN"]}},
            "Operator": "IS_TRUE",
        }}],
    },
]

For the complete create_traffic_policy call with all parameters, see the companion repository.

API reference: CreateTrafficPolicy

Rule set: the processing pipeline

Messages that pass the traffic policy enter the rule set (attachment-pipeline-rules), which evaluates two rules in order.

Rule 1 — Virus scan and bounce

This rule checks the Trend Micro add-on result. If Trend Micro reports isPassed = FALSE (infected attachment detected) — note that Mail Manager has already accepted the message by this point — the rule fires a Bounce action, which generates a non-delivery report (NDR) back to the sender with SMTP 550 (permanent failure) and status 5.7.1 (security/policy reason). It then Drops the message. No further rules run.

This after-the-fact NDR prevents infected messages from entering your processing pipeline while still providing clear guidance to legitimate senders.

Rule 2 — Process clean email

This rule has no conditions, so it applies to every message that passed the virus scan. It runs four actions in sequence:

  1. Archive — Mail Manager stores a copy in the archive for compliance and electronic discovery (eDiscovery).
  2. WriteToS3 — Mail Manager writes the raw MIME object to the amzn-s3-demo-bucket-general-receiving S3 bucket, keyed by message ID.
  3. InvokeLambda (EmailCategorizer, REQUEST_RESPONSE) — Mail Manager invokes the categorizer, which classifies the email with Amazon Bedrock and writes results to Amazon DynamoDB.
  4. InvokeLambda (AttachmentProcessor, REQUEST_RESPONSE) — Mail Manager invokes the processor, which extracts attachments and routes them to per-recipient S3 locations.

The categorizer fires before the attachment processor by design: the attachment processor deletes the original MIME from Amazon S3 after successfully extracting attachments. By running first, the categorizer is guaranteed to find the MIME in Amazon S3.

Because the Bounce and Drop actions fire in Rule 1, the Lambda functions in Rule 2 are never invoked for infected messages. There is no risk of malicious content reaching your Amazon S3 buckets or Amazon Bedrock.

API reference: CreateRuleSet

How Amazon Bedrock classifies inbound email

The MailManager-EmailCategorizer function uses Amazon Nova Micro (amazon.nova-micro-v1:0) to classify each email. Amazon Nova Micro is a fast, lightweight text-only model optimized for classification and structured output tasks. Access to all Amazon Bedrock foundation models, including Amazon Nova Micro, is available by default in all commercial AWS Regions. No access request is needed.

The function performs the following steps:

  1. Parses the recipient, message ID, and subject from the Mail Manager event.
  2. Retrieves the raw MIME from the amzn-s3-demo-bucket-general-receiving S3 bucket.
  3. Extracts the plain-text or HTML body from the MIME structure.
  4. Sends the subject (capped at 500 characters) and body (capped at 4,000 characters) to Amazon Bedrock with a classification prompt.
  5. Writes the structured result to the EmailCategories DynamoDB table.

The classification prompt returns a structured JSON response:

{
    "category": "invoice | contract | hr | unknown",
    "urgency": "urgent | non-urgent",
    "summary": "<50-word summary>"
}

If Amazon Bedrock returns an error or malformed JSON, the function falls back to category: unknown, urgency: non-urgent and continues. It never blocks the attachment processor.

Choosing a classification model

To customize the classification categories for your use case, update the SYSTEM_PROMPT in the categorizer Lambda function. The prompt uses a structured instruction format that you can extend with additional categories, urgency levels, or routing rules. For example, an insurance carrier could add categories like claim_new, claim_status, document_submission, and complaint to automatically triage patient email. You can also update the COMPANY_NAME environment variable to inject your organization’s name into the classification prompt without modifying the function code.

To switch the model, update the BEDROCK_MODEL_ID environment variable. The following table compares supported options:

Model Model ID Best for Latency Relative cost
Amazon Nova Micro amazon.nova-micro-v1:0 Fast structured classification, low latency ~200ms Lowest
Amazon Nova Lite amazon.nova-lite-v1:0 Richer summaries, multi-label classification ~400ms Moderate
Anthropic Claude 3 Haiku anthropic.claude-3-haiku-20240307-v1:0 Complex reasoning, nuanced categorization ~600ms Higher

Attachment extraction and routing

The MailManager-AttachmentProcessor function handles MIME parsing, recipient-based routing, and cleanup. It performs the following steps:

  1. Parses the recipient email address and message ID from the Mail Manager event information.
  2. Retrieves the raw MIME message from the amzn-s3-demo-bucket-general-receiving S3 bucket using the message ID from the event as the S3 key.
  3. Looks up the recipient’s S3 destination in the RecipientBucketLookup DynamoDB table, or creates a new entry if this is the first email for that recipient.
  4. Extracts attachment parts from the MIME message, skipping plain-text and HTML body parts that have no file name.
  5. Copies each attachment to the recipient’s S3 bucket at the prefix {local_part}/ (for example, invoices/ for [email protected]).
  6. Deletes the original MIME object from the landing-zone bucket, but only if every attachment copy succeeded. If any copy failed, the MIME is retained for retry.
  7. Returns a response to Mail Manager indicating success or failure.

This synchronous invocation pattern allows the rule set to make routing decisions based on the Lambda function’s response. If attachment extraction fails, subsequent rules can bounce the message or route it to a quarantine location.

Attachment detection logic

The function detects attachments using three criteria:

  1. Content-Disposition containing attachment.
  2. Any MIME part with a file name (even if disposition is inline or missing).
  3. Non-text, non-multipart parts (such as application/pdf or image/*).

For parts without a file name, the function generates one from the content type (for example, attachment.pdf).

Input validation and security

The pipeline implements the following input validation to protect against malicious content and unexpected inputs:

  • messageId validation — the messageId from the Mail Manager event is validated against an alphanumeric-plus-hyphen pattern ([a-zA-Z0-9\-]+) before use as an S3 key. Unexpected formats raise a ValueError, which causes Mail Manager to apply the ActionFailurePolicy.
  • Attachment filename sanitization — filenames from MIME Content-Disposition headers are attacker-controlled. Before use as S3 key components, each filename is processed through os.path.basename() to strip directory components, leading-dot stripping to prevent hidden-file creation, and a character allowlist ([\w.\- ]). Filenames are also truncated to 255 characters.
  • Prompt size caps — the email body sent to Amazon Bedrock is capped at 4,000 characters. The subject line is capped at 500 characters, preventing oversized prompts and excessive token usage.

The following additional controls are recommended before adapting this pipeline for production:

  • Validate attachment file types against an approved allowlist (such as .pdf, .docx, .xlsx). Reject or quarantine messages with disallowed file types.
  • Implement per-attachment size limits in addition to the overall 35 MB message size limit.
  • Verify MIME structure integrity before parsing. Handle malformed MIME structures as error conditions.
  • Log validation failures to Amazon CloudWatch for security monitoring and audit purposes.

AWS CloudFormation and CDK support for Mail Manager rule actions

The InvokeLambda and Bounce rule actions are supported natively in AWS::SES::MailManagerRuleSet as of March 2026. The companion CDK stack uses CfnMailManagerRuleSet directly. No Custom Resource is required.

When using the Python CDK L1 bindings, note that typed property classes for Bounce and InvokeLambda are not yet exposed in the Python bindings. Pass these actions as plain dicts with camelCase keys matching the AWS CloudFormation property names. RuleActionProperty accepts Dict[str, Any] for each field:

ses.CfnMailManagerRuleSet.RuleActionProperty(
    bounce={
        "smtpReplyCode": "550",
        "statusCode": "5.7.1",
        "diagnosticMessage": "Your attachment was infected.",
        "sender": "[email protected]",
        "roleArn": role.role_arn,
        "actionFailurePolicy": "CONTINUE",
    }
)

API reference: AWS::SES::MailManagerRuleSet | AWS CDK API Reference

Prerequisites

This post and companion GitHub project assume familiarity with SMTP protocols, email infrastructure concepts, AWS Lambda, Amazon S3, Amazon DynamoDB, and AWS IAM.

Estimated time: 20–30 minutes to deploy and test.

Estimated cost: This pipeline uses a Mail Manager open ingress endpoint that costs $50/mo in addition to various AWS services that are charged based on actual usage. In a low-volume test environment (fewer than 1,000 email messages per day), costs should typically be under $60 USD per month driven primarily by Mail Manager archiving, S3 storage, Lambda invocations, and Amazon Bedrock token usage. Use the AWS Pricing Calculator to estimate costs for your expected volume.

AWS IAM permissions: The deploying user needs permissions to create and manage AWS CloudFormation stacks, Lambda functions, S3 buckets, DynamoDB tables, AWS IAM roles, and Amazon SES Mail Manager resources. For testing, AdministratorAccess is sufficient. For production, scope permissions to the specific actions required: cloudformation:CreateStack, lambda:CreateFunction, s3:CreateBucket, dynamodb:CreateTable, iam:CreateRole, iam:PassRole, ses:CreateTrafficPolicy, ses:CreateRuleSet, and ses:CreateAddressList. (Separately, the Lambda functions’ own execution roles, created by the stack, grant bedrock:InvokeModel at runtime; that permission is not needed by the person deploying the stack.)

To deploy this pipeline, you need the following:

  1. An active AWS account.
  2. AWS Command Line Interface (AWS CLI) version 2.x or later installed and configured with credentials and default region.
  3. AWS CDK version 2.x or later installed (npm install -g aws-cdk) and Python 3.12 or later.
  4. Amazon SES configured with production access in the target region with a verified Amazon SES identity for the bounce sender address.
  5. Ability to administer the DNS entries for the Amazon SES identity to add an MX record pointing to the Mail Manager ingress endpoint’s A record.

Deployment

Tip: Whichever path you choose, review the Prerequisites section first to make sure your AWS account has the necessary permissions and that you have a verified domain available in Amazon SES. The complete solution is available as an open-source reference implementation. To deploy it in your AWS account, clone the companion repository:

git clone https://github.com/aws-samples/sample-amazon-ses-mail-manager-attachment-pipeline.git
cd sample-amazon-ses-mail-manager-attachment-pipeline

From here, you have two paths to get up and running:

Option 1: Deploy manually

Follow the step-by-step instructions in the repository’s README.md. At a high level, you will:

  1. Install prerequisites (AWS CDK, Node.js, Python).
  2. Configure your environment variables (AWS account, region, verified domain).
  3. Bootstrap your CDK environment.
  4. Deploy the stack with cdk deploy.
  5. Complete post-deployment verification (confirm email receiving rules are active and test with a sample message).

Option 2: Deploy with a coding agent

If you use an AI-powered coding assistant (such as Amazon Q Developer CLI or Kiro), install the AWS MCP server and SES/Mail Manager skills to empower your AI assistants with deep context on Amazon SES and Mail Manager. These resources give your assistant live access to AWS APIs and CDK documentation, which significantly reduces trial-and-error during deployment. The repository’s AGENTS.md file contains machine-readable guidance, deployment failure recovery patterns, and region handling notes specifically for AI assistants. Simply point your AI assistant at the AGENTS.md file in the repository root. This file provides structured, machine-readable instructions that guide the agent through the full deployment, from prerequisite checks through stack deployment and validation, without manual intervention.

# Example: point your agent at the instructions
@agent follow AGENTS.md

Validating the deployment

Once your stack is deployed and the MX record is in place, send a test email with an attachment to one of your approved recipient addresses. Then confirm each stage of the pipeline executed successfully:

1. Check Lambda execution

Open Amazon CloudWatch Logs for both functions and confirm they completed without errors:

aws logs tail /aws/lambda/MailManager-EmailCategorizer --follow
aws logs tail /aws/lambda/MailManager-AttachmentProcessor --follow

You should see log entries showing the message ID being processed by each function in sequence: the categorizer first, then the attachment processor.

2. Confirm email classification

Query the EmailCategories DynamoDB table to verify Amazon Bedrock classified your test message:

aws dynamodb scan --table-name EmailCategories --max-items 1

A successful record includes category, urgency, and a short summary, all generated by Amazon Nova Micro from the email’s subject and body.

3. Verify attachment extraction

Look up your recipient’s S3 destination in the RecipientBucketLookup table, then list the bucket contents to confirm the attachment arrived:

aws dynamodb get-item --table-name RecipientBucketLookup \
  --key '{"recipient": {"S": "[email protected]"}}'

aws s3 ls s3://<bucket-name>/<prefix>/ --recursive

If all three checks pass, your pipeline is fully operational. Email messages are being scanned, classified, and routed to per-recipient storage without any external orchestration.

Troubleshooting

If your test email does not flow through the pipeline as expected, start with these common issues:

Symptom Likely cause Resolution
Bounce action fails silently — infected emails are dropped without notification The bounce_sender identity is not verified in the deployment region. Amazon SES identities are regional. Verify the domain in your target region: aws sesv2 create-email-identity --email-identity example.com --region <region>, add the DKIM CNAMEs to DNS, and wait for verification. No redeployment required.
Bounce action returns a validation error bounce_sender is set to a bare domain instead of an email address Use a full address like [email protected], not just example.com

For CDK deployment issues, stack rollback errors, and teardown conflicts, see the repository troubleshooting guide.

General debugging tip: Both Lambda functions log to /aws/lambda/MailManager-EmailCategorizer and /aws/lambda/MailManager-AttachmentProcessor in Amazon CloudWatch Logs. Start there for any runtime failures.

Clean up

To avoid ongoing charges, destroy the stack when you are done:

AWS_DEFAULT_REGION= cdk destroy

Note: If the destroy fails with a ConflictException, detach the ingress point from the traffic policy first. Amazon DynamoDB tables created with RETAIN policies may also need manual deletion. See the repository’s Common failure modes table for details.

Do not forget to remove the MX record from your domain’s DNS once the ingress point is deleted. After completing the clean up, verify on the AWS Management Console that the Mail Manager ingress endpoint, Amazon S3 buckets, Amazon DynamoDB tables, and Lambda functions no longer appear in your account.

Conclusion

The Lambda action and Bounce action in Amazon SES Mail Manager support multi-step inbound email processing without complex orchestration workarounds. This pipeline demonstrates how these capabilities work together in production: scanning attachments for malware, classifying email content with AI, extracting and routing files to per-recipient storage, and providing immediate RFC-compliant feedback to senders. The modular architecture supports extension: add new classification categories, integrate additional scanning engines, or chain Lambda functions for multi-stage processing. The synchronous invocation pattern means that every processing step completes before the next begins, giving you full control over the pipeline flow. Get started by cloning the sample-amazon-ses-mail-manager-attachment-pipeline repository and deploying to your account. For an overview of the four new Mail Manager capabilities used in this pipeline, see Four new Amazon SES Mail Manager capabilities, explained.

FAQ

Q: Can I use a different Amazon Bedrock model for email classification?

Yes. Update the BEDROCK_MODEL_ID environment variable on the MailManager-EmailCategorizer Lambda function. No changes to code are required. See the preceding model comparison table for supported options.

Q: Do I need to request access to Amazon Nova Micro?

No. In all commercial AWS Regions, access to Amazon Bedrock foundation models including Amazon Nova Micro is available by default. AWS GovCloud (US) regions require an explicit access request through the Amazon Bedrock console.

Q: What happens if the Lambda function times out or fails?

A REQUEST_RESPONSE invocation is time-bounded to approximately 30 seconds, or sooner if your function’s own configured timeout is shorter. In either case, Mail Manager applies the ActionFailurePolicy configured on the rule action. If set to CONTINUE, the pipeline moves to the next action. If set to DROP, the message is discarded. This pipeline uses CONTINUE, so a transient classification failure does not block attachment delivery.

Q: Can I add more classification categories?

Yes. Edit the SYSTEM_PROMPT in the categorizer Lambda function. The function writes whatever categories the model returns to Amazon DynamoDB. No schema changes are needed.

Q: How does the pipeline handle email messages with no attachments?

The AttachmentProcessor detects zero attachment parts, skips extraction, deletes the raw MIME from the landing-zone bucket, and returns success. The EmailCategorizer still classifies the message normally.

Q: What is the maximum attachment size supported?

The traffic policy enforces a 35 MB maximum message size (total MIME payload including all attachments and base64 encoding overhead). Individual attachments are not size-limited beyond this total cap.

Q: Can I deploy this with an AI coding agent?

Yes. The repository includes an AGENTS.md file with machine-readable deployment instructions. Point your AI assistant (Kiro, Claude Code, Amazon Q Developer CLI) at this file and it handles the full deployment without manual intervention.

Q: Is the Bounce action RFC-compliant?

Yes, with one clarification: it is not a live SMTP-transaction rejection. Mail Manager first accepts the message, then the rule set runs. If the Bounce action fires, it generates a non-delivery report (NDR) back to the sender with an RFC 5321-compliant SMTP reply code and an RFC 3463-compliant enhanced status code.


About the authors

Scaling patterns for self-organizing multi-agent clusters with Kiro

Post Syndicated from Ivo Kammerath original https://aws.amazon.com/blogs/architecture/scaling-patterns-for-self-organizing-multi-agent-clusters-with-kiro/

Most multi-agent systems today follow the same shape: a supervisor agent breaks a task down, hands the pieces to subagents, and stitches the results back together. This is how Kiro CLI delegates to subagents, and what the Strands Agents SDK gives you primitives for with graphs and agents-as-tools. It is a good default. One process holds the plan, so behavior stays predictable and every result passes through a single gate.

That one process is also the limit. Every assignment and every result flows through the supervisor, so its context window caps how much work the system can hold at once. If it dies, the run dies with it. And because a single planner fixes the decomposition upfront, you get one take on the problem, multiplied by N workers.

Plenty of distributed systems still coordinate centrally, and should. But the alternative has been around for decades: let participants converge through shared state instead. We wanted to know what happens when you apply that move to artificial intelligence (AI) agents, so we built kiro-flock, an open-source reference implementation. It runs clusters of Kiro CLI agents on Amazon Elastic Compute Cloud (Amazon EC2) with nothing between them but an Amazon Simple Storage Service (Amazon S3) bucket. No orchestrator, no message bus. Agents coordinate by reading each other’s append-only logs. This post explains the pattern and gives you enough to deploy the sample and watch a cluster converge yourself.

When to use this pattern

Architecture has to match the task. In the 2025 study “Towards a Science of Scaling Agent Systems” of 260 agent-system configurations they found exactly that: task performance ran from +80.8 percent on decomposable financial
reasoning to -70.0 percent on sequential planning, against a single-agent baseline.
Neither the supervisor nor this pattern wins everywhere.

A self-organizing cluster fits work that splits into many quasi-independent contributions toward one goal: reviewing a large code base, migrating hundreds of modules against a known target, generating tests or design alternatives at scale. It also suits brainstorming where you want real variety instead of one planner’s take. Parallelism matters more than ordering. Agents can join, fail, and leave without ceremony.

A supervisor fits the opposite profile. The task tree is known upfront, steps depend on each other, or you need a verification gate before results ship. The same study found that architectures without centralized verification propagate more errors. That is a real cost of removing the arbiter, though you can still gate the finished result the way the migration example ends in a full test pass. What a cluster will not give you is a gate between every step, and if you need that, the supervisor earns its bottleneck.

Workload profile Better fit
Many independent contributions, one goal Cluster
Decomposition should emerge from the work Cluster
Diversity of approaches is an asset Cluster
Long-running, agents come and go Cluster
Known task tree, strict ordering Supervisor
Central verification gate required Supervisor
Interactive, latency-sensitive Supervisor

The pattern

The core decision: coordination lives in shared state. No component plans, assigns, or aggregates for the rest. Three parts make it work:

  • Agents. Independent processes that read and write a shared store and never connect to each other. One agent failing stops only its own log.
  • A shared environment. A single store holds a direction file, one append-only log per agent, and a working area for artifacts.
  • A direction. A markdown file that states the goal and leaves the path to the agents.

Each agent runs a loop. It starts a fresh session, reads the direction and the logs of a bounded set of peers, decides on one contribution that moves the goal forward, writes its artifacts, and appends one line to its own log:

{"ts":"2026-07-21T14:12:42Z","iteration":0,"action":"wrote discussion on coordination topologies","result":"Created discussion-coordination-topologies.md covering ring vs mesh vs swarm trade-offs with analysis of convergence/diversity tension.","next_intent":"read neighbour contributions and either deepen topology discussion or explore a second angle"}

That line is the entire coordination message. No broker delivers it, no acknowledgment comes back. The next agent that reads it decides for itself what to do about it. Remove an agent and its neighbors read one fewer log. Add one mid-run and it joins the division of labor already underway.

The bounded peer set is deliberate. Agents sit in a logical ring and each reads a fixed number of neighbors on either side, set by a radius parameter. Give every agent full visibility and the cluster collapses onto whatever the first agent wrote, because each later agent reads that as consensus. Limited visibility lets signals spread gradually, and agents working from different context get room to develop alternatives.

In kiro-flock, each agent is a headless Kiro CLI session on its own Amazon EC2 instance, and the shared environment is an Amazon S3 bucket. Which tools an agent may use without review, and what each iteration reads and writes, are design decisions you make once per cluster. We think of them as harness engineering and loop engineering, and the drift failure mode in the following section shows why the fresh session per iteration matters.

What a run looks like

Reference architecture: EC2 agent instances, an S3 shared environment holding the direction, logs, and artifacts, an API Gateway and Lambda control plane behind Amazon Cognito, CloudWatch metrics, and Amazon Bedrock post-run analysis.

Figure 1. Reference architecture for a kiro-flock cluster on AWS.

Agents run as headless Kiro CLI sessions on EC2 instances, each reading and writing the S3 bucket that holds the direction, one log per agent, and the shared artifacts. An Amazon API Gateway and AWS Lambda control plane behind Amazon Cognito starts, stops, and steers clusters from the dashboard. Agents publish metrics to Amazon CloudWatch, and Amazon Bedrock backs the post-run analysis.

Take a concrete run: sixteen agents in a ring, directed to hold a distributed discussion on AI agent clustering and converge on a shared synthesis. The operator writes one direction file and starts the cluster. Nothing else is assigned.

The following lines are from that run (result fields shortened for print). In the first iteration the agents fanned out with no assignment: failure modes, coordination topologies, distributed-systems parallels, and several overlapping stigmergy pieces, all written in parallel within one minute of start. The later lines show an agent correcting course after reading its neighbors, the synthesis forming, and the cluster winding itself down:

{"ts":"2026-07-21T14:12:42Z","iteration":0,"action":"wrote discussion-failure-modes.md","result":"covered 5 failure modes (premature convergence, echo chambers, hot-spot collapse, idle cascades, stale reads) with mitigations","next_intent":"read neighbours next iteration, look for topics to challenge or extend"}
{"ts":"2026-07-21T14:14:31Z","iteration":1,"action":"wrote discussion-scaling-laws.md","result":"explored scaling regimes and communication cost vs value; diverged from neighbours who all wrote stigmergy in iteration 0","next_intent":"read neighbour updates, check if any uncovered angles remain"}
{"ts":"2026-07-21T14:16:04Z","iteration":2,"action":"wrote synthesis.md","result":"structured synthesis: 5 core claims, 3 key tensions, 5 open questions; integrated topologies, scaling, economics, failure modes","next_intent":"read neighbours' reactions to synthesis, refine, move toward idle if coverage is solid"}
{"ts":"2026-07-21T14:22:24Z","iteration":7,"action":"idle","result":"All 6 visible neighbours idle (iterations 4-6). Synthesis.md stable and comprehensive. Direction unchanged. Cluster fully converged.","next_intent":"remain idle unless direction changes or operator requests further work"}

The cluster converged on a shared synthesis covering the angles in the direction, and by iteration 7 all sixteen agents had declared themselves idle. Nobody assigned the topics, arbitrated the synthesis, or told the cluster it was done. Even done is only a signal read from the logs, since an agent goes idle when its neighbors are idle and the output is stable.

kiro-flock dashboard showing a single six-agent cluster running the amorphous algorithm at radius 1, each agent card listing its neighbors, health, and iteration log, with the shared environment and direction on the right.

Figure 2. A single cluster in the kiro-flock dashboard. Six agents run at radius 1, each on its own EC2 instance. Every agent card shows its neighbors and its latest log line, the “did / result / next intent” message its neighbors read. The right panel shows the shared environment in S3 and the direction the cluster works toward.

Three ways to answer “whose work do I read?”

Every iteration starts with that question, and the answer defines the coordination algorithm. kiro-flock ships three, swappable at runtime.

Amorphous (ring). Each agent reads a fixed window of neighbors set by radius R. An agent at radius 2 reads four neighbors whether the cluster holds 8 agents or 800, so per-agent work stays constant as the cluster grows. The ceiling is your EC2 vCPU quota, not the algorithm. The largest system we have run so far totaled 184 agents across 11 cooperating clusters, creating a programming language. Rings beyond the low hundreds are extrapolation from that constant per-agent cost, not tested territory. The price is speed: a signal moves one hop per iteration. That slowness is also what lets dissenting agents mature alternatives before the neighborhood locks in. Use it for parallel work, or as the opening phase before consensus.

Mesh (full visibility). Every agent reads every other agent’s latest entry. Alignment is fast and context grows linearly with the cluster, so mesh stays comfortable to about 30 agents and workable to about 50. Diversity collapses, because agents reacting to the same first signal agree instead of exploring. Use it when a small group must converge quickly.

Swarm (recency). Each agent reads the K most recently active peers, so the cluster reorganizes around where the action is. Good for ideation, runs well past 100 agents. If K stays small while N grows, most agents read the same few peers and pile onto one subtask. Raise K or switch to amorphous.

A productive sequence uses all three: open amorphous to explore, switch to swarm as a direction forms, finish in mesh to align on the output.

How long does convergence take? In a ring, one iteration carries a signal 2R positions, so full propagation takes ceil(N / 2R) iterations, and consensus roughly two to three times that, because agents observe, react, and confirm. The wall-clock column assumes an iteration interval of 30 seconds per agent loop, the default interval in the reference implementation. The interval is configurable per cluster.

Agents (N) Radius (R) Propagation ceil(N/2R) Consensus (2-3x) Wall clock to propagate
8 1 4 iterations 8-12 iterations about 2 minutes
100 2 25 iterations 50-75 iterations about 12 minutes
1,000 4 125 iterations 250-375 iterations about 62 minutes
1,000 20 25 iterations 50-75 iterations about 12 minutes

Radius trades per-agent context for convergence speed, as the last two rows show. For parallel map-style work, propagation barely matters. Agents only need to avoid duplicating each other. It costs you when the task needs consensus, so match radius and cluster size to the context you are in. The cost model follows the same logic: no always-on orchestrator and no broker. You pay for Kiro credits and the EC2 instances while they run, plus S3 storage and requests. You also pay for the AWS Lambda, Amazon API Gateway, and Amazon Bedrock usage the control plane and post-run analysis incur. See AWS Pricing.

None of this is new theory. Identical unreliable parts producing coherent global behavior through local reads is amorphous computing. Coordinating through traces left in a shared medium instead of messages is stigmergy, described by Grassé for termites in 1959 and formalized for artificial systems by Theraulaz and Bonabeau. And a set of append-only logs is a grow-only conflict-free replicated data type (CRDT) spreading gossip-style: replicas converge without locks, which is all the consistency this workload needs.

Where it breaks

Self-organizing clusters fail in ways orchestrated systems do not. With no supervisor to arbitrate, a bad signal can spread before anyone corrects it. Four failure modes recur, and each maps to a design choice rather than a safeguard bolted on afterward.

Failure mode Where it comes from Design choice that addresses it
Groupthink Mesh visibility collapses the cluster onto the first signal Open amorphous to build diversity, switch to mesh only to align
Drift Persistent session history builds behavioral momentum Fresh session per iteration. State lives only in shared logs
Hot spots Swarm with K too small for N starves subtasks Raise K, or switch to amorphous
Carry-over Stale files from a previous run read as current context Archive environment/ and store/ to history/ on every start

Drift deserves one more sentence, because it is the least obvious. An agent that keeps its session history carries a narrow reading of the direction forward even after its neighbors move on. Starting every iteration with no conversational memory sounds wasteful. It is actually the control that keeps a thousand independent loops steerable. The only state an agent carries is what it reads back from the shared logs.

Composing clusters

The same decision works one level up: clusters coordinate by reading each other’s shared environment, the way agents read each other’s logs. We run a structure we call WeltenBuilder: a feature cluster implements against an agreed interface, a shared-infrastructure cluster owns common services, a QA cluster reads across the others and reports inconsistencies as artifacts. A coordinator cluster writes conflict-resolution notes the others pick up on their next iteration. A resolution note is a trace, not a command. Remove the coordinator and you remove a signal, not a dependency.

Because the shared environment is the coordination plane, all clusters launch at the same time with no dependency graph to sequence. Contract bottlenecks dissolve the same way: a small mesh cluster converges on interface definitions in a few iterations while other clusters build against its latest stable output. This is where the pattern points: standing clusters, each producing one class of artifact, composed into a factory whose unit of work is a direction file and a topology.

kiro-flock WeltenBuilder dashboard showing several specialized clusters running at once, each with its own algorithm and agent count, beside the shared environment tree on the right.

Figure 3. Multiple specialized clusters in the WeltenBuilder dashboard, each with its own algorithm and agent count, coordinating only through the shared S3 environment on the right.

Try it

The kiro-flock reference implementation is open source under Apache 2.0. It is a sample to study and adapt, not a production system.

One setup script provisions the stack with the AWS Cloud Development Kit (AWS CDK): Amazon S3 for the shared environment, Amazon EC2 for the agents, and AWS Lambda with Amazon API Gateway as a control plane behind a dashboard. The dashboard starts and stops clusters, changes the algorithm, and updates the direction mid-run. Amazon Cognito handles access, and Amazon Bedrock backs a post-run analysis that summarizes how the cluster converged.

You need an AWS account with the AWS CDK bootstrapped. Install kiro-cli and create a Kiro API key for headless mode (requires a Kiro subscription). Then:

cp install.config.template install.config   # set REGION and PROFILE
./setup.sh

Direct a cluster in plain language: “Start a flock of 8 agents to review the files in my project and suggest improvements.” The default runs 8 agents at radius 1 and converged in 5 to 7 iterations in our runs. Before wider use, scope each agent’s EC2 AWS Identity and Access Management (IAM) role, restrict security-group egress to the endpoints agents should call, and add AWS Budgets alerts.

Conclusion

The supervisor pattern remains the right default for bounded task trees, whether you build it with Strands, Kiro CLI subagents, or any of the coding agents that delegate this way. When the work decomposes into many independent contributions and diversity matters more than a central gate, moving coordination into shared state helps remove the throughput ceiling and the single point of failure in one move. The convergence math and the failure modes both follow from that decision, and the distributed-systems results they rest on have been known for decades. Deploy the sample, read the logs as a cluster converges, and decide where your own multi-agent workloads belong.


About the authors

Centralized CloudTrail monitoring across 100+ AWS accounts

Post Syndicated from Jagdish Komakula original https://aws.amazon.com/blogs/big-data/centralized-cloudtrail-monitoring-across-100-aws-accounts/

Organizations running workloads across dozens or hundreds of AWS accounts face a common challenge: centralized security monitoring at scale. Security teams need to search through hundreds of gigabytes of AWS CloudTrail logs daily to detect threats and satisfy compliance requirements for SOC 2, PCI DSS, and HIPAA audits. They also need to provide role-based access to multiple teams with different responsibilities.

Without purpose-built infrastructure, this often involves manual log searching that takes hours and compliance report generation that takes days. It also leads to fragmented code bases of custom AWS Lambda functions managing index lifecycles across environments. A single shared search domain without consistent access control compounds the problem further.

In this post, we show you how to build a centralized CloudTrail monitoring solution on Amazon OpenSearch Service. Terraform manages the full stack, from domain provisioning to access control and lifecycle policies. The solution handles 200 GB/day of CloudTrail logs, provides automated threat detection alerts, and gives 4 different teams isolated, role-appropriate access to the data.

Solution overview

The following diagram shows the architecture. CloudTrail logs flow from 100+ AWS accounts through an organization trail into a centralized S3 bucket. Amazon Simple Queue Service (Amazon SQS) notifications trigger the OpenSearch Ingestion pipeline. The pipeline auto-scales between 2 and 10 OpenSearch Compute Units (OCUs) to parse and index the logs into the OpenSearch domain. Four team-specific roles access the data through OpenSearch Dashboards with tenant isolation.

CloudTrail logs flow from 100+ accounts into S3, then through Amazon SQS and OpenSearch Ingestion into the OpenSearch domain used by 4 team roles

The key components are:

  • CloudTrail aggregation. An organization trail sends logs from 100+ accounts into a centralized Amazon Simple Storage Service (Amazon S3) bucket.
  • Ingestion. An Amazon OpenSearch Ingestion pipeline picks up new logs through Amazon SQS notifications on the S3 bucket. It automatically scales between 2 and 10 OCUs based on queue depth. Throttling on lower-environment queues prevents development and test spikes from starving production ingestion.
  • Amazon OpenSearch Service domain. 6 or1.4xlarge data nodes (OpenSearch Optimized instances) with 3 dedicated r8g.large master nodes, fine-grained access control, encryption at rest, and node-to-node encryption.
  • Infrastructure as code. Index templates, Index State Management (ISM) policies, roles, role mappings, tenants, alerting monitors, and dashboards are all declared in Terraform and applied consistently across environments.

Prerequisites

To implement this solution, you need the following:

  • An organization in AWS Organizations with CloudTrail enabled across member accounts.
  • Terraform v1.5+ with the AWS provider and the OpenSearch provider.
  • A virtual private cloud (VPC) with private subnets for the OpenSearch domain.
  • IAM roles for each team that will access the OpenSearch domain.
  • An Amazon Simple Notification Service (Amazon SNS) topic for security alert notifications.
  • Familiarity with Amazon OpenSearch Service, Terraform, and AWS CloudTrail.
  • Sample Terraform code is available in the GitHub repository

Implementation

This section walks through the Terraform code for each component of the solution, starting with the workload profile that informed our sizing decisions.

Workload profile

Before sizing the cluster, we defined the workload characteristics and SLAs for the centralized CloudTrail monitoring platform:

Metric Value
Index throughput 200 GB/day (~18,000 docs/sec)
Search queries ~2,000 queries/day (~0.023 QPS)
Average search latency < 100 ms (achieved: 76 ms)
Saved searches 600+
Dashboards and visualizations 100+
User teams 4 (Security Ops, Incident Response, Compliance, DevOps)
Retention 30 days (hot tier)
Availability target 99.9%

This is a write-heavy ingestion workload. The primary use case is automated alerting and periodic compliance queries rather than continuous interactive search. This workload profile informed the decision to use OR1 (storage-optimized) instances with zero replicas, prioritizing indexing throughput over search parallelism.

Domain provisioning

Start by provisioning the Amazon OpenSearch Service domain with encryption, fine-grained access control, and VPC placement:

resource "aws_opensearch_domain" "cloudtrail" {
  domain_name    = var.domain_name
  engine_version = "OpenSearch_3.3"

  cluster_config {
    instance_type          = "or1.4xlarge.search"
    instance_count         = 6
    zone_awareness_enabled = true
    zone_awareness_config {
      availability_zone_count = 3
    }
  }

  dedicated_master_config {
    dedicated_master_enabled = true
    dedicated_master_type    = "r6g.large.search"
    dedicated_master_count   = 3
  }

  ebs_options {
    ebs_enabled = true
    volume_type = "gp3"
    volume_size = 500
    iops        = 3000
    throughput  = 125
  }

  encrypt_at_rest { enabled = true }
  node_to_node_encryption { enabled = true }

  domain_endpoint_options {
    enforce_https       = true
    tls_security_policy = "Policy-Min-TLS-1-2-PFS-2023-10"
  }

  advanced_security_options {
    enabled                        = true
    internal_user_database_enabled = false
    master_user_options {
      master_user_arn = var.master_user_arn
    }
  }

  vpc_options {
    subnet_ids         = var.vpc_subnet_ids
    security_group_ids = var.vpc_security_group_ids
  }

  tags = {
    Environment = "production"
    Project     = "centralized-cloudtrail-monitoring"
    ManagedBy   = "terraform"
  }
}

This solution was built on OR1 instances, which are storage-optimized and use Amazon Elastic Block Store (Amazon EBS) (gp3 or io1) for local storage, with data copied synchronously to Amazon S3 as it arrives. This storage structure provides increased indexing throughput because indexing is performed exclusively on primary shards. Replicas are backed by Amazon S3 through segment replication, eliminating the CPU overhead of document replication on replica nodes. For new deployments, we recommend OR2 instances, which offer up to 26% higher indexing throughput compared to OR1 while maintaining the same storage-optimized architecture.

Ingestion pipeline

The Amazon OpenSearch Ingestion pipeline provides serverless, auto scaling ingestion from Amazon S3 into the OpenSearch domain. It picks up new CloudTrail logs through Amazon SQS notifications on the centralized S3 bucket and scales between 2 and 10 OpenSearch Compute Units (OCUs) based on queue depth:

resource "aws_iam_role" "osis_pipeline" {
  name = "cloudtrail-osis-pipeline-role"
  assume_role_policy = jsonencode({
    Version = "2012-10-17"
    Statement = [{
      Action    = "sts:AssumeRole"
      Effect    = "Allow"
      Principal = { Service = "osis-pipelines.amazonaws.com" }
    }]
  })
}

resource "aws_iam_policy" "osis_pipeline" {
  name = "cloudtrail-osis-pipeline-policy"
  policy = jsonencode({
    Version = "2012-10-17"
    Statement = [
      {
        Action   = ["s3:GetObject", "s3:ListBucket"]
        Effect   = "Allow"
        Resource = [var.cloudtrail_bucket_arn, "${var.cloudtrail_bucket_arn}/*"]
      },
      {
        Action   = ["sqs:ReceiveMessage", "sqs:DeleteMessage", "sqs:GetQueueAttributes"]
        Effect   = "Allow"
        Resource = var.cloudtrail_sqs_queue_arn
      },
      {
        Action   = ["es:DescribeDomain", "es:ESHttp*"]
        Effect   = "Allow"
        Resource = "${aws_opensearch_domain.cloudtrail.arn}/*"
      }
    ]
  })
}

resource "aws_iam_role_policy_attachment" "osis_pipeline" {
  role       = aws_iam_role.osis_pipeline.name
  policy_arn = aws_iam_policy.osis_pipeline.arn
}

resource "aws_cloudwatch_log_group" "osis_pipeline" {
  name              = "/aws/vendedlogs/OpenSearchIngestion/cloudtrail-pipeline"
  retention_in_days = 30
}

resource "aws_osis_pipeline" "cloudtrail" {
  pipeline_name = "cloudtrail-ingestion"
  pipeline_configuration_body = <<-EOT
    version: "2"
    cloudtrail-pipeline:
      source:
        s3:
          notification_type: "sqs"
          codec:
            json:
          compression: "gzip"
          sqs:
            queue_url: "${var.cloudtrail_sqs_queue_url}"
          aws:
            sts_role_arn: "${aws_iam_role.osis_pipeline.arn}"
            region: "${data.aws_region.current.name}"
      processor:
        - date:
            from_time_received: true
            destination: "@timestamp"
      sink:
        - opensearch:
            hosts: ["https://${aws_opensearch_domain.cloudtrail.endpoint}"]
            index: "cloudtrail-%{yyyy.MM.dd}"
            aws:
              sts_role_arn: "${aws_iam_role.osis_pipeline.arn}"
              region: "${data.aws_region.current.name}"
  EOT
  min_units = 2
  max_units = 10
  log_publishing_options {
    is_logging_enabled = true
    cloudwatch_log_destination {
      log_group = aws_cloudwatch_log_group.osis_pipeline.name
    }
  }
  tags = {
    Environment = "production"
    Project     = "centralized-cloudtrail-monitoring"
    ManagedBy   = "terraform"
  }
}

The pipeline uses the S3 source plugin with SQS-based notifications. When new CloudTrail log files land in S3, an SQS message triggers the pipeline to fetch and parse them. The min_units and max_units parameters control auto scaling. The pipeline starts at 2 OCUs and scales up to 10 based on queue depth, handling ingestion spikes without manual intervention. For lower environments (development and testing), you can apply throttling on the Amazon SQS queue to prevent non-production spikes from starving production ingestion capacity.

Index template

Define index templates up front to avoid painful reindexing later. The following template sets explicit mappings for CloudTrail fields, optimizes for write throughput with async translog durability, and integrates with ISM for automatic rollover:

resource "opensearch_index_template" "cloudtrail" {
  name = "cloudtrail-template"
  body = jsonencode({
    index_patterns = ["cloudtrail-*"]
    priority       = 100
    template = {
      settings = {
        number_of_shards                                  = 6
        number_of_replicas                                = 0
        "index.refresh_interval"                          = "10s"
        "index.translog.durability"                       = "async"
        "index.translog.sync_interval"                    = "30s"
        "plugins.index_state_management.rollover_alias"   = "cloudtrail"
      }
      mappings = {
        properties = {
          "@timestamp"        = { type = "date" }
          eventSource         = { type = "keyword" }
          eventName           = { type = "keyword" }
          awsRegion           = { type = "keyword" }
          sourceIPAddress     = { type = "ip" }
          errorCode           = { type = "keyword" }
          errorMessage        = { type = "text" }
          recipientAccountId  = { type = "keyword" }
          userIdentity = {
            properties = {
              type      = { type = "keyword" }
              arn       = { type = "keyword" }
              accountId = { type = "keyword" }
              userName  = { type = "keyword" }
              sessionContext = {
                properties = {
                  sessionIssuer = {
                    properties = {
                      type     = { type = "keyword" }
                      arn      = { type = "keyword" }
                      userName = { type = "keyword" }
                    }
                  }
                }
              }
            }
          }
          requestParameters = { type = "object", enabled = true }
          responseElements  = { type = "object", enabled = true }
        }
      }
    }
  })
}

Defining mappings before ingestion prevents mapping conflicts and avoids the need to reindex data after the fact.

Lifecycle management (ISM policy)

The following ISM policy replaces custom Lambda functions with a single declarative policy. The rollover action uses two OR conditions: min_index_age and min_primary_shard_size. Whichever threshold is reached first triggers the rollover. This keeps shard sizes bounded while ensuring timely rotation even during low-volume periods:

resource "opensearch_ism_policy" "cloudtrail_lifecycle" {
  policy_id = "cloudtrail-lifecycle"
  body = jsonencode({
    policy = {
      description   = "CloudTrail lifecycle - rollover, retain 30d, delete"
      default_state = "hot"
      ism_template  = [{ index_patterns = ["cloudtrail-*"], priority = 100 }]
      states = [
        {
          name    = "hot"
          actions = [{ rollover = { min_primary_shard_size = "30gb", min_index_age = "1d" } }]
          transitions = [{ state_name = "delete", conditions = { min_index_age = "30d" } }]
        },
        {
          name        = "delete"
          actions     = [{ delete = {} }]
          transitions = []
        }
      ]
    }
  })
}

This approach reduces lifecycle management code by approximately 60% compared to per-environment Lambda functions, and changes deploy in a single terraform apply.

Note: For indexes ingesting more than 100 GB/day (such as CloudTrail at 200 GB/day in this deployment), you can override min_index_age to 12h to roll over more frequently. The two conditions are OR-based in OpenSearch ISM. If a shard reaches 30 GB before 1 day, it rolls over on size. If 1 day passes before 30 GB, it rolls over on age.

Multi-team access control

When multiple teams need different access levels to the same data, define all roles declaratively and use for_each to create them consistently. The following example defines 4 team roles with varying permissions:

locals {
  team_roles = {
    security_ops = {
      description         = "Security Operations - full read, alert management"
      cluster_permissions = ["cluster_monitor", "cluster:admin/opendistro/alerting/*"]
      index_permissions = [
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search", "get"] },
        { index_patterns = [".opendistro-alerting-*"], allowed_actions = ["read", "write", "search", "get", "delete"] }
      ]
    }
    incident_response = {
      description         = "Incident Response - full read for investigation"
      cluster_permissions = ["cluster_monitor"]
      index_permissions = [
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search", "get"] }
      ]
    }
    compliance_auditors = {
      description         = "Compliance - read-only"
      cluster_permissions = []
      index_permissions = [
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search"] }
      ]
    }
    devops = {
      description         = "DevOps - infra metrics and limited CloudTrail"
      cluster_permissions = ["cluster_monitor"]
      index_permissions = [
        { index_patterns = ["infra-metrics-*"], allowed_actions = ["read", "search", "get"] },
        { index_patterns = ["cloudtrail-*"], allowed_actions = ["read", "search"] }
      ]
    }
  }
}

resource "opensearch_role" "teams" {
  for_each            = local.team_roles
  role_name           = each.key
  description         = each.value.description
  cluster_permissions = each.value.cluster_permissions

  dynamic "index_permissions" {
    for_each = each.value.index_permissions
    content {
      index_patterns  = index_permissions.value.index_patterns
      allowed_actions = index_permissions.value.allowed_actions
    }
  }

  dynamic "tenant_permissions" {
    for_each = [each.key]
    content {
      tenant_patterns = [each.key]
      allowed_actions = ["kibana_all_write"]
    }
  }
}

resource "opensearch_roles_mapping" "teams" {
  for_each      = local.team_roles
  role_name     = opensearch_role.teams[each.key].role_name
  backend_roles = var.team_iam_roles[each.key]
}

resource "opensearch_tenant" "teams" {
  for_each    = local.team_roles
  tenant_name = each.key
  description = "Dashboard workspace for ${replace(each.key, "_", " ")}"
}

This approach maps IAM roles (not individual users) to OpenSearch roles. Adding a new team means adding one entry to the locals block and running terraform apply. Each team gets an isolated tenant in OpenSearch Dashboards, preventing cross-team interference with saved searches, visualizations, and dashboard configurations. We chose OpenSearch Dashboards because tenants, roles, visualizations, and saved objects can all be managed programmatically through the Terraform OpenSearch provider, keeping the entire stack under infrastructure-as-code governance. For teams building new visualizations outside of Terraform-managed workflows, we recommend OpenSearch UI. This next-generation analytics interface supports multiple data sources, provides workspaces for team isolation, and remains available during cluster upgrades.

Alerting

Define alerting monitors in Terraform to detect security-critical events automatically. The following monitor catches CloudTrail tampering attempts (StopLogging, DeleteTrail) and sends alerts through Amazon SNS:

resource "opensearch_monitor" "cloudtrail_tampering" {
  body = jsonencode({
    name     = "CloudTrail Tampering Detection"
    type     = "monitor"
    enabled  = true
    schedule = { period = { interval = 1, unit = "MINUTES" } }
    inputs = [{
      search = {
        indices = ["cloudtrail-*"]
        query = {
          size = 5
          query = {
            bool = {
              must = [{ terms = { eventName = ["StopLogging", "DeleteTrail",
                "UpdateTrail", "PutEventSelectors", "DeleteEventDataStore"] } }]
              filter = [{ range = { "@timestamp" = { gte = "now-1m" } } }]
            }
          }
        }
      }
    }]
    triggers = [{
      name     = "trail_tampering_detected"
      severity = "1"
      condition = { script = {
        source = "ctx.results[0].hits.total.value > 0"
        lang   = "painless"
      } }
      actions = [{
        name             = "notify_security"
        destination_id   = var.sns_destination_id
        message_template = { source = "CRITICAL: CloudTrail tampering detected." }
      }]
    }]
  })
}

Results and performance

After deploying the solution, we measured steady-state performance against the SLAs defined in the workload profile:

Metric Target Achieved
Index throughput 200 GB/day 200 GB/day sustained (~18,000 docs/sec)
Search latency (avg) < 100 ms 76 ms
Search availability 99.9% 99.95%+ (no unplanned downtime in 145 days)
Alert detection time < 2 minutes ~1 minute (monitor interval)
Compliance report generation < 5 minutes On-demand via saved searches

Key outcomes:

  • Threat detection dropped from hours to minutes. Automated alerting replaced manual log searching. The CloudTrail tampering monitor detects suspicious activity within one minute of the event.
  • Compliance reports generate on demand. With 600+ saved searches and 100+ dashboards, compliance teams produce SOC 2, PCI DSS, and HIPAA audit evidence in minutes rather than days.
  • Four teams operate independently. Each team has its own isolated tenant in OpenSearch Dashboards, preventing cross-team interference with saved searches and dashboard configurations.
  • Zero custom Lambda functions. ISM policies, index templates, and access control are all managed declaratively through Terraform, eliminating the previous fragmented code base.

Best practices

  • Define index templates before ingesting anything. Changing mappings on existing indices means reindexing. Get this right first.
  • Set rollover thresholds based on your actual ingestion rate. At 200 GB/day, rolling over at 30 GB keeps shard counts manageable while balancing query performance.
  • Test ISM transitions in a lower environment first. Warm and cold migrations on large indices take time.
  • Map IAM roles, not users. People change teams. Roles stay stable. This simplifies access management.
  • Put an Amazon SQS queue between S3 and the ingestion pipeline. This gives you per-environment throttling control without modifying pipeline configuration.
  • Use for_each aggressively. Roles, tenants, index patterns, and monitors all follow a pattern across teams or environments, so use for_each to eliminate copy-paste drift.
  • Consider OpenSearch UI for new visualization workflows. This solution uses OpenSearch Dashboards for Terraform-managed tenants and roles. OpenSearch UI is a next-generation interface that supports multiple data sources, stays available during cluster upgrades, and includes workspaces for team isolation. It is the recommended interface for creating new dashboards and visualizations going forward.

Optional: Extending with cold storage for longer retention

For organizations with compliance requirements mandating longer retention (for example, 7 years for PCI DSS or HIPAA), you can extend the ISM policy with warm and cold tiers. The following example adds tiered storage that moves data through hot, warm, cold, and delete states:

states = [
  {
    name    = "hot"
    actions = [{ rollover = { min_primary_shard_size = "30gb", min_index_age = "1d" } }]
    transitions = [{ state_name = "warm", conditions = { min_index_age = "30d" } }]
  },
  {
    name = "warm"
    actions = [
      { warm_migration = {} },
      { force_merge = { max_num_segments = 1 } }
    ]
    transitions = [{ state_name = "cold", conditions = { min_index_age = "365d" } }]
  },
  {
    name    = "cold"
    actions = [{ cold_migration = { timestamp_field = "@timestamp" } }]
    transitions = [{ state_name = "delete", conditions = { min_index_age = "2555d" } }]
  },
  {
    name        = "delete"
    actions     = [{ cold_delete = {} }]
    transitions = []
  }
]

Warm storage uses force-merge to reduce segment count (lowering query overhead), while cold storage moves data entirely to Amazon S3 for minimal cost. This tiered approach keeps hot-tier performance high while meeting long-term audit requirements.

Cleanup

To avoid incurring ongoing charges, remove the resources created in this post by running:

terraform destroy

This removes the OpenSearch domain, ingestion pipeline, IAM roles, SQS queues, and all associated configurations. Verify that you have exported any data or dashboards you want to retain before running destroy.

Conclusion

In this post, we showed you how to build a centralized CloudTrail monitoring solution on Amazon OpenSearch Service with Terraform managing the entire stack. The approach moves from fragile, manually configured systems with redundant Lambda code to a version-controlled, peer-reviewed, consistently deployed infrastructure.

Threat detection drops from hours to minutes with automated alerting. Compliance reports that took days now generate on demand. And your team spends time on security analysis instead of infrastructure maintenance.

To get started, use the AWS Terraform provider aws_opensearch_domain resource for the domain, then use the Terraform OpenSearch provider for index templates, ISM policies, roles, and monitors. Configure your ingestion pipeline to transform and enrich incoming CloudTrail logs before indexing, building a modern, scalable security foundation that grows with your organization.

The complete source code for this solution is available in the GitHub repository: GitHub repository

For more on the services used in this solution:


About the authors

Jagdish Komakula

Jagdish Komakula

Jagdish is a Senior Delivery Consultant at AWS Professional Services, focused on Amazon OpenSearch Service and Infrastructure automation. He has spent the last several years guiding financial services customers through building data platforms that scale.

Aditya Ambati

Aditya Ambati

Aditya is a Delivery Consultant at AWS Professional Services, focused on DevOps and infrastructure as code. He works with customers on automating cloud operations and implementing GitOps practices.

AI-powered cost optimization agent for Amazon Kinesis Data Streams

Post Syndicated from Masudur Rahaman Sayem original https://aws.amazon.com/blogs/big-data/ai-powered-cost-optimization-agent-for-amazon-kinesis-data-streams/

Customers running multiple Amazon Kinesis Data Streams often struggle to estimate the cost impact of switching between capacity modes. As accounts grow to tens or hundreds of streams, manually reviewing Amazon CloudWatch metrics for each stream and comparing pricing across Provisioned, On-demand Standard, and On-demand Advantage becomes impractical. Teams often stay on their current mode, unsure whether switching would save money or cost more, leaving potential savings unquantified. On-demand Advantage is an account-level setting that unlocks additional capabilities and a different pricing structure for on-demand streams in an AWS Region. However, without a clear, data-driven comparison, the decision to enable it remains difficult to justify.

In this post, we show you how to deploy an AI-powered agent built on Amazon Bedrock. The agent automatically analyzes every Kinesis Data Stream in your account and compares costs across all three capacity modes. It tells you exactly which streams to move to On-demand and whether your account qualifies for On-demand Advantage pricing, all on a daily or weekly schedule with zero manual intervention. As we showed in Kinesis On-demand Advantage saves 60%+ on streaming costs, choosing the right mode can save over 60 percent on streaming costs. This agent automates that analysis for you.

What is the Kinesis Mode Optimizer Agent?

The Kinesis Mode Optimizer Agent is an open source, serverless solution that uses Amazon Bedrock AgentCore, a platform to build, connect, and optimize agents at scale, with any framework or model. The agent autonomously analyzes your Kinesis Data Streams usage. It collects 7 days of Amazon CloudWatch metrics for every stream in the Region and discovers Enhanced Fan-Out (EFO) consumers. It then computes a three-way cost comparison (On-demand Standard, On-demand Advantage, Provisioned) and generates per-stream recommendations along with an account-level On-demand Advantage assessment.

The agent strongly prefers on-demand modes for their operational simplicity (automatic scaling, no capacity planning, and no throttling risk).

Results are stored as both a visual HTML report and machine-readable JSON in Amazon Simple Storage Service (Amazon S3).

Architecture

The solution uses the following architecture:

Architecture diagram of the Kinesis Mode Optimizer Agent showing Amazon EventBridge, a Scheduler Lambda, Amazon Bedrock AgentCore, a Tool Lambda, and downstream Kinesis Data Streams, CloudWatch, and Amazon S3

Figure 1: Architecture of the Kinesis Mode Optimizer Agent

The architecture flow includes the following steps:

  1. Amazon EventBridge Schedule triggers the Scheduler Lambda, an AWS Lambda function, on your configured cadence (daily, weekly, or custom cron).
  2. Scheduler Lambda invokes the Amazon Bedrock AgentCore harness with the instruction to analyze streams and generate a report.
  3. Amazon Bedrock AgentCore Gateway, a capability of Amazon Bedrock AgentCore powered by Claude Sonnet, interprets the request and routes it to the appropriate Model Context Protocol (MCP) tools exposed by the Tool Lambda.
  4. Tool Lambda performs the heavy lifting, fanning out to three downstream services:
    1. Kinesis Data Streams – Lists streams in the Region, describes each stream (shard count, mode, retention), and discovers Enhanced Fan-Out consumers per stream.
    2. CloudWatch – Pulls 7 days of metrics per stream (IncomingBytes, OutgoingBytes, throttle events) and computes three-way cost comparison.
    3. Amazon S3 – Generates per-stream recommendations and an account-level Advantage assessment, then stores the final HTML and JSON reports.
  5. Amazon Bedrock AgentCore harness summarizes the findings and returns them to the caller.

The entire stack is deployed using AWS Cloud Development Kit (AWS CDK) with a single cdk deploy command.

Prerequisites

Before you begin, verify that you have the following:

  • AWS CDK – npm install -g aws-cdk.
  • Python 3.12+.
  • aws-cdk-lib >= 2.251.0 – for AgentCore L2 constructs.
  • AWS Command Line Interface (AWS CLI) configured with credentials that have permissions to deploy the required resources.
  • Amazon Bedrock model access – verify you have access to Claude Sonnet 4.5 (or your chosen model) in the target Region. Check in the Amazon Bedrock console under Model access.

Walkthrough

In the following sections, you deploy the Kinesis Data Streams Mode Optimizer Agent and test it against your Kinesis streams.

Step 1: Clone the repository

git clone https://github.com/aws-samples/sample-kinesis-optimizer-agent.git
cd sample-kinesis-optimizer-agent

Step 2: Install CDK dependencies

cd infra
pip install -r requirements.txt

Step 3: Set your target Region

The stack deploys to whatever Region is set in AWS_DEFAULT_REGION. Set it before running any CDK commands:

# Linux/macOS
export AWS_DEFAULT_REGION=us-east-1

# Windows PowerShell
$env:AWS_DEFAULT_REGION="us-east-1"

Step 4: Bootstrap CDK (first time per account/Region)

cdk bootstrap aws://<ACCOUNT_ID>/<REGION>

Step 5: Deploy

cdk deploy

Optionally customize the schedule and bucket name:

# Weekly instead of daily
cdk deploy --parameters ReportSchedule="rate(7 days)"

# Custom bucket name
cdk deploy --parameters ReportBucketName=amzn-s3-demo-bucket

Step 6: Test the agent

You can test the agent with the AWS CLI:

aws bedrock-agentcore invoke-harness \
    --harness-arn <HARNESS_ARN> \
    --runtime-session-id $(uuidgen) \
    --messages '[{"role":"user","content":[{"text":"Generate an optimization report"}]}]' \
    --region us-east-1

Step 7: View reports

Reports are stored in Amazon S3 at:

s3:// amzn-s3-demo-bucket/kinesis-optimization-reports/YYYY/MM/DD/HHMMSS-<report-id>.html
s3:// amzn-s3-demo-bucket/kinesis-optimization-reports/YYYY/MM/DD/HHMMSS-<report-id>.json

The HTML report includes a per-stream action table with priority indicators, detailed cost breakdowns, and the account-level Advantage recommendation.

Sample HTML report showing the per-stream action table with priority indicators and cost breakdowns

Figure 2: Sample HTML report with the per-stream action table

Sample report showing the account-level On-demand Advantage recommendation

Figure 3: Account-level On-demand Advantage recommendation in the sample report output

Multi-Region deployment

The agent is Region-specific. When deployed to a Region, it analyzes only the streams in that Region. To cover multiple Regions, change the environment variable and repeat the deployment:

export AWS_DEFAULT_REGION=eu-west-1
cdk bootstrap aws://<ACCOUNT_ID>/eu-west-1
cdk deploy

Each deployment is independent, with its own agent, Amazon S3 bucket, and schedule.

Clean up

To remove the stack from a Region:

cd infra
cdk destroy

Note: The Amazon S3 bucket has a RemovalPolicy.RETAIN setting and isn’t deleted with the stack. Delete it manually if you no longer need the reports.

Conclusion

In this post, you deployed an AI-powered agent that autonomously analyzes your Amazon Kinesis Data Streams and recommends the optimal capacity mode for each stream. The agent alleviates the manual effort of reviewing CloudWatch metrics across dozens or hundreds of streams and produces actionable, cost-aware recommendations on a recurring schedule.

By shifting from manual capacity reviews to autonomous, scheduled optimization, you gain three key benefits. First, you can reduce streaming costs by identifying streams that should switch modes. Second, you alleviate throttling risk by catching under-provisioned streams before they impact performance. Third, you free your team from repetitive operational work. All of this is achievable with a single cdk deploy.

To get started, clone the sample-kinesis-optimizer-agent repository and deploy it to your account today.


About the authors

Masudur Rahaman Sayem

Masudur Rahaman Sayem

Masudur is a Streaming Data Architect at AWS with over 25 years of experience in the IT industry. He collaborates with AWS customers worldwide to architect and implement sophisticated data streaming solutions that address complex business challenges. As an expert in distributed computing, Sayem specializes in designing large-scale distributed systems architecture for maximum performance and scalability. He has a keen interest and passion for distributed architecture, which he applies to designing enterprise solutions at internet scale.

Roy (KDS) Wang

Roy (KDS) Wang

Roy is a Senior Product Manager with Amazon Kinesis Data Streams. He is passionate about learning from and collaborating with customers to help organizations run faster and smarter. Outside of work, Roy strives to be a good dad to his new son and builds plastic model kits.

Automate SageMaker HyperPod incident triage and root-cause-analysis with AWS DevOps Agent

Post Syndicated from Tomonori Shimomura original https://aws.amazon.com/blogs/devops/automate-sagemaker-hyperpod-incident-triage-and-root-cause-analysis-with-aws-devops-agent/

Introduction

Large-scale machine learning workloads: training, fine-tuning, and inference run on clusters of hundreds to thousands of GPU instances for days or weeks at a stretch. Keeping operational visibility across a fleet of this size is a constant challenge: hardware health events, node lifecycle transitions, capacity fluctuations, and workload-level issues appear in the event stream around the clock, including nights and weekends.

Amazon SageMaker HyperPod is a purpose-built managed cluster service that lets you run distributed model training, fine-tuning, and inference across hundreds of accelerated instances. It provides built-in resiliency that automatically detects and replaces faulty hardware, so long-running jobs can continue with minimal interruption.

For teams operating these clusters, the scale still creates a fundamental tension: you need continuous visibility into your fleet, but you can’t afford to keep engineers watching the event stream 24/7. 

What HyperPod resiliency already handles

SageMaker HyperPod’s built-in resiliency layer automatically detects and self-heals instance-level GPU failures. When the Health Monitoring Agent (HMA) identifies a bad GPU, the HyperPod resiliency layer drains, reboots, or replaces the node depending on the error type, and the job resumes without human intervention. This is exactly what you want: routine hardware failures are handled automatically so your training runs keep going. 

This solution does not replace HMA or any part of HyperPod’s resiliency. It adds an autonomous investigation layer on top, using the cluster events and health signals that HMA and HyperPod already produce as its input. 

Operational conditions where a human still wants to be in the loop 

With that self-healing in place, there are operational conditions where a human still wants to be in the loop or decide: 

  • Configuration issues: a lifecycle-script change you made, a misconfigured mount, or a networking/security change that causes provisioning failures on every new node. 
  • Capacity conditions: a replacement waiting on capacity in the pool, where the operator needs to know recovery is in flight and can decide whether to intervene. 
  • Recurring hardware faults: each fault self-heals correctly, but the same GPU error signature recurring across three or more replacements on one instance group in a week is a pattern worth surfacing to an operator as a single signal. 
  • Workload-level conditions: Pods stuck in CrashLoopBackOff for hours, nodes sitting NotReady, or GPU allocation chronically low. 

Without automation, these conditions push operators into round-the-clock manual triage: correlating events across the SageMaker control plane, Amazon EKS, and Amazon CloudWatch, and deciding whether HyperPod is still recovering or needs a hand. 

Opportunity: AWS DevOps Agent as a 24/7 companion 

AWS DevOps Agent provides an autonomous incident-response platform that can be taught a domain’s operational model through custom skills. By wiring your HyperPod cluster into DevOps Agent, you get a 24/7 companion that complements HyperPod’s self-healing. It watches for the operational conditions that still need a human decision, triaging them, root-causing them, and delivering a clear verdict with recommended actions. 

By design, DevOps Agent is configured to run in observe-and-report mode for this integration – it is not granted SSM, SSH, or action-taking permissions against your cluster or its nodes. The agent reads cluster events, control-plane state, Kubernetes objects, and CloudWatch logs to reconstruct what happened; every corrective action (node reboots, replacements, drains) continues to be performed by HyperPod’s own resiliency layer or by an operator responding to the emailed verdict. This read-only boundary is deliberate: it keeps the agent’s blast radius zero while still delivering the correlation and triage value. 

​​In this post, you will learn how to connect any SageMaker HyperPod cluster (either the EKS or Slurm Orchestrator option) to AWS DevOps Agent. Conditions are auto-detected, triaged, root-caused from cluster state and CloudWatch logs, and emailed as a clear verdict. You will also see how the solution can be extended to detect additional conditions specific to your workloads.​ 

Solution Overview 

What this solution delivers 

This solution wires any SageMaker HyperPod cluster into AWS DevOps Agent so that operational conditions calling for a human decision are auto-detected, triaged, root-caused, and delivered as a human-readable verdict email. Specifically, you get: 

  • Autodetection of HyperPod conditions that complement resiliency self-healing, from the live SageMaker event stream and a periodic Kubernetes-state audit. 
  • Triage + root-cause analysis by the DevOps Agent, taught HyperPod’s operational model via two custom skills. It reconstructs the incident timeline and decides whether HyperPod is still recovering or needs an operator. 
  • Human-readable verdict emails: Monitor (recovery in flight, here’s the ETA), Escalate (you need to act, here’s why and what to do), or Resolved (auto-recovery closed the loop). Noise is filtered out. 
  • Extensibility: customize what conditions are detected (by modifying the periodic-audit Lambda) and how the agent reasons about them (by editing the plain-English skills). 

The following screenshot shows the DevOps Agent incident response dashboard with example verdict emails for three common fault types: 

DevOps Agent incident response dashboard showing investigation list and timeline, with three email verdict examples for GPU NVLink fault, lifecycle-script bootstrap failure, and insufficient-capacity errors

DevOps Agent incident response dashboard showing investigation list and timeline, with three email verdict examples for GPU NVLink fault, lifecycle-script bootstrap failure, and insufficient-capacity errors

Architecture 

The whole solution deploys one AWS CloudFormation stack per cluster. Two event paths feed the DevOps Agent, and one path carries its verdicts back out to you. 

Architecture diagram showing the event flow from HyperPod Health Monitoring Agent through EventBridge to DevOps Agent and email notification

Architecture diagram showing the event flow from HyperPod Health Monitoring Agent through EventBridge to DevOps Agent and email notification

This architecture shows a 1:1 relationship between a HyperPod cluster and a DevOps Agent space, and the deployment instructions in this post follow that model. If you need to associate multiple clusters with a single Agent Space, you can customize the CloudFormation template and the ClusterFilter parameter to widen the allowlist of cluster names forwarded by the webhook bridge.

Event flow 

  1. Event-driven issue detection: HyperPod emits cluster-state, node-health, and capacity events to Amazon EventBridge. The webhook bridge Lambda drops routine Info-level noise, maps the rest into a DevOps Agent investigation payload, signs it with HMAC-SHA256 using a shared secret stored in AWS Secrets Manager, and POSTs it to the agent’s generic webhook. 
  1. Polling-based issue detection: A periodic-audit Lambda checks Kubernetes state (CrashLoopBackOff pods, NotReady nodes) every 15 minutes and fires only when it finds a real issue, plus a daily heartbeat confirming the pipeline is alive. On a healthy cluster, nothing is POSTed, so no investigation runs and no cost is incurred. 
  1. Investigation: DevOps Agent receives the payload and runs two custom skills: the triage skill decides whether to link (duplicate), skip (noise), or proceed (investigate). The RCA skill reconstructs the timeline using describe-cluster, list-cluster-nodes, list-cluster-events, kubectl, and CloudWatch logs (HMA health monitoring, lifecycle scripts), then classifies the incident as Suppress, Monitor, Escalate, or Resolved. 
  1. Notification: An Amazon Lambda function sends notification emails via Amazon SES. It listens on the aws.aidevops event stream for investigation completions, reads the verdict from the agent’s journal, and sends an email with the headline, what happened, likely cause, and recommended action. Suppress verdicts are filtered to avoid noise on healthy clusters. 

​​Getting started​ 

For a step-by-step walkthrough to deploy this solution, visit the DevOps Agent Integration guide. Once you have the solution running, the following sections explain how to customize detection, reasoning, and notifications for your environment. 

Prerequisites 

  • An AWS account with AWS CLI v2 configured for the target region. 
  • An existing SageMaker HyperPod cluster (EKS or Slurm orchestrator). 
  • IAM permissions to create roles, deploy CloudFormation, manage Secrets Manager, and call devops-agent:* and eks:CreateAccessEntry. 
  • For email notifications: a verified Amazon SES sender identity. You can verify an email address in the Amazon SES console or with the AWS CLI. After running the command below, the address owner will receive a verification email and must click the confirmation link:

    aws ses verify-email-identity --email-address [email protected]

    Recipients must also be verified if your SES account is still in sandbox mode. 

Deploying with CloudFormation

The solution deploys as a single CloudFormation stack. Clone the awsome-distributed-ai repository, create a params.json with your cluster name and email settings, and run: 

cd 1.architectures/5.sagemaker-hyperpod/tools/devops-agent

# 1. Set up a Python env with boto3 >= 1.43.25
python3 -m venv .venv && source .venv/bin/activate && pip install 'boto3>=1.43.25'

# 2. Fill in your cluster name and email addresses
cp deploy/params.example.json deploy/params.json
# edit: HyperPodClusterName, EmailSender, EmailRecipients 

# 3. Deploy
make deploy

This provisions the Agent Space with read-only EKS access (auto-discovered from the cluster’s orchestrator ARN), the EventBridge rule and webhook bridge Lambda, the periodic-audit scheduler, and the email notifier. For Slurm-orchestrated clusters, the EKS access step is skipped automatically. 

The webhook bridge — mapping HyperPod events to DevOps Agent 

An EventBridge rule captures HyperPod events and invokes a Lambda function. The Lambda forwards all Warn and Error level events, normalizing each into a DevOps Agent investigation payload. It extracts the failure message, instance group, and event metadata, then signs it with HMAC using a shared secret stored in AWS Secrets Manager, and POSTs it to the agent’s generic webhook endpoint. Info-level events are dropped at the bridge to avoid creating investigations for routine status updates. 

A cluster allowlist parameter lets you scope which HyperPod clusters trigger investigations, useful when multiple clusters share the same account and region. 

How the skills are defined — teaching the agent HyperPod’s operational model 

AWS DevOps Agent skills are plain-English instructions that teach the agent how to reason about a domain. This solution includes two complementary skills: 

Triage skill — LINKED / SKIPPED / PROCEED (view the skill document) 

The triage skill runs first on every incoming task. It decides whether to link the event to an existing investigation, skip it, or proceed to a full investigation. 

  • Why triage matters — a concrete example: When a single node fails, HyperPod’s replacement process emits multiple events in quick succession: “lost orchestration-ready status,” “provisioning started,” “capacity request initiated.” Without triage, each event would spawn a separate investigation. The triage skill recognizes these events belong to the same incident (same instance group + overlapping time window) and links them, so only one investigation runs. This saves investigation compute and avoids duplicate emails. 
  • When to SKIP: When a node is already being replaced and a follow-up “lost orchestration-ready status” event arrives with a generic “Request to service failed” message, the triage skill recognizes that a replacement is already in progress for that instance group and skips the event. No new investigation is created for what is simply a progress update of an existing recovery. 

RCA skill — timeline reconstruction and verdict (view the skill document) 

When triage produces PROCEED, the RCA skill takes over. It reads cluster state, events, and logs, reconstructs an incident timeline, and classifies the situation into one of four verdicts:

 RCA Flowchart showing the four phases of root-cause analysis: data gathering, timeline reconstruction, classification, and recurrence check

RCA Flowchart showing the four phases of root-cause analysis: data gathering, timeline reconstruction, classification, and recurrence check

  • Phase 1 — Data gathering: The skill reads describe-cluster, list-cluster-nodes, list-cluster-events, and CloudWatch log streams (HMA health monitoring, lifecycle scripts) to collect the raw facts. 
  • Phase 2 — Timeline reconstruction: It orders events chronologically and identifies the fault chain: what triggered what, which nodes were affected, and what recovery actions HyperPod took. 
  • Phase 3 — Classification: Based on the timeline, recurrence statistics, and HyperPod’s resiliency behavior, it assigns a verdict: 
    • Suppress — a non-issue (for example, a transient event that has already resolved). 
    • Monitor — recovery is in flight; here’s the expected resolution window. 
    • Escalate — you need to act; here’s the root cause and recommended action. 
  • Resolved — auto-recovery closed the loop; no action needed. 
  • Phase 4 — Recurrence check: The skill computes sliding-window statistics over the one week cluster event history. When thresholds are crossed, the verdict escalates to alert the operator of a systemic pattern. For example, the same GPU error signature on the same instance group three or more times in a week, or five or more replacements fleet-wide in 24 hours. 

The verdict is written to the agent’s investigation journal along with a human-readable report containing what happened, the likely cause, and recommended operator actions. 

The periodic-audit Lambda — Kubernetes state monitoring 

The periodic-audit Lambda fires every 15 minutes and inspects Kubernetes Pod/Node state directly (via the EKS API server). It checks for: 

  • Pods in CrashLoopBackOff (default: flagged when restart count reaches five and the last crash is within 15 minutes) 
  • NotReady nodes (default: flagged when a node has been NotReady for at least 15 minutes and at least 10% of nodes are affected) 

Namespace-aware filtering controls which pods are checked: 

  • Pods in kube-public and kube-node-lease are ignored entirely by default. 
  • Pods in kube-system, aws-hyperpod, and amazon-cloudwatch are tagged as system-workload issues (distinct from user-workload issues in the verdict). 

All thresholds and namespace lists are configurable via the CloudFormation stack parameters. 

The Lambda POSTs a webhook event to DevOps Agent only when a real issue is found. On a healthy cluster, nothing is POSTed, so no investigation runs and no cost is incurred. A separate daily heartbeat schedule confirms the monitoring pipeline itself is alive. The heartbeat is visible in the DevOps Agent console but deliberately not emailed on healthy runs — so silence in your inbox means the cluster is healthy, not that the pipeline is broken. 

Note: HyperPod infrastructure faults (node health, capacity errors, lifecycle-script failures) are handled event-driven by the webhook bridge. They come from the native HyperPod event stream in EventBridge. The periodic audit deliberately does not duplicate that path; it only covers Kubernetes workload state, which is not in the HyperPod event stream. 

Closing the loop — the email notifier

An EventBridge rule on the aws.aidevops event stream captures investigation lifecycle events. The email-notifier Lambda processes these events through the following steps: 

  1. Event filtering: Only “Investigation Completed” events are processed (one email per investigation lifecycle). The event payload contains the agent_space_id, task_id, and execution_id. 
  2. Dedup: The Lambda checks an S3 marker at s3://<bucket>/emailed/<execution_id>. If present, this investigation has already been emailed and the event is dropped. This prevents duplicate emails when the same completion event is re-emitted. 
  3. Fetching the investigation context: The Lambda calls two DevOps Agent APIs: 
    • get_backlog_task(agentSpaceId, taskId) — retrieves the task metadata (title, priority, timestamps). 
    • list_journal_records(agentSpaceId, executionId) — retrieves the investigation’s findings, symptoms, and investigation gaps from the agent’s journal. 
  4. Suppress-verdict filtering: If the investigation produced a Suppress verdict or no findings at all, no email is sent. 
  5. Email composition: The Lambda composes a single HTML email from the journal records: a short headline followed by a one-paragraph summary covering what happened, the likely cause, and the recommended action. 
  6. Send via SES: The formatted email is sent to the configured recipients. After successful delivery, the S3 dedup marker is written. 

The operator also has access to the full investigation in the DevOps Agent web console (see following “Viewing investigations” section). 

Viewing investigations in the DevOps Agent console 

For readers new to AWS DevOps Agent, here’s how to navigate to your investigations: 

  1. Open the AWS DevOps Agent console. 
  2. Select your Agent Space (named hyperpod-<cluster-name>-devops-agent by default). 
  3. From the Launch web app drop-down, choose an option to open the DevOps Agent web app. 
  4. Select Incidents from the left navigation pane to open the Incident Response Dashboard. It lists all investigations with their subject, status, and timestamp. 
  5. Select any investigation to see its full timeline, journal records, and the verdict report. 

Asking the agent directly — the DevOps Agent Chat UI 

Beyond the automated emails, you don’t have to wait for the next investigation to get answers about your cluster. You can open the DevOps Agent’s AI chat at any time and ask follow-up questions in plain English. The agent answers from the live cluster state, the investigation history, and the skills it has been taught. 

For example: 

  • “I got an email about a GPU failure in my cluster. Did it get resolved now with HyperPod’s resiliency?” — The agent checks the current cluster state, confirms whether the replacement succeeded, and provides a timeline of what happened (HMA detection → replacement initiated → node back in service), along with anything to watch for. 
  • “Are there unhealthy Pods on my cluster?” — The agent inspects the Kubernetes state and reports any CrashLoopBackOff pods or NotReady nodes. 
  • “I just triggered scaling up. Check if it is progressing well.” — The agent looks at the cluster’s current node counts vs. target counts and reports whether provisioning is on track. 
AWS DevOps Agent chat interface showing a natural language query about cluster health

AWS DevOps Agent chat interface showing a natural language query about cluster health

The chat conversations are stored per Agent Space, so you can revisit past interactions alongside the automated investigations. This makes the Agent Space a single pane of glass for both automated incident response and ad-hoc troubleshooting of your HyperPod cluster. 

Extending the solution — detection vs. reasoning 

The solution has two extension points, which serve different purposes: 

  1. Extending detection (what conditions are caught): 
    • Event-driven path: The webhook bridge Lambda drops Info-level events and forwards all Warn and Error level HyperPod events to DevOps Agent. This typically does not need modification. It already catches all actionable events. 
    • Polling-based path: The periodic-audit Lambda checks Kubernetes state. To detect additional conditions (for example, GPU allocation below a threshold or specific Pod labels stuck in error states), add that logic to the Lambda code. 
  2. Extending reasoning (how the agent investigates and classifies): edit the plain-English skill definitions. For example, you can teach the RCA skill new classification rules, add domain-specific context about your workload’s expected behavior, or adjust the recurrence thresholds. 

Detection is code; reasoning is natural language. Both are in the repo and designed to be customized independently. 

Investigation feedback 

After each investigation completes, a Feedback button appears in the DevOps Agent console. Clicking it opens the Investigation feedback dialog, where you can: 

  • Rate whether the root cause was correct 
  • Indicate whether human steering was needed during the investigation 
  • Provide written feedback explaining what could be improved 

This structured feedback is stored per investigation. An auto-learning mechanism that uses this feedback to improve future investigations is actively being developed. 

DevOps Agent APIs used by this solution 

For readers interested in the programmatic integration, here are the key DevOps Agent APIs this solution calls:

Component API Purpose
Webhook provisioner (deployment) register_service Register the generic webhook service with DevOps Agent
Webhook provisioner (deployment) associate_service Associate the webhook with the Agent Space
Skill uploader (deployment) list_assets Check if a skill already exists
Skill uploader (deployment) create_asset / update_asset Upload or update the triage and RCA skill definitions
Email notifier (runtime) get_backlog_task Retrieve task metadata (title, priority, timestamps)
Email notifier (runtime) list_journal_records Retrieve findings, symptoms, and gaps from the investigation journal
Teardown disassociate_service / deregister_service / delete_asset Clean up on stack deletion

Cleaning up 

To remove all resources created by this solution, run: 

make teardown-stack

This deletes the CloudFormation stack, removes the Agent Space, EKS access entries, secrets, and email configuration.

Additionally, if you no longer need the prerequisite resources, you can revert their setup, for example, deleting the verified Amazon SES email address identities you created for notifications.

Cost considerations 

​​​This solution is designed to be near-zero cost on a healthy cluster and scales proportionally with fault volume. Cost scales with fault volume, not node count directly. At large scale (100+ nodes), the triage skill becomes critical. A single hardware fault can generate 5-10 correlated EventBridge events, most of which are filtered by the webhook bridge Lambda before reaching the agent. Where triage adds value is linking and deduplicating across similar faults that affect multiple instances, or repeated faults on the same instance over time, consolidating them into a single investigation instead of many. As an example, a 500-node training cluster might see 20-50 investigations per month after filtering and deduplication.​​ 

  1. ​​​Filtering and triage are your cost savers at scale. The webhook bridge filters correlated events from a single node failure (5-10 EventBridge events reduced to 1 forwarded event), eliminating redundant investigations at the source. Triage then links similar faults across multiple instances into a single investigation. For example, if 5 nodes hit the same GPU error in a window, triage consolidates them into 1 investigation instead of 5 (saving 4 × $4 = $16). The bigger the cluster, the more both layers save.​​ 
  2. Investigation duration grows sub-linearly. A 1000-node cluster investigation doesn’t take 100x longer than a 10-node one. The agent queries describe-cluster and list-cluster-events once regardless of size. The data returned is bigger, but the API call count is similar. 
  3. CloudWatch Logs queries are the variable. On large clusters, the agent may query more HMA log streams, which takes longer agent-seconds AND incurs CloudWatch Logs Insights charges on your account (not part of DevOps Agent pricing). 

DevOps Agent (the primary cost driver): Estimates based on 2 accelerator instances in a cluster 

Component Pricing Your cluster estimate
Investigations $0.0083/agent-second ~$4/investigation (at 8 min avg)
Chat (on-demand SRE tasks) $0.0083/agent-second ~$0.25/chat query (at 30 sec avg)
Daily heartbeat $0.0083/agent-second ~$1-2/day (short investigation confirming health)

On a healthy cluster with no faults, only the daily heartbeat fires, approximately $30-60/month in DevOps Agent time. On a cluster experiencing 5 real faults per week (typical for a large GPU fleet), expect ~20 investigations/month × $4 each = $80/month in investigation costs.

Free tier and credits:

New DevOps Agent customers receive a 2-month free trial (20 hours of investigations, 20 hours of chat per month). Enterprise Support customers receive monthly credits equal to 75% of their AWS Support charge toward DevOps Agent usage. 

Supporting infrastructure (secondary costs): 

Component Monthly Cost Estimates
Lambda invocations ~96/day (15-min audit) + event-driven = well within free tier
S3 (skills + dedup markers) < $0.01 (a few MB total)
Secrets Manager (1 secret) $0.40
EventBridge rules Negligible (per-event pricing)
SES emails $0.10/1000 emails — at most 1 per investigation
CloudWatch Logs (Lambda) < $1 (minimal log volume)

Total estimated monthly cost: 

Scenario DevOps Agent Infrastructure Total
Healthy cluster (no faults) ~$30-60 (heartbeat only) < $2 ~$32-62/month
Moderate faults (5/week) ~$80-120 < $2 ~$82-122/month
Heavy faults (20/week) ~$320-400 < $5 ~$325-405/month

How cluster size impacts cost 

Factor Small cluster (1-10 nodes) Large cluster (100-1000 nodes)
Fault frequency Rare (maybe 1-2/week) Constant (NVIDIA reports ~1 fault/2-3 hours at 10K GPU scale)
Events per fault Few (1 node replacement = 3-5 events) More (cascading replacements, capacity queuing)
Investigation duration Shorter (less state to read, fewer events in timeline) Longer (more nodes to describe, more events to correlate, larger CloudWatch log groups to query)
Triage value Low (few duplicates) High (one fault generates many correlated events — triage links them into 1 investigation)
Periodic audit Fast (few pods/nodes to check) Slower (more K8s state to inspect)
Cluster Size Faults/month Investigations Est. Agent Cost
1-10 nodes (your test) 2-5 2-5 + heartbeat $8-20/mo + ~$30 heartbeat
10-50 nodes (typical prod) 5-20 5-15 (triage dedup) $20-60/mo + ~$30 heartbeat
100-500 nodes (large training) 50-200 20-50 (heavy triage) $80-200/mo + ~$45 heartbeat
1000+ nodes (frontier) 200-700 50-100 (massive dedup) $200-500/mo + ~$60 heartbeat

Cost control levers: 

  1. Disable the periodic audit (EnablePeriodicAudit: false) to eliminate the heartbeat cost. Live event bridging still works. 
  2. Triage (LINK/SKIP decisions) runs at task creation time. No investigation cost is billed for deduplicated or skipped events. 
  3. Suppress verdicts filter email notifications but the investigation still runs. If you want to eliminate that cost, tune your EventBridge rule to drop more event types at the bridge level. 

Comparison to manual monitoring: 

Without automation, each fault requires an on-call engineer to manually correlate events across CloudWatch, EKS, and the SageMaker console, typically 30-45 minutes of triage before they even know whether HyperPod is self-healing or needs intervention. This solution delivers a root-caused verdict in minutes at ~$4 per investigation, while providing 24/7 coverage without human wake-ups. The cost savings compound with cluster scale: at 20 faults per month, that’s 10-15 hours of engineering triage replaced by automated verdicts. 

Conclusion 

In this post, we showed how to build an end-to-end agentic incident-response pipeline for SageMaker HyperPod using AWS DevOps Agent. The solution complements HyperPod’s built-in resiliency by watching for the operational conditions where a human still wants to be in the loop: configuration issues affecting provisioning, capacity-bound recoveries, recurring hardware fault patterns, and workload-level conditions. It delivers clear, root-caused verdicts to the operator’s inbox. 

The broader takeaway is a reusable pattern: teaching an AI agent a domain’s operational model through plain-English skills, so it can distinguish “the system is recovering on its own” from “this needs a human decision.” This pattern applies beyond HyperPod to any event-driven AWS service where operational conditions benefit from automated correlation and triage. 

What’s next 

​​​To deploy the solution, follow the step-by-step instructions in the DevOps Agent Integration guide on the AI on SageMaker HyperPod site. Once it’s running, you can customize it for your environment:​​ 

  • Adjust the CloudFormation parameters: tune the periodic-audit schedule, CrashLoopBackOff thresholds, NotReady node percentages, namespace filtering, and email recipients. No code changes required. 
  • Extend detection: modify the periodic-audit Lambda to check for additional Kubernetes conditions specific to your workloads (for example, GPU allocation below a threshold, specific Pod labels stuck in error states). 
  • Extend reasoning: edit the triage or RCA skill definitions to adjust classification rules, add domain context about your expected cluster behavior, or tune the recurrence thresholds. 
  • Add notification channels: connect Slack or PagerDuty via DevOps Agent’s built-in integrations or via a sibling EventBridge rule on the same aws.aidevops event stream. 

The skills are plain English. Iterate on them the same way you’d iterate on a runbook. 

About the authors

Tomonori Shimomura is a Principal Solutions Architect on the Amazon SageMaker AI team, where he provides in-depth technical consultation to SageMaker AI customers and suggests product improvements to the product team. Before joining Amazon, he worked on the design and development of embedded software for video game consoles, and now he leverages his in-depth skills in Cloud side technology. In his free time, he enjoys playing video games, reading books, and writing software.

Mayank Gupta is a Senior AI/ML Specialist with deep expertise in machine learning frameworks and enterprise AI architecture. He brings strong hands-on experience with AWS AI services, including SageMaker AI and SageMaker AI HyperPod, and leads the design and delivery of end-to-end AI solutions spanning model development, distributed training, and production-scale deployment. With deep experience in performance optimization and scalable ML architectures, Mayank partners with customers to translate complex business challenges into secure, high-impact, production-ready AI systems that drive measurable outcomes.

Deepthi Madamanchi is a Principal Technical Account Manager at AWS focused on AI Models, where she leads frontier AI segment through building and operating multi-thousand-node GPU clusters for foundation model training and inference. She specializes in distributed training, high-throughput networking, GPU fleet optimization, Amazon Bedrock adoption, helping them optimize performance, reliability, and cost efficiency from experimentation through production. In her free time, Deepthi explores functional health, experiments with new recipes, and travels with her family.

Dushyant Dubaria is a Senior Technical Account Manager on the AWS Frontier AI Startup team, where he supports frontier AI model builder companies deploying and operating large-scale GPU training infrastructure on Amazon SageMaker HyperPod and Amazon EKS. He specializes in distributed training orchestration, storage at petabyte scale (Amazon FSx for Lustre, Amazon S3), high-throughput networking, and operational resilience including cluster health monitoring, capacity planning, and proactive incident management for multi-thousand-node clusters. He helps organizations achieve reliable, high-performance ML workloads from initial cluster deployment through sustained production training. In his free time, he enjoys building automation tools, exploring new AI technologies, and playing cricket.

Shreyas Adiyodi is a Product Manager at AWS based out of Seattle. He is focused on enabling Gen AI model development on SageMaker HyperPod, partnering with customers to simplify cluster provisioning, accelerate foundation-model training, and strengthen security and compliance. Outside of work, he enjoys chess, MMA and watching movies.

Scaling organizational knowledge in Kiro with Amazon Bedrock Knowledge Bases, LangChain, and MCP

Post Syndicated from Sakshi Singh original https://aws.amazon.com/blogs/devops/scaling-organizational-knowledge-in-kiro-with-amazon-bedrock-knowledge-bases-langchain-and-mcp/

“A pull request comes back with a single comment: “This doesn’t follow our circuit breaker pattern. Check the Architectural Decision Record .” 

You know the architecture decision record exists somewhere. You open your team’s wiki, search “circuit breaker,” scroll past six irrelevant results, find the document, read through it, switch back to your editor, and fix the code. Fifteen minutes are gone. Not because the problem was hard, but because the knowledge lived in one place and the code lived in another.

This plays out multiple times a day across engineering teams. Developers face several recurring challenges when working with organizational knowledge:

  • Context switching – Retrieving coding standards, API specs, or architecture decisions means leaving the editor to search wikis, shared drives, or documentation portals
  • Knowledge fragmentation – Team knowledge lives across multiple systems, making it difficult to find the right document at the right time
  • Onboarding friction – New team members spend days navigating unfamiliar documentation structures before becoming productive
  • Stale compliance – Code reviews catch standards violations after the fact, instead of surfacing the correct pattern during development

The documentation exists and is well structured. But it is not accessible from where development happens.

In this post, we show how to connect Amazon Bedrock Knowledge Bases to Kiro through the Model Context Protocol (MCP), enabling developers to query team documentation directly from their editor and get cited answers quickly. Kiro is an agentic IDE that uses MCP to connect developers to external knowledge sources beyond the local workspace. Whether you already have a Knowledge Base or are building one from scratch, setup typically takes a few minutes.

Why MCP with Knowledge Bases When Kiro Already Has Steering and Agent Skills

Kiro provides several built-in mechanisms to give context to the agent:

  • Steering files (.kiro/steering/*.md) deliver static instructions and project-level context. They can be included, conditionally matched by file pattern, or manually referenced. Ideal for coding standards, team conventions, and project-specific rules that fit in a few files.
  • Agent Skills (.kiro/skills/) offer reusable instructions that users activate to guide agent behavior for specific workflows like code reviews, testing strategies, or deployment procedures.
  • File references (#File, #Folder) provide explicit references to local workspace files for point-in-time context.

The MCP with Knowledge Bases approach is complementary, not a replacement. Use Steering for the ten rules every commit must follow. Use Agent Skills for workflow guidance. Use MCP with Knowledge Bases when your organization maintains hundreds of Architectural Decision Records, API specs, runbooks, security guidelines, and onboarding documents. No developer can internalize all of it. Semantic search surfaces the right answer at the right moment.

Together these serve distinct roles: Steering governs Kiro’s behavior, Knowledge Bases hold your organization’s collective knowledge, and MCP provides the connective layer that makes that knowledge accessible to Kiro on demand.

Solution overview

Amazon Bedrock Knowledge Bases has powered RAG workloads for multiple teams since well before Kiro launched. If your team already has a Knowledge Base, you have completed the foundational setup: documents curated, vectors indexed, knowledge layer built. What follows is a five-minute integration that brings all of it into the editor.

The question is not whether to start from scratch. It is simpler than that: how do you bring what you already have into Kiro?

In this integration, the awslabs.bedrock-kb-retrieval-mcp-server bridges the gap between Kiro and your Knowledge Base, translating natural language queries into vector search operations and returning cited passages directly in the editor.

The answer is a single configuration file and an MCP server that takes less than few minutes to connect.

The use cases that change daily workflows

Before we dive into the how, consider what becomes possible when your Knowledge Base lives inside your editor:

Coding standards enforcement in real time. A developer asks Kiro: “What’s our error handling pattern?” and gets back the exact custom error class structure your team agreed on six months ago, complete with the code snippet from your standards document.
API specifications at your fingertips. Instead of opening a browser tab to check authentication requirements, a developer types: “What authentication does the Orders API require?” and immediately sees the JWT scope requirements, header format, and rate limits pulled directly from your OpenAPI spec stored in the Knowledge Base.

Architecture decisions with full context. When someone needs to understand why a decision was made, not just what was decided, they ask Kiro. The Architectural Decision Record comes back with the rationale, the alternatives considered, and the tradeoffs, all cited with source documents.

Kiro CLI in CI/CD. Run headless queries against your Knowledge Base in pipelines. Validate that generated code matches team patterns. Automate compliance checks against your security guidelines during pull request reviews.

Two paths: bring what you have or start fresh

You already have a Knowledge Base

If your team already uses Amazon Bedrock Knowledge Bases, whether it was built for a chatbot, an internal search tool, or a customer-facing assistant, you don’t need to rebuild anything. Your existing Knowledge Base works with Kiro out of the box.

Here’s the approach:

  1. Tag your existing Knowledge Base with mcp-multirag-kb=true. This is how the MCP server discovers it.
  2. Configure the MCP server in Kiro (covered in the next section). Your documents, your embeddings, your vector store, all stay exactly where they are.

The official awslabs.bedrock-kb-retrieval-mcp-server auto-discovers Knowledge Bases with that tag. If you have multiple Knowledge Bases (one for API docs, another for architecture decisions, a third for runbooks), tag them all. Kiro can query across your tagged Knowledge Bases.

You don’t have a Knowledge Base yet

If you’re starting fresh, the accompanying sample repository provides a complete AWS CDK application that deploys everything you need: an Amazon S3 bucket for your documents, an Amazon OpenSearch Serverless collection for vector search, and an Amazon Bedrock Knowledge Base that ties it together. The setup script handles deployment in few minutes.

For the full infrastructure deployment walkthrough, including CDK stack details, document ingestion, and monitoring setup, see the repository README.
After the setup script completes, you see the following output confirming the deployment and providing next steps:

Setup script completion output showing MCP config ready, Knowledge Base tag set for auto-discovery, and sample queries
Figure 1: Setup script completion output. The script confirms the MCP config is ready, the Knowledge Base tag is set for auto-discovery, and provides sample queries to test immediately.  

How it works

The Model Context Protocol (MCP) is what connects Kiro to your Knowledge Base. It acts as a bridge: Kiro connects via MCP on one side, Amazon Bedrock Knowledge Bases uses its Retrieve API on the other, and the MCP server translates between them.

The Architecture Diagram in Repository shows the end-to-end integration.

When you ask Kiro a question, the following sequence occurs:

  1. Developer asks a question – You type a natural language query in Kiro (IDE or CLI).
  2. MCP request – Kiro sends your query to the MCP server running as a local child process over stdio.
  3. Retrieve API call – The MCP server calls the Amazon Bedrock Knowledge Bases Retrieve API (not RetrieveAndGenerate).
  4. Vector search – Amazon Bedrock embeds your query using Amazon Titan Text Embeddings v2 and searches the Amazon OpenSearch Serverless vector store.
  5. Ranked chunks returned – The MCP server receives ranked document chunks with relevance scores and passes them back to Kiro.
  6. Kiro generates the response – Kiro’s own LLM synthesizes the retrieved chunks into a cited answer and presents it directly in your editor.

The official MCP server handles retrieval only. Kiro handles the generation, which means the quality of the response benefits from Kiro’s full conversation context and reasoning capabilities.You get cited answers directly in your editor, no context switching required.

Prerequisites

You need the following to connect the MCP server to Kiro:

  • Kiro IDE or CLI installed on your machine
  • uv package manager (provides uvx for running the server without installation)
  • AWS CLI v2 configured with credentials that have bedrock:Retrieve permissions
  • An existing Amazon Bedrock Knowledge Bases (or deploy one using the sample repository)

Connect your Knowledge Base to Kiro

Create or update .kiro/settings/mcp.json in your project root

{ 
  "mcpServers": { 
    "awslabs.bedrock-kb-retrieval-mcp-server": { 
      "command": "uvx", 
      "args": ["awslabs.bedrock-kb-retrieval-mcp-server@latest"], 
      "env": { 
        "AWS_PROFILE": "default", 
        "AWS_REGION": "<YOUR_REGION>", 
        "FASTMCP_LOG_LEVEL": "ERROR", 
        "KB_INCLUSION_TAG_KEY": "mcp-multirag-kb", 
        "BEDROCK_KB_RERANKING_ENABLED": "false" 
      }, 
      "disabled": false, 
      "autoApprove": [] 
    } 
  } 
} 

Replace <YOUR_REGION> with the region where your Knowledge Base lives.

– BEDROCK_KB_RERANKING_ENABLED controls whether the server applies Amazon Bedrock’s reranking model to re-score retrieved chunks by relevance before returning them. Set to “true” to enable reranking for higher-quality results at the cost of additional latency and reranking model charges. The default is “false”, which returns results ranked by vector similarity only.

– Note on permissions: Kiro inherits the same AWS permissions as the profile specified in AWS_PROFILE. The MCP server runs as your local process, so it uses your configured credentials directly. If your profile has broad permissions, Kiro can exercise all of them. For production Knowledge Bases, use a profile with least-privilege access – bedrock:Retrieve is sufficient for read-only queries.

Key settings:

  • command: “uvx” runs the server without installing anything permanently. It downloads, executes, and cleans up automatically.
  • KB_INCLUSION_TAG_KEY tells the server to auto-discover any Knowledge Bases tagged with mcp-multirag-kb=true.
  • autoApprove is empty by default. Add “ListKnowledgeBases” and “QueryKnowledgeBases” to skip confirmation prompts for read-only queries. Both tools are read-only — they retrieve data from your Knowledge Base without modifying it, so auto-approving them is appropriate for read-only workflows.

Restart Kiro. The MCP server connects and discovers your tagged Knowledge Bases automatically.

What this looks like in practice

Same pull request. Same reviewer comment about the circuit breaker pattern. But this time, you do not open a browser. You ask Kiro:
"What's our circuit breaker pattern?"
Kiro calls the MCP server, queries the Knowledge Base, and returns the result directly in your editor:

Kiro querying the Knowledge Base for the circuit breaker pattern, showing ListKnowledgeBases discovery, local ADR file reading, and QueryKnowledgeBases returning the full parameter table from ADR-001 with source attribution
Figure 2: Kiro querying the Knowledge Base for the circuit breaker pattern. It calls ListKnowledgeBases to discover tagged Knowledge Bases, reads the local ADR file, and calls QueryKnowledgeBases to return the full parameter table from ADR-001 with source attribution.

The response includes the architecture decision record, the specific parameters (failure threshold, reset timeout, success threshold), and the source file reference. You fix your code quickly — no context switch, no browser tab, no searching.

Example: Querying API specifications

A developer types: "What authentication does the Orders API require?"

Kiro returns:

All requests require a valid JWT in the Authorization: Bearer <token> header. Tokens are issued by the Auth Service and must include the orders:read or orders:write scope.
Source: api-spec-orders.md 

Example: Discovering documentation gaps

A teammate asks Kiro: "What security headers should our APIs return?" 
The MCP server queries the Knowledge Base and returns the security guidelines document, which covers authentication, input validation, and secrets management — but does not mention HTTP response security headers. Kiro recognizes this gap in the retrieved content and, using its own workspace context (Kiro can read local files like security-guidelines.md independently of the MCP server), recommends the headers that should be added based on the existing security posture documented elsewhere.

Kiro querying security guidelines from the Knowledge Base, showing the MCP server returning existing security posture including JWT handling, input validation, and secrets management, with Kiro identifying the missing HTTP response security headers section
Figure 3: Kiro querying security guidelines from the Knowledge Base. The MCP server returns the existing security posture (JWT handling, input validation, secrets management), and Kiro identifies the missing HTTP response security headers section, recommending additions based on the documented security context.

This illustrates how Kiro combines Knowledge Base retrieval with its native workspace awareness. The MCP server handles the retrieval; Kiro handles the reasoning across all available context.

The LangChain alternative: a cloud-agnostic approach with more control

The official MCP server covers most use cases. For advanced scenarios – provider portability (swap between Amazon Bedrock, OpenAI, or local models), server-side RAG with built-in relevance filtering, or custom LCEL chain composition, see the LangChain alternative section in the repository README.
You can run both servers simultaneously. Kiro selects the right tool based on your query.

Both MCP servers running simultaneously, showing Kiro calling ask_knowledge_base on the LangChain server and ListKnowledgeBases on the official server in parallel, then falling back to QueryKnowledgeBases to retrieve security guidelines for API authentication from kiro-dev-knowledge-base
Figure 4: Both MCP servers running simultaneously. Kiro calls `ask_knowledge_base` on the LangChain server and `ListKnowledgeBases` on the official server in parallel, then falls back to `QueryKnowledgeBases` to retrieve the full security guidelines for API authentication from the kiro-dev-knowledge-base.

For the complete LangChain setup, including provider swapping (OpenAI, Ollama, local models) and LCEL chain details, see the LangChain alternative section in the repository.
The Architecture Diagram for Langchain alternative in Repository shows the end-to-end integration.

Best practices for your Knowledge Base content

The quality of answers depends on the quality of your documents:

  • Write Markdown with clear headings. The 512-token chunking works best with self-contained sections under each heading.
  • Include code examples. Developers use returned snippets immediately. An error handling standard with a code sample is ten times more useful than one without.
  • Use consistent naming. If your API is called “Orders API” in one document and “Order Service” in another, retrieval suffers.
  • Keep documents current. Stale docs erode trust faster than missing docs. Set a quarterly review cadence.

Kiro CLI: Knowledge Base queries in your terminal and CI/CD

The same MCP configuration works for both Kiro IDE and Kiro CLI:

# Interactive
kiro-cli chat
# Headless (for scripts and pipelines)
kiro-cli chat --no-interactive --trust-tools=read \
"What's our circuit breaker pattern?" 

The --no-interactive runs without a session, and – --trust-tools=read auto-approves read-only tool calls (like QueryKnowledgeBases) without prompting. Headless mode requires the KIRO_API_KEY environment variable. To generate an API key, follow the steps in the Kiro Documentation.

Use headless mode in CI/CD pipelines to validate generated code against team standards, or in onboarding scripts that walk new developers through your architecture decisions.

Cleanup

The MCP server is an open-source tool; costs apply to the underlying AWS resources (Amazon OpenSearch Serverless, Amazon S3 storage, and Amazon Bedrock API calls). The primary ongoing cost is Amazon OpenSearch Serverless, which charges for OCU (OpenSearch Compute Unit) capacity even when idle. Amazon S3 storage and Amazon Bedrock API calls are pay-per-use. For detailed pricing, see the Amazon S3 Pricing page and Amazon Bedrock Pricing page. Destroy resources when you’re done experimenting:

cd kiro-bedrock-kb-mcp/infrastructure
npx cdk destroy --all

For detailed cleanup instructions, see the repository README.

Conclusion

In this blog post, we showed how to connect Amazon Bedrock Knowledge Bases to Kiro through MCP, turning organizational documentation into an in-editor knowledge assistant. This integration addresses the challenges outlined at the beginning of this post:

  • No more context switching – Developers query coding standards, API specs, and architecture decisions without leaving their editor
  • Unified knowledge access – A single MCP configuration connects to multiple Knowledge Bases, regardless of where the original documents live
  • Faster onboarding – New team members get cited answers to questions quickly, without navigating unfamiliar documentation systems
  • Proactive standards enforcement — Team standards surface during development rather than after a code review catches a violation.

Two paths to get started:

  • Existing Knowledge Base – Tag it with mcp-multirag-kb=true, add the MCP configuration to Kiro, and start querying after few minutes.
  • Starting fresh – Deploy the sample infrastructure using the repository, upload your team documents, and connect.

Your documentation already held the answers. Now developers get them quickly, without leaving their workflow.

About the author

Sakshi Singh

Sakshi Singh

Sakshi is an Associate Delivery Consultant at AWS Professional Services GCC, specializing in mainframe modernization and generative AI solutions. She helps organizations transform legacy systems into modern, cloud-native architectures on AWS, leveraging AI-driven approaches to accelerate migration. Her work bridges traditional enterprise infrastructure and cutting-edge cloud technologies, delivering scalable solutions that drive business value.

Nishtha Yadav

Nishtha Yadav

Nishtha is an Associate Delivery Consultant at AWS Professional Services, specializing in DevOps and AI-powered developer tooling. She works on infrastructure automation and generative AI solutions, helping customers streamline DevOps workflows and accelerate delivery. With a passion for solving complex automation challenges, she brings creativity and technical depth to every engagement. Outside work, she loves her dogs and gaming.

Yashika Baranwal

Yashika Baranwal

Yashika is an Associate Delivery Consultant at AWS Professional Services GCC, helping enterprises design and deliver modern cloud solutions. She specializes in cloud-native application development, with expertise in serverless architectures and generative AI. Her work focuses on building scalable, AI-powered applications that modernize enterprise infrastructure, turning complex challenges into production-ready implementations on AWS.