Tag Archives: Advanced (300)

Introducing AWS Glue 6.0 for faster and more cost-effective data integration

Post Syndicated from Aarthi Srinivasan original https://aws.amazon.com/blogs/big-data/introducing-aws-glue-6-0-for-apache-spark/

Organizations running large data processing pipelines want lower costs, faster job runtimes, and dependable support for open table formats, without adding operational overhead. AWS Glue, a serverless, scalable data integration service that you can use to discover, prepare, move, and integrate data from multiple sources, has now launched AWS Glue 6.0, the new version of AWS Glue that addresses these needs. This version upgrade lowers AWS Glue pricing by 30% and improves performance with AWS optimized Apache Spark 4.1. It also augments developer experience with new features and adds support for Apache Iceberg V3 specifications that are suitable for enterprise adoption. The newly available AWS Glue 6.0 makes data processing workloads more manageable, faster to run, and easier to operate.

In this post, we cover the key capabilities of AWS Glue 6.0 and their performance benefits. We share code examples to help you take full advantage of the release, and we show you how to get started.

AWS Glue 6.0 highlights

AWS Glue 6.0 brings together four major improvements designed to transform how you build and run data integration workloads.

First, it upgrades the underlying runtime to Apache Spark 4.1.1, Python 3.13, Scala 2.13, and AWS SDK for Java 2.x, delivering performance improvements that can help with faster job completion times and lower costs.

Second, this release reduces current AWS Glue usage rate by 30%, and when combined with the performance improvements, you may realize an even lower effective cost.

Third, AWS Glue 6.0 introduces support for more capabilities of Apache Iceberg V3. This includes the VARIANT data type with automatic shredding, deletion vectors, row lineage tracking, nanosecond timestamps, and geo types. With these capabilities, you can build modern lakehouse architectures on the latest open table format standards.

Finally, new features like Spark Declarative Pipelines, Real-Time Mode for streaming and Python virtual environments with S3 caching are designed to further improve performance and developer experience. The following sections dive deeper into each of these areas.

Runtime upgrades

AWS Glue 6.0 upgrades the core runtime stack across the board, bringing newer versions of Apache Spark, Python, Scala, and the AWS SDK to your serverless data integration workloads.

  • Apache Spark 4.1.1 – AWS Glue 6.0 runs an AWS optimized build of Apache Spark 4.1.1, a major generational leap from Spark 3.5 on AWS Glue 5.1. This release introduces improvements focused on intent-driven data engineering, real-time streaming with sub-second latencies down to single-digit milliseconds for stateless tasks, faster PySpark performance, and expanded SQL features.
  • Python 3.13 – Supports Python 3.13, a stable release that brings interpreter changes, Python data model enhancements, standard library updates, and security updates.
  • Scala 2.13 – Upgrades to Scala 2.13 which includes a collections library overhaul, language and syntax feature changes, standard library additions, and compiler performance updates.

Reduced Pricing

AWS Glue 6.0 cuts current AWS Glue pricing by 30%. This means every job you run on AWS Glue 6.0 costs 30% less per DPU-hour compared to AWS Glue 5.1, with no changes required to your workload configuration. When you combine this pricing reduction with the performance improvements delivered by runtime upgrades, your effective cost savings can compound because jobs can complete faster and consume fewer DPU-hours on a lower price point. If you run large-scale Extract, Transform, and Load (ETL) pipelines or recurring batch workloads, this compounding effect can help reduce your monthly spend.

To quantify the comparison, we ran the industry-standard TPC-DS benchmark at 3 TB scale on Parquet data stored in Amazon Simple Storage Service (Amazon S3), using 30 G.2X workers on AWS Glue. The following table compares the results we obtained in our tests for AWS Glue 6.0 and AWS Glue 5.1. Thus, based on TPC-DS benchmark at 3 TB scale, AWS Glue 6.0 delivers up to 36% better price performance than AWS Glue 5.1.

. AWS Glue 6.0 AWS Glue 5.1
Estimated Cost ($) USD 5.61 USD 8.87

Table 1: 3TB TPC-DS benchmark comparison between AWS Glue 6.0 and AWS Glue 5.1

Updated Open Table Format (OTF) support

AWS Glue 6.0 ships with updated versions of all three major open table formats – Iceberg 1.11.0, Hudi 1.1.1, and Delta Lake 4.2.0 – providing better performance, improved merge-on-read capabilities, streamlined concurrency control, and expanded SQL compatibility.

Besides supporting the latest open table format versions, AWS Glue 6.0 delivers Apache Iceberg V3 specification that is suitable for enterprise use. The highlight is Variant shredding, which AWS Glue uses to automatically decompose semi-structured data into physically optimized, columnar sub-fields, which should result in faster query read performance. Combined with deletion vectors for efficient row-level updates, UNKNOWN column types, default column values, and richer data type support, AWS Glue 6.0 is designed to make your open data lakes faster, more flexible, and more cost-efficient. AWS Glue 6.0 also adds support for geospatial data types (Geometry and Geography) and nanosecond-precision timestamps from the Apache Iceberg V3 specification, neither of which are currently supported in open-source Apache Spark 4.1. Additional features like row lineage tracking round out the Apache Iceberg V3 capabilities available on AWS Glue 6.0.

In the following sections, we illustrate select capabilities from Apache Iceberg V3 specification on AWS Glue 6.0.

  1. VARIANT column type

Apache Iceberg V3 introduces the Variant type to store semi-structured data (think JSON, XML, logs, and deeply nested event data) in a compact binary format. Variant shredding is designed to automatically decompose VARIANT columns into physically optimized, columnar sub-fields, facilitating predicate pushdowns and reducing scan overhead. It aims to provide simpler management of semi-structured data, without the need for complex flattening logic. With Variant type, you get the flexibility of embedding a JSON data type in your table columns while shredding is designed to help accelerate read queries and reduce costs.

  1. UNKNOWN column type

The UNKNOWN type in Apache Iceberg V3 acts as a flexible placeholder for columns where the data type is not yet determined at the time of table creation or data ingestion. Tables can accept all-null data initially, and the column type can be upgraded later without breaking ingestion pipelines or consuming applications. This can simplify schema evolution for rapidly changing data sources. Apache Iceberg V3’s UNKNOWN column type maps to Spark 4.1’s VOID type.

  1. DEFAULT column values

Apache Iceberg V3’s DEFAULT column values allow specifying a default value for a column in the table metadata. When you add a new column, the query engine is designed to automatically apply this default to older rows, without rewriting data or running manual backfill operations.

The following code demonstrates creating an Apache Iceberg V3 table that uses VARIANT and UNKNOWN types, and DEFAULT values for a column.

Prerequisites

To get started with this code example, make sure you have the following prerequisites.

  1. An AWS account.
  2. An AWS Identity and Access Management (IAM) role with permissions for AWS Glue, the AWS Glue Data Catalog, and Amazon S3. For more information, see Setting up IAM permissions for AWS Glue. This will be the AWS Glue job execution role.
  3. An S3 bucket to store the Iceberg table data.

Steps

To create an AWS Glue 6.0 job, use the following steps.

  1. Log in to your AWS account and open the AWS Glue console.
  2. Create a new ETL job, with Script editor option.
    1. Choose engine as Spark in the drop-down menu.
    2. Start fresh, Create script and copy-paste the following code.
    3. Replace the demo S3 bucket name with your bucket name in the code.
# Example pySpark script for testing few Iceberg v3's new data types
from pyspark.sql import SparkSession

CATALOG = "glue_catalog"
DATABASE = "sample_glue6_iceberg_db"
TABLE_NAME = "sample_glue6_table"
TABLE = f"{CATALOG}.{DATABASE}.{TABLE_NAME}"
TABLE_LOCATION = "s3://amzn-s3-demo-table-bucket/glue6blog-newdatatypes/"

# Configure Spark to use Apache Iceberg with the AWS Glue Data Catalog.
spark = (
    SparkSession.builder
    .appName("Glue6NewDataTypes")
    .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
    .config(f"spark.sql.catalog.{CATALOG}", "org.apache.iceberg.spark.SparkCatalog")
    .config(f"spark.sql.catalog.{CATALOG}.catalog-impl", "org.apache.iceberg.aws.glue.GlueCatalog")
    .config(f"spark.sql.catalog.{CATALOG}.io-impl", "org.apache.iceberg.aws.s3.S3FileIO")
    .config(f"spark.sql.catalog.{CATALOG}.warehouse", "s3://amzn-s3-demo-table-bucket/glue6blog-newdatatypes")
    .config("spark.sql.defaultColumn.enabled", "true")
    .getOrCreate()
)

spark.sql(f"CREATE DATABASE IF NOT EXISTS {CATALOG}.{DATABASE}")

# Create an Iceberg v3 table with VARIANT, unknown, and a default value.
# Spark's VOID type is stored as the Iceberg v3 unknown type.
spark.sql(
    f"""
    CREATE TABLE {TABLE} (
        record_id BIGINT,
        payload VARIANT,
        reserved_field VOID,
        status STRING DEFAULT 'active'
    )
    USING ICEBERG
    LOCATION '{TABLE_LOCATION}'
    TBLPROPERTIES ('format-version' = '3')
    """
)

# Insert two rows. The omitted columns use null and the declared default.
spark.sql(
    f"""
    INSERT INTO {TABLE} (record_id, payload)
    VALUES
    (1, parse_json('{{"event_type":"created","score":98.5}}')),
    (2, parse_json('{{"event_type":"processed","score":87.2}}'))
    """
)

# Query the row and extract values from the VARIANT column.
spark.sql(
    f"""
    SELECT
        record_id,
        variant_get(payload, '$.event_type', 'string') AS event_type,
        variant_get(payload, '$.score', 'double') AS score,
        reserved_field,
        status
    FROM {TABLE}
    """
).show(truncate=False)

spark.stop()
  1. Provide the following details in the Job details tab.
    1. A Name for the job.
    2. The IAM role you have from Prerequisites (2) for the IAM role of the job.
    3. Choose Glue 6.0 for the Glue version.
    4. Leave the rest as defaults.The following screenshot shows the Job details tab with illustrated values in the AWS Glue console.

      AWS Glue Job details tab with the Glue version set to Glue 6.0 and other settings left as defaults

      Figure 1: Job details tab with the Glue version set to Glue 6.0

    5. Scroll down. Under Advanced properties, for Job parameters, add the following additional Job parameter key-value pair:--datalake-formats=icebergThe following screenshot shows the Job parameters with the illustrated key-value pair in the AWS Glue console.

      Advanced properties section showing the Job parameters key –datalake-formats set to the value iceberg

      Figure 2: Job parameters with the datalake-formats key set to iceberg

  2. Save the job and choose Run.
  3. After the job is completed successfully, from the Runs tab – Run details, you can inspect the Output logs that take you to the logs in the Amazon CloudWatch console. The following shows the sample output for the SELECT query in the script.
+---------+----------+-----+--------------+------+
|record_id|event_type|score|reserved_field|status|
+---------+----------+-----+--------------+------+
|1        |created   |98.5 |NULL          |active|
|2        |processed |87.2 |NULL          |active|
+---------+----------+-----+--------------+------+

Notice that we inserted two rows with values only in the record_id and the variant column. Variant column inserts were done using parse_json(). The reserved_field is of VOID type, hence returns NULL values. The status column is declared with a default active value and returns active, since the column was omitted during the insert operation.

  1. Deletion Vectors
    Apache Iceberg V3 replaces the traditional positional delete files used in Apache Iceberg V2 to deletion vectors. This change can help improve Merge-on-Read (MoR) performance. This shift replaces heavy, multi-file Parquet reads with highly compressed, direct binary bitmaps that can provide lower storage overhead and faster reads on delete-heavy tables. In scenarios with heavy table updates, such as streaming change data capture (CDC) from operational databases, Apache Iceberg V3 can offer read performance advantage over Apache Iceberg V2.

    To validate the performance of deletion vectors, we created two identical AWS Glue streaming jobs and ingested the events into two different Iceberg tables, one in Apache Iceberg V2 and another in Apache Iceberg V3 format. The streaming CDC events were approximately 150,000 events per second, merge-on-read, update-heavy. Every micro-batch writes row-level deletes. We froze both tables at the same delete-heavy state and disabled compaction, leaving the tables with roughly 1.7 million rows in valid state out of the 26.4 million physical rows. The following table summarizes the read performance latency of the two Iceberg tables. We observed in this testing that reading from the delete-heavy Apache Iceberg V3 is at least 1.5 times faster than the reading from a similar Apache Iceberg V2 table. For larger enterprise scale Apache Iceberg V3 tables, the read performance could improve further.

Read latency comparison showing Apache Iceberg V3 deletion vectors reading at least 1.5 times faster than Apache Iceberg V2 delete files

Table 2 – Read latency comparison between Apache Iceberg V2 delete files and Apache Iceberg V3 deletion vectors

New ETL features

AWS Glue 6.0 introduces several additional capabilities designed to simplify how you build and manage data pipelines, some of which are discussed in the following list.

  • Spark Declarative Pipelines (SDP) where you define the outcomes you want for your entire data pipelines in a declarative fashion with SQL statements or Python decorators while AWS Glue handles execution flow, dependency resolution, parallelism, checkpointing, CDC, and recovery automatically. This helps you focus on business logic rather than orchestration plumbing.
  • Real-Time Mode (RTM) for streaming delivers continuous execution for Structured Streaming with sub-second latencies, down to single-digit milliseconds for stateless tasks. This can help support real-time use cases like fraud detection, live dashboards, and event-driven architectures without managing dedicated streaming infrastructure.
  • Arrow-Native UDFs/UDTFs execute Python functions directly on PyArrow batches without Pandas conversion overhead, which can result in faster performance for custom transformation logic at scale.
  • Recursive Common Table Expressions (CTE) adds WITH RECURSIVE queries natively, allowing graph traversals and hierarchical queries without workarounds or external libraries.
  • Python data source filter pushdown evaluates filters at the data source designed to minimize data movement, reduce the volume of data scanned, and improve job performance.
  • With Python virtual environments and S3 caching, you can provide pre-built Python dependencies, which should result in reduced startup latency for AWS Glue jobs by eliminating runtime dependency resolution. For existing jobs that use --additional-python-modules, no action is required. AWS Glue automatically handles the conversion to virtual environments when your job runs on AWS Glue 6.0.

Dependent library upgrades

The following table summarizes the key runtime and library version upgrades on AWS Glue 6.0.

Feature Glue 6.0 Glue 5.1
Spark 4.1.1-amzn-0 3.5.6-amzn-1
Python 3.13.14 3.11.15
Scala 2.13.17 2.12.18
AWS SDK for Java 2.44.6
(Version 1.x removed)
2.35.5
Boto3 1.42.84 1.40.61
Java 17.0.20 17.0.19
Iceberg 1.11.0 1.10.0
Hudi 1.1.1 .0.2
Delta Lake 4.2.0 3.3.2
AWS Glue Data Catalog client 4.11.0 4.9.0
EMR DynamoDB connector 6.1.0 5.7.0
Arrow 18.3.0 2.0.1
Hive 2.3.10-amzn-1 2.3.9-amzn-4

Table 3: Runtime and library version comparison across AWS Glue 6.0 and AWS Glue 5.1

Getting started

To get started with AWS Glue 6.0, you can use one of the following methods.  

Clean up

To avoid incurring costs, clean up the resources you created for this post.

  1. Delete the Data Catalog database and the Iceberg table.
  2. Delete the data and metadata folders of the Iceberg table from your S3 bucket.
  3. Delete the AWS Glue job and the AWS Glue job execution IAM role.

Conclusion

AWS Glue 6.0 is designed to be faster, more cost-effective, and easier to use for building your open data lakehouse architectures and orchestrating your data pipelines. In this post, we discussed the key highlights of AWS Glue 6.0 and illustrated usage of Apache Iceberg V3 features with code samples. You can create new AWS Glue jobs on AWS Glue 6.0 or migrate your existing AWS Glue jobs to benefit from these improvements.

With Apache Spark 4.1.1, Apache Iceberg V3, Python 3.13, upgraded open table format libraries, and new streaming capabilities, AWS Glue 6.0 aims to help you build new data applications or to operate your existing data pipelines more efficiently and with less maintenance overhead.

We encourage you to test AWS Glue 6.0 in your development environment today. Check out this blog that talks about upgrading your AWS Glue jobs to AWS Glue 6.0. Also, in the coming days and weeks, look out for blogs on individual topics illustrating various features of Spark 4.1.1 and Apache Iceberg V3 on AWS Glue 6.0.

Acknowledgements: We thank the numerous engineers and leaders who helped build AWS Glue 6.0 to support customers with a highly performant Spark runtime and other value-added capabilities.


About the authors

Aarthi Srinivasan

Aarthi Srinivasan

Aarthi is a Senior Big Data Architect working on data, analytics and GenAI topics with the worldwide specialist org at AWS. She works with AWS customers and partners to architect open data lake solutions, enhance product features, and establish best practices for data governance and analytics services adoption.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He is focused on building and scaling data integration and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics and machine learning workflows.

Angel Conde Manjon

Angel Conde Manjon

Angel is a Senior Solutions Architect at AWS where he helps partners develop businesses centered on Data and AI. He has previously worked on research related to Data Analytics and Artificial Intelligence in diverse European research projects. Angel is also an Apache Iceberg contributor.

Peter Tsai

Peter Tsai

Peter is a Software Development Engineer at AWS, where he enjoys solving challenges in the design and performance of the AWS Glue runtime. In his leisure time, he enjoys hiking and cycling.

Danylo Prozorov

Danylo Prozorov

Danylo is a Software Development Engineer at AWS Glue, where he works on building data integration and generative AI solutions for analytics customers. Outside of work, he enjoys sports, hiking, riding motorcycles, and building his overland rig.

Bo Li

Bo Li

Bo is a Senior Software Development Engineer on the AWS Glue team. He is devoted to designing and building end-to-end solutions to address customers’ data analytic and processing needs with cloud-based, data-intensive and GenAI technologies.

Kartik Panjabi

Kartik Panjabi

Kartik is a Software Development Manager on the AWS Glue team. His team builds generative AI features for the Data Integration and distributed system for data integration.

Mohit Saxena

Mohit Saxena

Mohit leads AWS Glue and AWS Data Analytics agentic AI initiatives that help customers build and operate big data applications on Apache Spark, Amazon S3, and cloud data lakes and warehouses, spanning across AWS Glue, Amazon EMR, and Amazon Athena.

Upgrade AWS Glue jobs to Glue 6.0 with AI-powered Spark upgrades

Post Syndicated from Prasad Nadig original https://aws.amazon.com/blogs/big-data/upgrade-aws-glue-jobs-to-glue-6-0-with-ai-powered-spark-upgrades/

Upgrading PySpark jobs to a new Apache Spark major version can introduce breaking changes. Removed configuration keys, stricter type casting, and Python library incompatibilities can cause runtime failures or silent behavior differences. With AWS Glue 6.0 now running Apache Spark 4.1 and Python 3.13, you need a reliable way to migrate your existing jobs while validating correctness.

In this post, we walk through upgrading a PySpark ETL job from AWS Glue 5.1 to AWS Glue 6.0. We use the generative AI upgrades for Apache Spark in the AWS Glue console. The upgrade analysis automatically identifies incompatibilities, iteratively resolves them, validates the result with data quality checks, and presents recommended changes for your review. AWS Glue 6.0 also delivers up to 36% better price performance* along with Iceberg v3, Spark Declarative Pipelines, Real-Time Mode, and Arrow-native Python UDFs.

What changes with AWS Glue 6.0

AWS Glue 6.0 runs Apache Spark 4.1, which introduces several behavioral changes from the Spark 3.5 runtime used in AWS Glue 5.1:

Behavior Spark 3.5 (AWS Glue 5.1) Spark 4.1 (AWS Glue 6.0)
ANSI SQL mode Disabled by default Enabled by default
Legacy Parquet datetime configs Supported Removed (renamed)
Python runtime 3.11 3.13

Beyond version compatibility, AWS Glue 6.0 also introduces:

  • Apache Iceberg v3 with VARIANT Shredding for efficient semi-structured data handling.
  • Spark Declarative Pipelines — agent-authorable ETL.
  • Real-Time Mode — single-digit millisecond streaming latency.
  • Arrow-native Python UDFs (PyArrow) for improved performance.
  • Built-in observability with structured metrics and enhanced Spark UI.
  • Up to 36% better price performance compared to AWS Glue 5.1*.

These runtime changes mean your existing AWS Glue jobs might encounter removed configuration keys, stricter type casting behavior, or Python package version incompatibilities when running on AWS Glue 6.0. Fixing these manually is time-consuming and error-prone. The following sections show how the generative upgrade analysis handles this automatically.

The sample job

Our example is a daily ecommerce order analytics pipeline running on AWS Glue 5.1:

What the job does:

  • Ingests 10,000 orders from Parquet files with INT96 timestamps (including pre-1900 historical dates from a legacy system migration).
  • Computes revenue metrics by casting string prices to numeric values and calculating line totals with discounts and tax.
  • Segments customers using recency, frequency, and monetary (RFM) scoring through mapInPandas with pandas and scikit-learn.
  • Writes enriched results back to Amazon Simple Storage Service (Amazon S3).

Job configuration (AWS Glue 5.1):

Glue version: 5.1
Worker type: G.1X
Workers: 10
Python modules: pandas==2.2.2, scikit-learn==1.5.0, numpy==1.26.4
Spark configs:
  spark.sql.legacy.parquet.datetimeRebaseModeInWrite=LEGACY
  spark.sql.legacy.parquet.int96RebaseModeInWrite=LEGACY
  spark.sql.parquet.datetimeRebaseModeInRead=LEGACY
  spark.sql.parquet.int96RebaseModeInRead=LEGACY

This job runs successfully on AWS Glue 5.1. The following sections walk through how the upgrade analysis identifies and resolves incompatibilities when upgrading this job to AWS Glue 6.0. Before starting, confirm you have the prerequisites in place.

Prerequisites

  • An AWS account with access to the AWS Glue console.
  • An existing AWS Glue job on version 5.1 or earlier with at least one successful run.
  • An Amazon S3 path for storing the upgrade analysis results.

Running the upgrade analysis from the console

The following steps walk through the upgrade analysis workflow using the AWS Glue console.

Step 1: Select your job

Navigate to your job in the AWS Glue Studio console. Confirm the job has a successful run history on AWS Glue 5.1 before starting the upgrade analysis.

AWS Glue Studio job run history showing a successful run on AWS Glue 5.1

Figure 1: Job run status for the job on AWS Glue 5.1

Step 2: Start the upgrade analysis

From the job’s Actions menu, select Upgrade with generative AI. Configure the following:

  • Target AWS Glue version: 6.0.
  • Results S3 path: An S3 location where the analysis stores its artifacts and recommendations.
Actions menu in AWS Glue Studio with the Upgrade with generative AI option

Figure 2: The Upgrade with generative AI option in the Actions menu

Configure the target AWS Glue version and the S3 results path, then choose Run.

Upgrade window with the target AWS Glue version set to 6.0 and an Amazon S3 results path

Figure 3: The Upgrade with generative AI window for setting the target AWS Glue version and results path

Choose Run. The analysis begins by running your job on AWS Glue 5.1 to establish a baseline. It then iteratively tests the job on AWS Glue 6.0, identifies failures, applies recommended fixes, and validates the job. If the upgrade analysis cannot resolve an incompatibility within its attempt budget, the analysis stops and reports the unresolved issue for manual review. Your original job remains unchanged.

Note: The upgrade analysis executes your job multiple times (one baseline run plus one or more validation attempts), and each run consumes Data Processing Units (DPUs). For large or long-running jobs, consider using the run configuration option to specify fewer workers or a smaller dataset to optimize analysis cost.

Step 3: Monitor progress

The console displays the analysis progressing through multiple validation attempts. Each attempt either succeeds or fails with a specific error, and the upgrade analysis uses that error signal to determine and apply the appropriate fix for the next attempt.

Upgrade analysis progress showing multiple validation attempts with success and failure states

Figure 4: Upgrade analysis progress across multiple validation attempts

What the upgrade analysis found and fixed

The analysis completed in four validation attempts, identifying and resolving three distinct incompatibilities. The upgrade uses deterministic migration rules for known config changes, and automated diagnosis for runtime or code errors.

Iteration 1: Removed Parquet legacy configuration

The analysis first sanitizes any Spark configurations that were removed in Spark 4.1. Our job used spark.sql.legacy.parquet.datetimeRebaseModeInWrite and spark.sql.legacy.parquet.int96RebaseModeInWrite, which no longer exist.

Migration rule applied: The SQL configs with the spark.sql.legacy prefix were removed in Spark 4.1. They have been renamed to their non-legacy equivalents, preserving the original values.

Recommended change:

Before:
spark.sql.legacy.parquet.datetimeRebaseModeInWrite=LEGACY
spark.sql.legacy.parquet.int96RebaseModeInWrite=LEGACY

After:
spark.sql.parquet.datetimeRebaseModeInWrite=LEGACY
spark.sql.parquet.int96RebaseModeInWrite=LEGACY

The read-side configs (datetimeRebaseModeInRead, int96RebaseModeInRead) already used the correct non-legacy names and required no changes.

However, with this fix applied, the validation run still failed because the Python module installation encountered an error on the AWS Glue 6.0 image.

Iteration 2: Python module version incompatibility

The pinned module versions (pandas==2.2.2, scikit-learn==1.5.0, numpy==1.26.4) could not be installed in the AWS Glue 6.0 Python 3.13 environment.

Error:

LAUNCH ERROR | Installation of Additional Python Modules failed

Recommended change: The upgrade analysis updated the version specifications from exact pins to minimum version constraints, allowing pip to resolve compatible versions for Python 3.13:

Before: pandas==2.2.2, scikit-learn==1.5.0, numpy==1.26.4
After:  pandas>=2.1.0, scikit-learn>=1.3.0, numpy>=1.24.0

With modules installing successfully, the job launched on AWS Glue 6.0 but encountered a runtime error.

Iteration 3: ANSI mode strict type casting

Spark 4.1 enables ANSI SQL mode by default (spark.sql.ansi.enabled=true). Our revenue calculation casts string prices to double, but approximately 1.8% of records contain non-numeric placeholder values such as “N/A”, “pending”, or “null” from the upstream system. From a business perspective, this meant 1.8% of revenue orders were silently excluded from revenue metrics. This data quality issue was invisible to the original pipeline.

On AWS Glue 5.1 (ANSI mode off), cast("N/A" as double) silently returns null. On AWS Glue 6.0 (ANSI mode on), this throws an exception:

Error:

NumberFormatException: [CAST_INVALID_INPUT] The value 'null' of the type
"STRING" cannot be cast to "DOUBLE" because it is malformed. Correct the
value as per the syntax, or change its target type. Use try_cast to
tolerate malformed input and return NULL instead. SQLSTATE: 22018

Migration rule applied: As of Spark 4.1, spark.sql.ansi.enabled is on by default. Casting a malformed value now raises CAST_INVALID_INPUT instead of returning NULL. The upgrade analysis resolved this by updating the script to use try_cast(), which safely returns NULL for malformed input while preserving ANSI mode protections for the rest of the job.

Recommended change:

Before (AWS Glue 5.1):

F.col("unit_price").cast("double")

After (AWS Glue 6.0, fixed by the upgrade analysis):

F.expr("try_cast(unit_price as double)")

This is a targeted fix that handles the known dirty data without disabling ANSI mode globally, keeping overflow detection and type safety active throughout the job.

Final validation and data quality check

After applying all three fixes, the analysis ran the job on AWS Glue 6.0 one final time and performed a data quality comparison between the AWS Glue 5.1 baseline output and the AWS Glue 6.0 output.

Result: The job completed successfully and all data validations passed with no mismatches detected between the source and target outputs.

Completed upgrade analysis status with links to the results output path in Amazon S3

Figure 5: Final analysis status with links to the results output path in Amazon S3

Reviewing the upgrade summary

The analysis produces a detailed summary stored in your S3 results path. This summary documents each validation attempt, the errors encountered, the migration rules applied, and the recommended configuration changes:

s3://amzn-s3-demo-bucket/scripts/auto-upgrade/ja-{analysis-id}/
    summary/
        summary.md                       # Full iteration-by-iteration report
        data_validation_summary.md       # Data quality comparison results
    artifact/
        attempt_N/
            script/main.py               # Recommended script (if modified)
            job_config_modifications.json  # Recommended parameter changes
            requirements.txt             # Updated dependency versions

The following is a snippet from the upgrade summary (summary.md) showing the recommended changes and validation attempt details:

The summary documents each validation attempt, the changes applied, and the data quality results.

Upgrade summary showing validation attempts, applied changes, and data quality comparison results

Figure 6: Upgrade summary snippet showing validation attempt details, data quality, and analysis results

After reviewing the recommendations, accept the changes to upgrade your job to AWS Glue 6.0. This updates your job definition with the recommended configuration, including the renamed Spark configs, updated module versions, and any script modifications. Because the analysis has already validated the job on AWS Glue 6.0 and confirmed data quality parity with the original, your job is ready for production.

After reviewing the recommendations, you can apply the upgraded script to your job.

AWS Glue Studio prompt to apply the upgraded script to the job

Figure 7: The option to apply the upgraded script to the job

Choose Apply to confirm the upgrade.

Confirmation dialog with the Apply button to upgrade the job to AWS Glue 6.0

Figure 8: The Apply button that confirms upgrading the job to AWS Glue 6.0

After applying, the job definition reflects the new AWS Glue version.

AWS Glue job details showing version 6.0 after applying the upgrade

Figure 9: The AWS Glue version for the job after applying the upgrade

Python virtual environments in AWS Glue 6.0

AWS Glue 6.0 introduces --python-virtual-env-storage-prefix, a service-managed virtual environment with S3 caching that simplifies Python dependency management.

For existing jobs that use --additional-python-modules, no action is required. AWS Glue automatically handles the conversion to virtual environments when your job runs on AWS Glue 6.0. Your jobs continue to work without any changes.

For new jobs on AWS Glue 6.0, we recommend using the virtual environment approach:

{
    "DefaultArguments": {
        "--python-virtual-env-storage-prefix": "s3://amzn-s3-demo-bucket/glue-venv-cache/",
        "--additional-python-modules": "pandas>=2.1.0,scikit-learn>=1.3.0,numpy>=1.24.0"
    }
}

How it works:

  • On the first run, AWS Glue installs your modules into a virtual environment, packages it, and caches the result to your specified S3 path (approximately 15–30 seconds of additional startup time).
  • On subsequent runs, AWS Glue downloads and extracts the cached virtual environment instead of running pip install.
  • The cache is automatically invalidated when your module list, versions, or AWS Glue version changes.

This approach provides faster cold starts after the first run, requires no Docker image management (unlike --python-virtual-env), and is entirely service-managed with no maintenance burden.

Conclusion

The generative upgrade analysis identified and resolved three distinct compatibility issues in our AWS Glue 5.1 job, so the job now runs successfully on AWS Glue 6.0 with Apache Spark 4.1:

  • The upgrade analysis renamed legacy Parquet datetime configuration keys (removed in Spark 4.1) to their current equivalents.
  • The upgrade analysis updated Python module version specifications that were incompatible with Python 3.13 to use flexible minimum version constraints.
  • The upgrade analysis addressed the new ANSI SQL mode default (which causes runtime failures on malformed data) with a targeted fix using try_cast() to safely handle non-numeric values while preserving ANSI mode protections.

The analysis validated that the upgraded job produces output consistent with the original, and presented all changes as recommendations for review before applying them to your job.

Next steps

After you have reviewed and accepted the upgrade changes, you can delete the analysis results stored in your S3 results path.

*Based on 3TB TPC-DS benchmark comparing AWS Glue 6.0 to AWS Glue 5.1.


About the authors

Prasad Nadig

Prasad Nadig

Prasad is a Senior Analytics Specialist Solutions Architect at Amazon Web Services (AWS), specializing in large-scale data analytics and AI. Prasad partners with customers to design, migrate, and modernize their analytics platforms on AWS into scalable, cost-effective solutions, with deep expertise in data lakes, data warehousing, distributed processing, and performance tuning at petabyte scale.

Shrey Malpani

Shrey Malpani

Shrey is a Senior Product Manager Technical at Amazon Web Services (AWS), where he works at the intersection of distributed data processing and data integration. He is focused on building and scaling data integration and data management capabilities across services like AWS Glue, Amazon EMR, and Amazon Redshift that help customers build AI-ready data platforms for their analytics and machine learning workflows.

Rishabh Nair

Rishabh Nair

Rishabh is a Software Development Engineer in the AWS analytics organization, where he combines generative AI with distributed systems to build agentic workflows that modernize large-scale data processing. He is passionate about the infrastructure that makes these workflows reliable and scalable for customers.

Keerthi Chadalavada

Keerthi Chadalavada

Keerthi is a Senior Software Development Engineer in the AWS analytics organization. She focuses on combining generative AI and data integration technologies to design and build comprehensive solutions for analytics and data engineering workloads.

Build a unified AI agent architecture with DynamoDB and Bedrock

Post Syndicated from Dhananjay Karanjkar original https://aws.amazon.com/blogs/architecture/build-a-unified-ai-agent-architecture-with-dynamodb-and-bedrock/

Teams building AI agents on AWS often face a fragmented data architecture: operational data lives in Amazon DynamoDB while vector embeddings for semantic search sit in a separate, purpose-built vector store. This duplication increases infrastructure cost, adds synchronization complexity, and widens the window for stale retrieval results. With the general availability of native vector search in Amazon DynamoDB (launched August 5, 2026), you can now store embeddings alongside your operational data in the same table. You query them using the SearchVectors API operation.

In this post, I show you how to build a unified AI agent architecture where an Amazon Bedrock agent uses a single DynamoDB table for both structured lookups and semantic similarity search. The agent calls AWS Lambda action groups that invoke SearchVectors for natural language retrieval and standard DynamoDB APIs for create, read, update, and delete (CRUD) operations. An Amazon DynamoDB Streams pipeline automatically generates embeddings using Amazon Titan Text Embeddings V2 whenever content changes. This keeps the vector index synchronized without manual intervention.

Use case

Consider a technical knowledge management platform where a team maintains hundreds of internal documents: runbooks, architecture decision records, and troubleshooting guides. Team members interact with a conversational agent to find relevant content (“What’s our retry strategy for payment failures?”), retrieve specific documents by ID, or update existing entries.

Without native vector search, this architecture requires a DynamoDB table for document storage plus a separate vector database (or Amazon OpenSearch Service cluster) for semantic retrieval. The Amazon DynamoDB Streams pipeline must write to both stores, and the agent must route requests to the correct backend. With DynamoDB vector search, you collapse this into a single table and reduce operational overhead.

Solution overview

This solution uses a single-table design in DynamoDB that serves two access patterns: key-value lookups for operational data and approximate nearest neighbor (ANN) search for semantic queries. A Bedrock agent orchestrates user interactions and routes requests to the appropriate action group function.

The following list summarizes the core components:

  • DynamoDB table with vector index stores documents, metadata, and 1,024-dimension embeddings in one place.
  • Bedrock agent handles conversation orchestration, tool selection, and response synthesis.
  • Action group Lambda executes semantic search (using SearchVectors) and CRUD operations against the same table.
  • Embedding pipeline Lambda (triggered by DynamoDB Streams) generates embeddings for new or modified content using Amazon Titan Text Embeddings V2.

Architecture

The following diagram illustrates the data flow through the unified architecture.

Architecture diagram showing a user query flowing to an Amazon Bedrock agent, which invokes action group Lambda functions that call the DynamoDB SearchVectors API and standard CRUD APIs, with DynamoDB Streams triggering an embedding pipeline Lambda that generates vectors with Amazon Titan Text Embeddings V2

Figure 1: Unified AI agent architecture using DynamoDB vector search and Amazon Bedrock

The numbered steps describe the data and request flow:

  1. A user sends a natural language query to the Bedrock agent.
  2. The agent analyzes the request and invokes the appropriate action group Lambda function.
  3. For semantic search, the action group Lambda generates a query embedding using Amazon Titan Text Embeddings V2.
  4. The Lambda function calls the DynamoDB SearchVectors API (or standard CRUD APIs for operational lookups) against the single table with vector index.
  5. When new content is written to the table, DynamoDB Streams captures the change.
  6. DynamoDB Streams triggers the embedding pipeline Lambda.
  7. The embedding pipeline Lambda calls Amazon Titan Text Embeddings V2 to generate a vector for the new content and writes it back to the same DynamoDB item, where the vector index automatically indexes it.

Prerequisites

To implement this architecture in your account, you need the following:

  • An AWS account with permissions to create DynamoDB tables, Lambda functions, Bedrock agents, and IAM roles.
  • DynamoDB Streams enabled on the table with StreamViewType set to NEW_AND_OLD_IMAGES (the embedding pipeline compares old and new content to prevent a write loop).
  • Access to the Amazon Titan Text Embeddings V2 model (amazon.titan-embed-text-v2:0) enabled in Amazon Bedrock model access.
  • Access to an Anthropic Claude or Amazon Nova model for the Bedrock agent foundation model (check model support by Region).
  • Python 3.12 or later (for Lambda function code).

Implementation

This section walks through the key components of the architecture.

Designing the single-table schema

The table uses a composite primary key (entity_id as partition key, sk as sort key) and stores embeddings as a list of numbers:

# Table schema overview
# PK: entity_id (S) - unique document identifier
# SK: sk (S) - sort key for item versioning
# Attributes: title, content, category, metadata, embedding (L of N)

The vector index partitions search results by the category attribute. Choose a partition key with moderate cardinality that matches your query patterns. A very low-cardinality key (a handful of values) concentrates data in few partitions and limits throughput scaling, while a unique-per-item key leaves no neighbors to compare. For multi-tenant workloads, tenant_id is usually the right partition key. For more information, refer to the DynamoDB vector search best practices.

The following AWS Command Line Interface (AWS CLI) command creates the vector index on an existing table:

aws dynamodb update-table \
    --table-name unified-agent-data \
    --stream-specification StreamEnabled=true,StreamViewType=NEW_AND_OLD_IMAGES \
    --attribute-definitions \
        AttributeName=category,AttributeType=S \
    --vector-index-updates \
    '[{"Create": {
        "IndexName": "content-embedding-index",
        "VectorAttribute": {"AttributeName": "embedding"},
        "Dimensions": 1024,
        "DistanceFunction": "COSINE",
        "SearchSchema": [
            {"AttributeName": "category", "SearchSchemaElementType": "HASH"}
        ],
        "Projection": {"ProjectionType": "INCLUDE", "NonKeyAttributes": ["title", "category"]}
    }}]'

After creating the index, wait for it to become searchable. Poll DescribeTable until IndexStatus is ACTIVE and Backfilling is no longer true. The first few searches after the index reports ACTIVE can still return ValidationException because SearchVectors is served by a dedicated search endpoint. Treat these as retryable rather than as a failure.

aws dynamodb describe-table --table-name unified-agent-data \
    --query 'Table.VectorIndexes[?IndexName==`content-embedding-index`].[IndexStatus,Backfilling]'

Key constraints to keep in mind:

  • DynamoDB vector indexes require on-demand capacity mode (provisioned mode isn’t supported).
  • Maximum five vector indexes per table, with up to 4,096 dimensions each.
  • The SearchSchema HASH attribute is mandatory in every SearchConditionExpression.
  • Only equality operators are supported in search conditions.
  • SearchVectors responses are limited to 16 MB and don’t support pagination. Project only the attributes you need and keep TopK modest to stay within this limit.
  • Items missing the SearchSchema HASH attribute (category in this example) are silently excluded from the vector index while remaining in the base table.

Building the action group Lambda

The action group Lambda handles both semantic search and operational lookups. The agent invokes it with a function name and parameters based on the tool definition.

The semantic search function generates a query embedding and calls SearchVectors. This index uses COSINE distance, where lower scores indicate greater similarity. Name the field accordingly so the agent doesn’t invert the ranking:

def semantic_search(query: str, category: str, max_results: int = 5):
    embedding = generate_embedding(query)
    results = dynamodb.search_vectors(
        TableName=TABLE_NAME,
        IndexName=INDEX_NAME,
        SearchVector=[{"N": str(v)} for v in embedding],
        TopK=min(max_results, 100),
        SearchConditionExpression="category = :cat",
        ExpressionAttributeValues={":cat": {"S": category}},
    )
    return [
        {"entity_id": r["Item"]["entity_id"]["S"],
         "title": r["Item"].get("title", {}).get("S", ""),
         "distance": r["Score"]}  # COSINE: lower = more similar
        for r in results.get("SearchResults", [])
    ]

The generate_embedding helper calls Amazon Titan Text Embeddings V2:

def generate_embedding(text: str) -> list[float]:
    response = bedrock_runtime.invoke_model(
        modelId="amazon.titan-embed-text-v2:0",
        body=json.dumps({
            "inputText": text,
            "dimensions": 1024,
            "normalize": True
        }),
    )
    return json.loads(response["body"].read())["embedding"]

The Lambda handler routes requests based on the function name passed by the Bedrock agent:

def handler(event, context):
    function = event.get("function")
    parameters = {p["name"]: p["value"] for p in event.get("parameters", [])}
    if function == "semantic_search":
        result = semantic_search(parameters["query"], parameters["category"])
        body = json.dumps({"results": result})
    elif function == "get_item_details":
        body = json.dumps(get_item_details(parameters["entity_id"]))
    else:
        body = json.dumps({"error": f"Unknown function: {function}"})
    return {
        "messageVersion": "1.0",
        "response": {
            "actionGroup": event["actionGroup"],
            "function": function,
            "functionResponse": {"responseBody": {"TEXT": {"body": body}}}
        }
    }

Automating embeddings with DynamoDB Streams

The embedding pipeline Lambda triggers on INSERT and MODIFY events. It generates an embedding for new or changed content and writes it back to the same item:

def handler(event, context):
    for record in event["Records"]:
        if record["eventName"] not in ("INSERT", "MODIFY"):
            continue
        new_image = record["dynamodb"]["NewImage"]
        old_image = record["dynamodb"].get("OldImage", {})
        content = new_image.get("content", {}).get("S")
        if not content:
            continue
        # Prevent infinite loop: skip if content hasn't changed
        if "embedding" in new_image and old_image.get("content") == new_image.get("content"):
            continue
        embedding = generate_embedding(content)
        dynamodb.update_item(
            TableName=TABLE_NAME,
            Key={"entity_id": new_image["entity_id"], "sk": new_image["sk"]},
            UpdateExpression="SET embedding = :emb",
            ExpressionAttributeValues={
                ":emb": {"L": [{"N": str(v)} for v in embedding]}
            },
        )

The infinite-loop guard is critical. Without it, the Lambda writes back an embedding, which triggers another Streams event, which triggers another embedding generation, and so on. The check compares the content field between old and new images, skipping processing when only the embedding attribute changed. This guard requires StreamViewType = NEW_AND_OLD_IMAGES. Without it, OldImage is empty and the guard never fires.

For production use, configure the event source mapping with ReportBatchItemFailures so that only failed records are retried. Add an Amazon Simple Queue Service (Amazon SQS) dead-letter queue (or on-failure destination) for records that repeatedly fail. Retry Amazon Bedrock InvokeModel calls with exponential backoff to handle throttling.

Defining the agent tool schema

The Bedrock agent needs a function schema that describes the available tools. This tells the agent when and how to call each function:

{
    "functions": [
        {
            "name": "semantic_search",
            "description": "Search documents by meaning using natural language. Returns results ranked by COSINE distance (lower = more similar).",
            "parameters": {
                "query": {"type": "string", "required": true,
                          "description": "Natural language search query"},
                "category": {"type": "string", "required": true,
                             "description": "Document category to search within"}
            }
        },
        {
            "name": "get_item_details",
            "description": "Retrieve a specific document by its unique ID.",
            "parameters": {
                "entity_id": {"type": "string", "required": true,
                              "description": "Unique document identifier"}
            }
        }
    ]
}

When to use this pattern

This unified architecture works best when your application already uses DynamoDB as its primary operational store and you want to add semantic search without managing a separate service. Consider the following decision points:

  • Use this pattern when your application meets these conditions:
    • Documents update frequently and must be immediately searchable.
    • Your dataset fits within the DynamoDB vector index constraints.
    • You want to minimize infrastructure components.
  • Use Amazon Bedrock Knowledge Bases when your source data lives in Amazon Simple Storage Service (Amazon S3), you need managed chunking and ingestion, or you don’t need real-time index updates tied to operational writes.
  • Use Amazon OpenSearch Service when you need advanced search features (range filters, aggregations, faceted search), your queries require more than equality-based filtering, or you need results beyond the 100-item TopK limit.

Security considerations

The following list highlights the key security aspects of this architecture:

  • Least-privilege IAM policies: Scope dynamodb:SearchVectors to the specific index ARN (arn:aws:dynamodb:{region}:{account}:table/{table}/index/{index}). The embedding Lambda needs only dynamodb:UpdateItem, not search permissions.
  • No fine-grained access control for SearchVectors: DynamoDB condition keys like dynamodb:LeadingKeys don’t apply to the SearchVectors API. For multi-tenant workloads, use the SearchSchema HASH partition key to scope queries by tenant, or use separate tables for strict isolation.
  • Encryption at rest: DynamoDB encrypts data including vector embeddings using your choice of AWS owned keys, AWS managed keys, or customer managed keys through AWS Key Management Service (AWS KMS).
  • Transport encryption: All SearchVectors traffic uses TLS. The API routes to a dedicated search endpoint that the AWS SDKs handle automatically.
  • Bedrock model access: Restrict bedrock:InvokeModel permissions to the specific embedding and agent foundation model ARNs required by the solution.
  • Agent-to-Lambda invocation: Grant lambda:InvokeFunction to bedrock.amazonaws.com on the action group Lambda, scoped with an aws:SourceArn condition matching the agent ARN. Without this resource-based policy, the agent can’t invoke the action group.

Clean up

To avoid ongoing charges, delete the resources in the following order:

  1. Delete the Bedrock agent and its action group.
  2. Delete the embedding pipeline Lambda function and its event source mapping.
  3. Delete the DynamoDB table (this also removes the vector index). If you want to keep the table but remove the vector index, run the following command first:
    aws dynamodb update-table \
        --table-name unified-agent-data \
        --vector-index-updates '[{"Delete": {"IndexName": "content-embedding-index"}}]'

  4. Delete the action group Lambda function and associated IAM roles.

Conclusion

With this pattern, you can build a unified AI agent architecture that uses a single DynamoDB table for both operational data and vector-based semantic search. The native vector search of DynamoDB combined with Bedrock agent action groups eliminates the need for a separate vector database. DynamoDB Streams-driven embedding generation keeps the index synchronized in real time.

This pattern reduces infrastructure complexity for applications that already rely on DynamoDB and need to add conversational AI capabilities. The automatic embedding pipeline keeps your vector index synchronized with operational writes, and the action group design gives the agent access to both semantic and structured query paths.

Adapt the table schema, embedding dimensions, and agent instructions to your domain. Clone the sample-dynamodb-vector-search-architecture repository to deploy the complete working implementation. For more information about DynamoDB vector search capabilities and limits, refer to the Amazon DynamoDB vector search documentation.

References

About the author

How AgentFlo built AI sales agents with Amazon Bedrock AgentCore – Part 2

Post Syndicated from Muhammad Musab Iqbal original https://aws.amazon.com/blogs/architecture/how-agentflo-built-ai-sales-agents-with-amazon-bedrock-agentcore-part-2/

If you’re building AI agents for commerce at scale, you face two critical challenges: handling unpredictable traffic spikes and ensuring your agents can be trusted with real customer transactions.

This post shows how AgentFlo solved these challenges using Amazon Bedrock AgentCore and AWS serverless architecture. You learn the architectural patterns behind their reliability and trust frameworks, see the measurable business results (including +12% net revenue uplift based on early deployment data), and explore their roadmap for voice agents and server-side tool execution.

This is Part 2 of a two-part series. Part 1 covers velocity, standardization, and scalability.

Pillar 4: Trust: guardrails for autonomous commercial action and real-time visibility into agent operations

AgentFlo enforces trust at every layer of the stack, from pre-request filtering to post-response privacy controls, so merchants can deploy autonomous agents with confidence.

The challenge

Enterprise customers won’t deploy autonomous agents unless they can trust them. An agent in production can’t expose sensitive data, offer unauthorized discounts, or access another customer’s information. The system also prevents price hallucination, unauthorized tool calls, opt-out violations, and credential exposure.

Merchants also need fine-grained control over who can interact with their AI agents and what data each segment can access. Enterprise customers require restricted access; B2C businesses need open access for broader reach. Without identity-based controls, deploying customer-facing AI is a non-starter.

Defense in depth

In AgentFlo, trust isn’t only about safe responses. It’s about safe action. Agents can create carts, place orders, apply discounts, access customer data, and interact with backend systems, so policy enforcement must sit outside the model’s reasoning loop. The model proposes. Deterministic policy decides.

AgentFlo applies trust controls across the full agent lifecycle: before the model sees the request, during tool execution, and after the model generates a response.

Three-layer guardrails

When users deploy an agent from the AgentFlo Portal, security enforcement happens at three stages.

First, the AWS Fargate layer detects prompt injection and handles opt-outs before requests reach the agent. WhatsApp messages are authenticated using phone numbers as unique identifiers. Enterprise customers like EBM restrict access to authorized users, while restaurant deployments stay open for broader reach.

Next, the AgentCore layer verifies identity and enforces order locks during tool execution. AgentCore Gateway, a capability of Amazon Bedrock AgentCore, enforces policies that prevent sales agents from accessing customer support tools. Cedar policies enforce business rules like maximum discount percentages independently of the model’s reasoning. Cedar is an open-source policy language developed by AWS for fine-grained, verifiable authorization decisions. Additionally, Policy in Amazon Bedrock AgentCore integrates with Amazon Bedrock Guardrails, so Cedar policies can invoke configurable safeguards for prompt attack detection, content filtering, and sensitive information blocking directly at the gateway boundary.

Finally, post-turn privacy filters screen outputs to block inadvertent token disclosure and unverified price claims before customers see responses. Secrets are managed through AWS Secrets Manager with OIDC (OpenID Connect)-authenticated continuous integration and continuous delivery (CI/CD) pipelines. No credentials are stored in agent code.

Infrastructure security

Trust at the application layer requires sound underlying infrastructure. AgentFlo combines the built-in isolation of AgentCore with application-level controls:

  • Session isolation: Session isolation through AgentCore runtime, a capability of Amazon Bedrock AgentCore, which provides complete separation between merchants’ agent sessions through dedicated microVMs.
  • AgentCore Gateway policies: Fine-grained Cedar policies control which agents can access which tools and data, enforced deterministically regardless of model reasoning.
  • AWS Identity and Access Management (IAM)-based access control: Fine-grained permissions for agent-to-service communication.
  • Amazon Virtual Private Cloud (Amazon VPC) integration: Agent sessions operate within AgentFlo’s VPC with domain-level network restrictions, so agents only communicate with approved endpoints.
  • Compliance: Data residency controls and audit trails for regulatory requirements across multiple jurisdictions.

Observability

Trust requires visibility. Merchants need to see what their agents are doing in real time, not only after something goes wrong. AgentFlo uses AgentCore Observability, a capability of Amazon Bedrock AgentCore, to provide end-to-end tracing of every agent interaction, from initial request through tool execution to final response.

AgentCore Observability captures structured traces for each agent turn, including model latency, tool invocation sequences, token usage, and error rates. These traces flow into Amazon CloudWatch, where AgentFlo builds dashboards showing active sessions, response times, and tool call patterns. Full request-to-response traces enable trace-level debugging. The system tracks P50/P95 latencies and throughput across agent types, alerting merchants when behavior deviates from baselines. Cost attribution provides per-merchant, per-agent breakdowns tied to specific conversations.

Results

Together, observability and trust give enterprises confidence to deploy autonomous agents at scale. Merchants benefit from safer execution, stronger compliance across jurisdictions, and controlled access to tools and data at the session level.

Pillar 5: Reliability: a data foundation that keeps agents grounded in reliable data

Reliable agents need reliable data. AgentFlo grounds every agent action in verified, current information through stateful sessions, merchant knowledge bases, and semantic product discovery.

The challenge

Enterprise customers won’t trust AI agents that forget context, hallucinate product details, or operate on stale data. An agent that quotes the wrong price, forgets a customer’s earlier request, or recommends discontinued products destroys confidence instantly. The system must make sure every agent action is grounded in verified, current information, from product catalogs and pricing to conversation history and business rules.

Merchants also need their agents to maintain continuity across long customer journeys, access up-to-date business-specific knowledge, and surface products through natural language, all without manual intervention or prompt engineering.

Data architecture overview

In AgentFlo, reliability isn’t only about accurate responses. It’s about accurate action grounded in verified data. Agents retrieve product information, manage carts, and complete transactions, so every data source must be authoritative and current. The model reasons. Structured data decides.

AgentFlo applies data reliability controls across three layers: stateful conversation management through Amazon DynamoDB, merchant-specific knowledge through Amazon Bedrock Knowledge Bases, and semantic product discovery through vector embeddings in Amazon S3 Vector.

State management architecture

A real sales journey can span 8 hours or 3 days. A customer might ask about a product in the morning, compare options at lunch, and complete the purchase that evening. The agent must remember context across all turns and maintain state, so customers don’t need to start over.

AgentCore runtime and Amazon DynamoDB manage this context storage. The agent replays relevant history, loads context based on intent, and continues transactions safely.

Each agent needs to be stateful (remembering earlier interactions), autonomous (deciding next steps without human intervention), and safe (operating within business rules without exposing sensitive data).

The two-table DynamoDB design covers all three requirements: session continuity through conversation replay, autonomous context loading based on detected intent, and data integrity through structured ground-truth storage that the model can’t hallucinate over.

Diagram of the per-message conversation flow between AgentCore runtime and the DynamoDB session and cart tables

Figure 1: Per-message conversation flow. AgentCore runtime loads the last 15 messages from the DynamoDB Session Table at the start of each turn. The Cart Table is loaded on demand only when intent detection flags the request as cart-related, preventing the model from generating incorrect prices and quantities.

Knowledge base system

Merchants upload business-specific data (restaurant menus, clinic policies, product specifications, promotion calendars) into Amazon Bedrock Knowledge Bases backed by Amazon Simple Storage Service (Amazon S3). Agents automatically retrieve and reason over this merchant-specific content. Responses stay grounded in accurate, up-to-date business information without merchants writing a single prompt.

Semantic search and vector retrieval

AgentFlo improves product discovery by giving every product in its Amazon Aurora database a lightweight vector embedding. Customers can find products by name or description.

AgentFlo also generates extra searchable tags for each product automatically, and merchants can add their own. The result: phrases like “the pink one,” “the smallest one,” “the new one,” or “the chocolate with the golden wrapper” all map to the right product. Customers ask for things naturally, and the platform finds what they’re looking for.

Observability and billing

The system captures all message interactions through Amazon Data Firehose to Amazon S3, so merchants can track cost per conversation and compare those costs against sales revenue. This pipeline shows merchants the return on investment (ROI) of their agent deployments and provides the data foundation for continuous agent improvement.

Results

The data architecture helps minimize context loss in conversations across multi-day customer journeys, with grounded responses that eliminate price and product hallucination. Merchants benefit from natural language product discovery without keyword dependency and merchant-specific knowledge retrieval without prompt engineering. Full cost visibility and ROI attribution per agent deployment give merchants clear measurement of platform value.

Business impact: measurable results across the customer lifecycle

Salesflo’s solution, powered by Strands Agents SDK and Amazon Bedrock AgentCore, delivered measurable improvements across the customer lifecycle:

Metric Improvement
Net revenue uplift +12%
Customer engagement +40%
Conversion rate +15%
Average order value +8%
Customer reactivation +20%

AgentFlo generated these results by comparing agent-assisted customer journeys against a control group over a 90-day early deployment period.

Note: The foundation models referenced in this post are available in select AWS Regions. For the latest information on model availability, see Supported Regions and models for Amazon Bedrock.

The compounding effect at scale

At AgentFlo’s scale, the impact is significant. With $300 billion in annual transacted value flowing through the Salesflo solution (based on platform transaction data), even single-digit percentage improvements translate to billions in incremental revenue for merchants.

Operational improvements

Beyond the metrics, merchants see operational improvements:

  • 24/7 coverage — AI agents engage customers at any hour, in any time zone, across any channel.
  • Consistent quality — Every customer interaction follows best-practice selling methodologies without the variability of human agents.
  • Scalable personalization — Thousands of concurrent agent sessions, each maintaining unique context per customer.
  • Rapid merchant onboarding — New merchants go live with customized AI agents in days, not months, through the self-service configuration platform.

What’s next: voice, server-side execution, and integration expansion

AgentFlo is actively extending the platform along three directions, each at a different stage of maturity.

Real-time voice agents (in pilot)

AgentFlo already supports voice within WhatsApp. Incoming voice notes go through a two-pass transcription process: a raw initial pass, followed by domain-aware correction that fuzzy-matches text against the live product catalog. The second pass catches brand names and SKUs even when partially misheard, across more than 90 supported languages. For outbound audio, merchants pick from multiple Text-to-Speech (TTS) providers per deployment.

Because Speech-to-Text (STT), reasoning on AgentCore runtime, and TTS are fully independent components, any one can be swapped without disrupting the rest of the pipeline.

The next step: BidiAgent. The next evolution is real-time voice agents built on the Strands SDK BidiAgent and Amazon Bedrock AgentCore WebRTC support. BidiAgent supports bidirectional audio streaming, natural interruptions, and concurrent tool execution. The agent can check inventory or apply a discount while continuing to listen and respond to the customer in the same call.

The AgentCore WebRTC protocol and Amazon Kinesis Video Streams handle peer-to-peer transport for mobile and browser interactions without requiring relay infrastructure. This pushes AgentFlo beyond text messaging into proactive outbound calls and high-value B2B sales, where real-time conversation is essential for building trust.

Server-side tool execution (under development)

AgentFlo is experimenting with server-side tool execution in Amazon Bedrock, which removes client-side orchestration entirely.

Traditional orchestration loops between model and tools repeatedly (model → execute → send result → repeat). With server-side execution, the agent makes a single API call to the Amazon Bedrock Responses API with an AgentCore Gateway Amazon Resource Name (ARN). The model then autonomously discovers, invokes, and processes tools through the Gateway Model Context Protocol (MCP) interface, all inside AWS infrastructure with no roundtrips back to the client.

Here’s what that one API call looks like:

from openai import OpenAI

# OPENAI_BASE_URL = https://bedrock-mantle.us-west-2.api.aws/v1
# OPENAI_API_KEY = <Amazon Bedrock API key>
client = OpenAI()

response = client.responses.create(
    model="openai.gpt-oss-120b",
    stream=True,
    background=False,
    store=False,
    input=[
        {
            "type": "message",
            "role": "user",
            "content": [{"type": "input_text", "text": user_message}],
        }
    ],
    tools=[
        {
            "type": "mcp",
            "server_label": "agentflo_gateway",
            "connector_id": GATEWAY_ARN,  # arn:aws:bedrock-agentcore:...:gateway/...
            "server_description": "AgentFlo commerce tools (cart, catalog, knowledge base)",
            "require_approval": "never",
        },
    ],
)

Note: Model availability varies by AWS Region. The model and endpoint shown in this example may not be available in all Regions. See Supported Regions and models for Amazon Bedrock for current availability.

One request handles the whole loop. The Gateway ARN goes in as an MCP connector, and Bedrock takes it from there; pulling the tool list, picking the right one, invoking it, and feeding the result back to the model without anything leaving AWS. The client never sees credentials, tool schemas, or the intermediate turns. See ShopAssist: E-Commerce Agent Demo for more details.

Early results: For specialist agents with short, focused tool loops, early measurements show approximately 30% lower latency. A sales agent calling three tools in sequence (check inventory, apply discount, update cart) collapses 50+ lines of orchestration into a single API call. All credentials stay server-side, and the simplified architecture makes onboarding new agent developers faster.

Integration ecosystem expansion (ongoing)

AgentFlo adds integrations based on merchant demand. Each new platform (payment processors, shipping providers, loyalty systems, or vertical-specific ERPs) becomes another MCP server connector in AgentCore Gateway. This pattern keeps expansion modular and quick.

Key takeaways

  1. Pick your pillars first. The right AWS stack follows. AgentFlo first defined what velocity, standardization, scalability, and trust meant for the production system. With clear requirements, the AWS stack choices became obvious: Strands for the agent layer, AgentCore runtime for stateful sessions, AgentCore Gateway for tool routing and policy, and AWS Fargate for message ingestion.
  2. Specialized agent recipes beat multi-agent complexity. Deploying a single, well-configured agent with domain-specific tools and knowledge outperforms multi-agent orchestration for most customer interactions. Multi-agent handoffs are reserved for cross-domain transitions where context boundaries are clearly defined.
  3. Commerce is stateful. Plan for that on day one. A real sales journey can span 8 hours or 3 days. Stateless chatbots can’t maintain context that long. AgentCore runtime stateful sessions and DynamoDB-backed context were chosen on day one for exactly this reason. Retrofitting state onto a stateless agent later is much more painful than designing for it up front.
  4. Three-layer security builds trust. Pre-turn guards on Fargate, per-tool guards through AgentCore Gateway policies with Cedar, and post-turn output filters work together to support safe deployment of autonomous agents handling commercial transactions.
  5. System design supports scale. By building agent customization as a software as a service (SaaS) layer on top of AgentCore infrastructure, AgentFlo serves hundreds of merchants on shared infrastructure while still delivering personalized agent behavior for each one.
  6. Serverless + AgentCore = elastic commerce. The combination of Fargate for message handling and AgentCore for agent execution means AgentFlo scales from normal traffic to 50x flash-sale spikes without pre-provisioning or capacity planning.
  7. The feedback loop is the product. Real customer conversations teach the agents how to close. Where customers hesitate, what language converts, when they want a human handoff — all of it feeds back into recipes, prompts, and tool definitions. A new merchant deployment benefits from every conversation that ran on the platform before it. That compounding loop is harder to copy than any single piece of the architecture.

Conclusion

Building production-scale AI agents for commerce requires scalability and trust from day one. By combining Amazon Bedrock AgentCore stateful sessions and microVM isolation with AWS serverless infrastructure, AgentFlo delivers autonomous agents that handle unpredictable traffic while maintaining the security controls enterprises require.

Next steps

For questions about implementing similar architectures, visit the AWS Architecture Center or contact your AWS account team. To start building, open the Amazon Bedrock console or explore the Amazon Bedrock service detail page.

We’d love to hear how you’re building agentic AI systems. Share your experiences in the comments.


About the authors

Amazon Linux default SSM parameter will now track the latest kernel

Post Syndicated from Gokul Govindaraju original https://aws.amazon.com/blogs/compute/amazon-linux-default-ssm-parameter-will-now-track-the-latest-kernel/

Today we are announcing that the Amazon Linux kernel-default AWS Systems Manager (SSM) parameter will now update to point to the latest Amazon Linux kernel version as new kernel versions get released. On August 17, 2026, for Amazon Linux 2023 (AL2023), the SSM parameter was updated from kernel 6.1 to kernel 6.18. As new kernel versions get released (expected annually), the parameter will continue to update to the latest kernel version after a validation period.

This post explains the default kernel behavior, what it means for your workloads, and how to manage the transition.

What’s changing?

Amazon Linux ships multiple kernel versions and has tracked a default kernel for each OS version. For example,
the AL2023 parameter:

ssm:/aws/service/ami-amazon-linux-latest/al2023-ami-{minimal}-kernel-default-{x86_64 arm64}

has remained on kernel 6.1 since launch. Going forward, the kernel-default SSM parameter will update to the latest kernel as new versions are released. Each new kernel will go through a 3- to 6-month validation period after GA before we update the default. This window gives you time to test the new kernel before the change. We will announce the kernel-default upgrade date before it takes effect.

SSM Parameter Resolved to (Before) Resolves to (Now)
al2023-ami-{minimal}-kernel-default-{x86_64, arm64} Kernel 6.1 AMI Kernel 6.18 AMI (what’s changed)
al2023-ami-{minimal}-kernel-6.18-{x86_64, arm64} Kernel 6.18 AMI Kernel 6.18 AMI (unchanged)
al2023-ami-{minimal}-kernel-6.1-{x86_64, arm64} Kernel 6.1 AMI Kernel 6.1 AMI (unchanged)

Note: Already-running instances will keep the kernel they booted with and are not affected by this change. Only new instances launched from the kernel-default parameter will boot kernel 6.18. If you already use a version-specific SSM parameter, nothing changes for you.

Why are we making this change?

The Linux kernel is the foundation of workloads you run on Amazon Elastic Compute Cloud (Amazon EC2) and other services. Each new kernel brings meaningful improvements. For example, kernel 6.18 includes the Earliest Eligible Virtual Deadline First (EEVDF) CPU scheduler for fairer CPU time distribution and improved latency in mixed workloads. The kernel also increases Transmission Control Protocol (TCP) receive buffer for better network throughput on high-bandwidth instances.

Previously, customers who wanted to run the latest Amazon Linux kernel had to manually update their SSM parameter references and redeploy each time a new kernel became available. With this change, you can receive these improvements without needing to manually upgrade.

Evaluating the default kernel upgrade

Staying on the default kernel is the recommended approach as it allows your new instances to always run the latest validated kernel with no manual intervention. However, because the default will now advance annually, you should build processes to validate that the new kernel works for your workload before each upgrade takes effect. If your workload has specific requirements that mandate a fixed kernel version, evaluate whether the new default is compatible or revert to a kernel version that suits your use case.

If you haven’t validated kernel 6.18 yet, we recommend launching test instances on kernel 6.18 using the version-specific SSM parameter al2023-ami-{minimal}-kernel-6.18-{x86_64, arm64}. For instructions on referencing SSM parameters in your launch configuration, see the AL2023 User Guide.

Staying on or reverting to a specific kernel version

If you experience issues with the new default, or if your workload requires a specific kernel version for additional validation time or any other reason, revert to the version-specific SSM parameter. Change your references from al2023-ami-{minimal}-kernel-default-x86_64 to al2023-ami-{minimal}-kernel-{kernel_version}-x86_64 (for example, al2023-ami-kernel-6.1-x86_64). This applies anywhere you resolve an AL2023 AMI, including AWS CloudFormation templates, launch templates, Amazon EC2 Auto Scaling groups, CI/CD pipelines, or CLI scripts. For examples, refer to the AL2023 User Guide.

Each of the supported kernels (6.1, 6.12, and 6.18) continue to receive updates as defined in AL2023 kernel lifecycle. When staying on a specific version, we recommend tracking the kernel lifecycle and planning upgrades before the kernel reaches end of support.

Note: For Federal Information Processing Standards (FIPS) workloads, the default kernel may not always be the FIPS-validated kernel. If you require FIPS mode, see AL2023 FIPS FAQ.

Conclusion

In this post, we announced that the Amazon Linux default SSM parameter will now upgrade to the latest kernel as new kernel versions are released. The AL2023 kernel-default parameter was updated from kernel 6.1 to kernel 6.18 on August 17, 2026. We explained how the new cadence works, how already-running instances are unaffected, and how to stay on a specific kernel version if your workload requires it.

To learn more, see the AL2023 Kernel documentation and the AL2023 release notes. For questions or issues, contact AWS Support.

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.

AI-powered clinical trial eligibility and safety using Amazon Bedrock AgentCore

Post Syndicated from Sachin Jain original https://aws.amazon.com/blogs/architecture/ai-agents-for-clinical-trial-screening/

AI agents built on Amazon Bedrock AgentCore let clinical trial teams make fast, accurate enrollment decisions while keeping clinicians in control through human-in-the-loop oversight. According to the Tufts Center for the Study of Drug Development, 80 percent of clinical trials miss their enrollment timelines, and each day of delay costs an estimated $500,000.

Today, eligibility decisions rely on manual chart review across fragmented sources — EHR notes, lab results, imaging reports, and medication histories. Study teams spend hours reconstructing each candidate’s history and mapping it to protocol criteria. As protocols grow more complex, this doesn’t scale: screen failure rates stay high and enrollment targets slip.

We show how to architect a Clinical Trial Eligibility and Safety Agent on AWS that assembles patient evidence, evaluates it against protocol criteria, and presents screening recommendations with citations, while clinicians retain final authority and full audit trails. It combines AWS HealthLake for FHIR-native data access, Amazon Bedrock AgentCore for multi-step reasoning, and Amazon Bedrock AgentCore Evaluations for scoring each decision via LLM-as-a-judge and human-in-the-loop. This post is for solution architects, engineering teams, and technology leaders applying AI to clinical trial operations on AWS.

AI agents for clinical trial screening

AI agents with Human-in-the-Loop (HIL) are well-suited for clinical trial eligibility and safety decisions because they address information fragmentation while preserving human clinical judgment. The core problem isn’t a lack of data, but that eligibility and safety signals are scattered across EHR notes, lab portals, imaging reports, and medication histories, forcing study teams to reconstruct each participant’s clinical picture. A knowledge graph addresses this by storing clinical data as entities and the relationships between them, representing each patient, molecule, endpoint, and market as a node with relationships stored as edges. To answer an eligibility or safety question, the agent traverses these edges, going from a diagnosis to its associated labs or a medication to its known interactions, rather than re-querying and joining disconnected sources each time. This structure supports the agent’s preparatory work:

  • Organizing evidence from fragmented sources into a knowledge graph, linking patients, molecules, endpoints, and markets as interconnected nodes.
  • Mapping patient information against protocol criteria.
  • Surfacing relevant passages with citations for clinician review.
  • Highlighting uncertainties that require human judgment.

Critically, the clinician remains the decision-maker. The agent organizes the supporting information. These systems augment rather than replace clinical reasoning — proposing preliminary assessments, flagging edge cases, providing confidence scores, and learning from feedback.

As protocols grow more complex with precision oncology and biomarker-driven eligibility, agents manage multi-step logic and maintain consistency across sites, while deferring final judgment to clinical staff.

Architecture overview

This proposed architecture illustrates how core AWS services can be combined to create an end-to-end clinical trial screening pipeline. AWS HealthLake serves as the FHIR-native clinical data foundation, ingesting and normalizing patient records from disparate EHR systems, lab portals, and imaging archives into a unified, queryable data store. Amazon Bedrock AgentCore orchestrates the multi-step workflow assembling patient profiles, matching them against trial protocols, detecting safety signals, and generating evidence-backed screening recommendations. An Amazon Bedrock Knowledge Bases stores trial protocols, inclusion/exclusion criteria, and safety guidelines. The entire pipeline feeds into a clinician review dashboard where investigators examine agent reasoning, verify citations, and render final decisions. Actions are captured in an immutable audit trail for regulatory compliance.

Architecture diagram showing the clinical trial screening pipeline with AWS HealthLake, Amazon Bedrock AgentCore, and Amazon CloudWatch

Architecture workflow

The screening pipeline operates in the following steps. Each step maps to a distinct phase of the eligibility and safety assessment, from data ingestion through clinician review and continuous monitoring.

Step 1: Clinical data ingestion

AWS HealthLake ingests patient records from EHR systems, lab portals, imaging reports, and medication histories, then normalizes them into FHIR R4 resources for standardized, queryable access.

Step 2: Agent orchestration

Amazon Bedrock AgentCore orchestrates three specialized agents, each scoped to a distinct phase of the screening pipeline. They operate within the Amazon Bedrock AgentCore Runtime, which connects to tools through MCP Gateway, maintains session memory so agents reference earlier findings without re-querying, and enforces identity-based access control for least-privilege data access. A built-in code interpreter handles dynamic calculations such as eGFR or BMI derivation.

Pre-screening agent: The first gate. It resolves three threshold questions: Is the patient’s informed consent valid and current? Does their high-level profile (age, diagnosis category, geography) align with basic enrollment parameters? Have they completed any required washout period? Patients who clear all three advance. Those who don’t receive a documented rejection citing the failing criterion.

Detailed screening agent: The core clinical reasoning engine. It walks through all inclusion and exclusion criteria, retrieving the relevant FHIR resources — Observation for labs, Condition for diagnoses, MedicationStatement for medications — and evaluating each against the protocol threshold. It also reviews organ function, adverse drug reactions, and contraindicated conditions, cross-references medications against the investigational product for interactions, and assesses the overall comorbidity profile for risk combinations no single criterion would catch. The output is a structured determination (Eligible, Ineligible, or Requires Review) with a per-criterion evidence matrix, confidence scores, and a reasoning summary citing source records.

Site & enrollment agent: Once a patient clears screening, it handles operational logistics — matching the patient to the most appropriate site by proximity, capabilities, and investigator availability, then confirming open enrollment capacity. If the preferred site is full, it identifies alternatives and flags the study coordinator.

All three agents operate behind Amazon Bedrock Guardrails, which enforce:

  1. PII/PHI filtering to protect patient health information.
  2. Content safety controls to help prevent clinically inappropriate outputs.
  3. Grounding checks to keep responses anchored in retrieved evidence rather than model parametric knowledge.
  4. Denied topic boundaries to keep agents within their screening scope.

Step 3: LLM-as-judge evaluation

Amazon Bedrock AgentCore Evaluations scores every screening decision using a combination of built-in and custom evaluators across three dimensions:

  1. Clinical accuracy: Correctness of the eligibility determination against patient data, faithfulness to source evidence (not hallucinated justifications), logical coherence across reasoning steps, and context relevance confirming the right protocol and patient records were retrieved.
  2. Operational effectiveness: Response completeness and clarity for coordinators reviewing dozens of patients daily, appropriate use of FHIR queries and knowledge base tools, and end-to-end goal success (did the agent complete the full screening workflow?).
  3. Safety compliance: Custom evaluators verify that safety-critical criteria (lab thresholds, restricted medications, contraindicated conditions) were never skipped, that uncertainties are explicitly acknowledged rather than resolved with false confidence, and that all safety flags route to the appropriate review tier.

Decisions that pass evaluation with high confidence proceed to the clinician dashboard. The system flags those that fall below quality thresholds and routes them to human review with the specific evaluation concern highlighted.

Step 4: Human-in-the-loop review and enrollment

Flagged cases and agent recommendations flow into a tiered clinical review structure:

  1. PI review queue: Principal Investigators review flagged decisions from the LLM Judge, examining the agent’s reasoning chain, verifying citations against source records, and rendering a final determination.
  2. Study coordinator dashboard: Coordinators manage trial logistics, scheduling, and the day-to-day enrollment pipeline, using the agent’s structured outputs to accelerate their workflow.
  3. Patient communication: Outreach and consent updates are coordinated through the dashboard, keeping patients informed of their screening status.
  4. Escalation to medical director: Complex or high-risk cases that exceed the PI’s comfort level are escalated to the Medical Director for final adjudication.

Clinicians retain complete override capability at every stage. When a clinician overrides an agent recommendation, approving a patient the agent flagged or rejecting one it cleared, the system captures the corrected decision and the clinician’s reasoning. These corrections expand the ground truth dataset used by Amazon Bedrock AgentCore Evaluations and surface patterns that inform prompt and retrieval tuning, creating a continuous learning loop where human judgment directly improves agent performance over time.

Step 5: Observability and continuous monitoring

Amazon CloudWatch provides end-to-end observability across all agents, surfacing agent traces (step-by-step execution logs), latency metrics, error rates (failed tool calls, guardrail blocks), judge scores (pass/flag rates per agent), HITL metrics (override rates, review latency), and alarm-based escalation when safety thresholds are breached.

Although the current implementation focuses on screening and enrollment, the same agent orchestration framework, evaluation pipeline, and compliance infrastructure support future post-enrollment monitoring agents such as adverse event detection from lab results and clinical notes, protocol deviation tracking, retention risk prediction, and re-screening triggers when clinical changes affect ongoing eligibility. Each inherits the existing scoring, logging, and auditability without requiring a separate governance framework.

Evaluating agent performance in clinical trial screening with human oversight

The screening pipeline’s credibility rests on two layers: an automated evaluation layer that scores every decision, and a human-in-the-loop (HITL) layer that gives clinicians final authority. LLM-as-Judge (Step 3) decides which cases clinicians see and how they’re prioritized. The HITL workflow (Step 4) decides how clinicians act. Together they form a continuous loop where human judgment both safeguards and improves agent performance. Using Amazon Bedrock AgentCore Evaluations, you build a framework spanning three dimensions: clinical accuracy, operational effectiveness, and safety compliance with built-in and custom evaluators that run continuously.

Clinical accuracy and reasoning

Built-in evaluators check whether the agent gets the determination right and whether its reasoning holds up: Correctness (accurate against the patient’s labs, diagnoses, and medications), Faithfulness (reasoning stays grounded in patient data and protocol, not plausible-sounding invention), Coherence (no logical contradictions across steps), Context relevance (the right protocol and records were retrieved), and Goal success rate (the full workflow ran end to end). Custom LLM-as-Judge evaluators add clinical specifics: Eligibility accuracy (each inclusion/exclusion criterion evaluated correctly) and Criteria coverage (no criteria skipped, especially safety-critical lab thresholds and restricted medications).

Operational effectiveness

Accuracy alone is insufficient, output must fit workflows where coordinators review dozens of patients daily. Helpfulness, conciseness, and relevance confirm a clear, scannable, on-topic determination. Instruction following verifies the expected structured format (patient summary, criteria checklist, determination, justification, safety flags, next steps). Tool selection and parameter accuracy check the agent invoked the right tools with correct inputs.

Safety and responsible behavior

Safety carries the strictest thresholds. Harmfulness detection flags clinically dangerous content; Stereotyping detection makes sure decisions aren’t influenced by demographics beyond protocol requirements. Both trigger immediate review. Custom evaluators target the highest-risk failures: Safety flag detection confirms every significant concern surfaced (contraindicated medications, out-of-range labs, disqualifying conditions, drug interactions), with a single miss treated as critical; Uncertainty acknowledgment makes sure the agent recommends human review on missing or ambiguous data rather than making an overconfident call.

The human-in-the-loop safeguard

When a wrong eligibility call can affect patient safety, human judgment is the final safeguard. A score below threshold routes the case to the HITL workflow.

The three agents together produce an eligibility determination with a confidence score. At trial onset, the clinician sets a confidence threshold. Cases below it or flagged by evaluation reach the clinician dashboard with the specific concern highlighted. Clinicians review the full reasoning and approve, reject, or request more information from the same interface. Their corrections are stored alongside machine-approved records, feeding back into future determinations and continuously improving accuracy.

Review and approval workflow

Review is tiered by complexity: automated pre-screening filters clearly ineligible candidates. Low-complexity cases get expedited review, medium-complexity follow standard protocols, and high-complexity edge cases escalate to senior clinicians. Cases unreviewed beyond set timeframes escalate automatically. Final enrollment decisions, low-confidence cases, experimental therapies, and complex histories require human approval. Routine high-confidence checks proceed automatically.

Audit trails

The system generates immutable audit records in Amazon DynamoDB for every decision, capturing clinician ID, timestamp, patient and trial IDs, outcomes, AI recommendations, and complete workflow execution history. These records are designed to support FDA 21 CFR Part 11 requirements for electronic records and signatures, providing documentation for regulatory inspections and quality assurance. Readers should consult their compliance team and conduct their own assessment. See the AWS compliance resources for further guidance.

Security and compliance

Clinical trial data is among the most sensitive in healthcare. HIPAA, FDA 21 CFR Part 11, GxP, and GDPR require strict controls over how patient data is stored, accessed, and processed, and AI agents reasoning over that data introduce new security considerations. This solution protects data at every layer while maintaining the audit trails and privacy standards regulators require.

AWS HealthLake is HIPAA-eligible with encryption at rest and in transit, access controls, and SMART on FHIR authorization. Amazon Bedrock is HIPAA-eligible, SOC 2 attested, ISO and CSA STAR Level 2 certified, and never shares customer data with model providers. AWS PrivateLink keeps traffic off the public internet.

Amazon Bedrock AgentCore enforces agent boundaries at runtime through declarative authorization policies — readable, deterministic rules, outside application code, defining what the agent can access, invoke, and retrieve. AgentCore runs within your Amazon Virtual Private Cloud (Amazon VPC) for network isolation, and AWS CloudTrail records API calls for an immutable audit trail that can support FDA compliance requirements.

Amazon Bedrock AgentCore Evaluations scores each decision using built-in and custom evaluators with an LLM-as-a-Judge approach. Continuous sampling detects drift, and Amazon CloudWatch alerts teams when quality drops below thresholds — ongoing evidence the agent performs within validated parameters, supporting GxP with minimal manual testing.

Conclusion

In this post, we showed how combining the FHIR-native data foundation of AWS HealthLake
with the multi-step reasoning capabilities of Amazon Bedrock AgentCore turns manual,
fragmented clinical trial screening into an AI-assisted workflow that reduces patient matching
time from days to minutes. Clinical trial enrollment remains one of drug development’s most
resource-intensive bottlenecks, and delayed starts carry heavy financial consequences from lost
patent-protected sell time and operational burn. Clinicians receive organized evidence,
transparent reasoning, and actionable recommendations while retaining full decision authority
and audit traceability.

The impact extends beyond speed: more consistent criteria interpretation across sites, earlier
detection of safety contraindications, and lower screen failure rates. As oncology trial eligibility
criteria grow in complexity — with fewer than 5% of cancer patients enrolling under strict
requirements — this human-in-the-loop approach offers a scalable, compliance-aligned path to
faster, higher-quality recruitment.

Call to action

Ready to accelerate your clinical trial operations? Take the next step:

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.

Querying raw log data using SQL and PPL with the optimized engine in Amazon OpenSearch Service

Post Syndicated from Kaushik Krishnan original https://aws.amazon.com/blogs/big-data/querying-raw-log-data-using-sql-and-ppl-with-the-optimized-engine-in-amazon-opensearch-service/

In this post, you learn how to run fast analytical queries directly against raw log and trace data in Amazon OpenSearch Service using PPL and SQL.

Amazon OpenSearch Service is a fully managed service that helps you deploy, scale, and operate OpenSearch, the open source suite for search, analytics, and observability in the AWS Cloud. OpenSearch Service powers search and real-time analytics workloads, from lexical and hybrid search to log analytics and observability. This post focuses on log analytics, and on a practical question: how much analytical work can you do directly against raw log and trace data, without moving it or reshaping it first?

The new optimized engine in OpenSearch Service answers that question: you can point Piped Processing Language (PPL) and Structured Query Language (SQL) queries at raw log and trace data. The engine returns aggregations, filters, and scans over billions of events on the data exactly as you ingested it. In this post, you follow a single incident investigation, one query at a time. You see how the engine answers each new question, from multi-dimensional breakdowns and latency distributions to error rates and fleet sizing. No precomputed structure sits behind the results.

How the optimized engine queries raw data

The optimized engine stores data in the columnar Apache Parquet format and runs queries through Apache DataFusion, a vectorized execution engine, with Apache Calcite planning each query. Because the engine stores data in columns, an analytical query reads only the columns it touches and processes their values in batches, instead of reading each matching document in full. Alongside the columnar format, the engine also keeps an inverted index on the same data, so the query planner routes each operation to the path that serves it best: the columnar engine for aggregations and analytical scans, and the inverted index for selective search and filtering.

You ingest your logs and traces through the same Bulk API and clients you use today, and you write PPL or SQL against them as they land.

An investigation, one query at a time

The following walkthrough traces a common observability use case, root-cause analysis during a live incident, from the perspective of a site reliability engineer (SRE). The engineer notices elevated latency and a handful of error alerts, with nothing that points to a clear cause. No existing dashboard covers this particular shape of problem, so the engineer opens Amazon OpenSearch Service and starts asking questions of the raw trace data, letting each answer decide the next one. PPL suits this work well. Each command transforms the data and passes it to the next, so the engineer reads a query left to right the same way they think through the investigation.

The walkthrough uses generated OpenTelemetry (OTEL) data from a synthetic load generator, at billion-document scale. The focus is the query capability, that is, what the engineer can express and retrieve directly from raw spans, rather than the specific values in each result.

Step 1: Assess the scope

The first question in any investigation is how widespread the signal is. The engineer breaks errors down across service, HTTP method, and cloud Region in a single pass over roughly 1.1 billion spans.

source=otel-traces
| where @timestamp >= timestamp("2026-05-15 00:00:00") and @timestamp < timestamp("2026-05-18 00:00:00")
| eval e = if(status_code = 2, 1, 0)
| stats sum(e) as errors, avg(durationInNanos) as avg_ns, count() as total_count
  by serviceName, http_method, cloud_region
| sort - errors
| head 8

In plain terms, this query answers the engineer’s first question: where are the failures happening? It counts the error spans and breaks them down by service, HTTP method, and AWS Region in a single pass. Rather than guessing which service to open first, the engineer gets a ranked list of the hardest-hit combinations to investigate.

errors total_count avg_ns serviceName http_method cloud_region
730 112,436 41,246,806 export-service GET us-west-2
722 111,215 41,000,227 catalog-service PUT eu-central-1
704 112,051 41,295,539 image-service PATCH us-west-2
612 93,214 41,451,145 healthcheck-service PUT us-east-1
609 94,314 41,418,897 auth-service POST us-east-1
609 94,414 41,447,444 email-service PATCH ap-northeast-1
593 89,726 41,047,114 payment-service PUT eu-central-1
581 89,854 41,195,643 file-service PUT ap-northeast-1

The errors spread across services, methods, and Regions, which points to a systemic pattern rather than a single misbehaving service.

Step 2: Check whether one host concentrates the failures

The spread could still reflect one saturated node or a fleet-wide condition. To tell the two apart, the engineer groups failures by exception type, service, and host across the entire index, with no time filter to narrow the scan.

source=otel-traces
| where isnotnull(exception_type)
| stats count() as total_count by exception_type, serviceName, host_name
| sort - total_count
| head 8
total_count exception_type serviceName host_name
6 DeadlockDetectedException notification-service ip-10-0-16-34
6 IllegalStateException api-gateway ip-10-0-180-234
6 FileNotFoundException cart-service ip-10-0-90-162
6 ConnectionRefusedException feature-flag-service ip-10-0-8-123
5 TimeoutException auth-service ip-10-0-97-78
5 ConcurrentModificationException order-service ip-10-0-165-15
5 TimeoutException coupon-service ip-10-0-158-25

In this sample the counts are low and every row lands on a different host, so no single node stands out. This points to a fleet-wide pattern rather than one bad machine. On production data the same query makes the distinction directly: a code-level bug shows up across many hosts, whereas a single failing node concentrates its errors on one host_name.

Step 3: Quantify the latency distribution per service

Next, the engineer pulls a latency profile for each service. This includes count, average, minimum, and maximum duration, to see how each one behaves and how wide the spread runs.

source=otel-traces
| where @timestamp >= timestamp("2026-05-15 00:00:00") and @timestamp < timestamp("2026-05-18 00:00:00")
| stats count() as total_count, avg(durationInNanos) as avg_ns, min(durationInNanos) as min_ns, max(durationInNanos) as max_ns
  by serviceName
| sort - total_count
| head 8
serviceName total_count avg (ns) min (ns) max (ns)
event-bus 11,087,263 41,249,552 26,113 9,304,132,159
scheduler-service 9,175,964 41,251,927 21,919 13,432,040,933
cdn-service 9,173,572 41,225,385 23,468 13,768,293,306
ml-inference 9,036,753 41,289,101 40,410 14,625,084,517
compliance-service 8,274,635 41,294,694 41,915 7,462,983,016
metrics-collector 7,804,234 41,334,728 16,535 23,228,217,669
notification-service 7,688,714 41,204,635 51,562 8,695,311,374
image-service 7,674,406 41,248,069 47,473 15,350,500,299

This gives the engineer a latency fingerprint for each service: the averages sit near 41 milliseconds. But the multi-second maxima reveal a long tail consistent with requests queuing behind a slow dependency.

Step 4: Measure the error rate per service

To track a service-level objective, the engineer computes the error rate (errors against total requests) per service. The query uses an inline conditional, followed by a grouped sum and count, and a final division to produce the error rate.

source=otel-traces
| eval is_err = if(status_code = 2, 1, 0)
| stats sum(is_err) as errors, count() as total_count by serviceName
| eval error_pct = round(100.0 * errors / total_count, 2)
| sort - error_pct
| head 8
errors total_count error_pct serviceName
699,358 22,415,308 3.12 payment-service
647,811 26,880,140 2.41 checkout-service
562,811 30,096,860 1.87 auth-service
316,192 24,510,990 1.29 cart-service
288,314 30,671,704 0.94 order-service
202,612 28,140,552 0.72 search-service
186,012 33,820,415 0.55 catalog-service
134,722 35,453,247 0.38 image-service

The engineer defines the error-rate metric in the query itself, and the engine computes it across the full index. The busiest paths, payment and checkout, run near 3 percent, whereas some services stay below 1 percent.

Step 5: Size the fleet footprint with SQL

Finally, the engineer sizes how much of the fleet each service spans, a capacity and impact question, and switches from PPL to SQL to express it.

SELECT serviceName,
       COUNT(*) AS total_count,
       COUNT(DISTINCT host_name) AS hosts
FROM otel-traces
GROUP BY serviceName
ORDER BY total_count DESC
LIMIT 8
serviceName total_count hosts
ml-inference 35,481,688 2,535
image-service 35,453,247 2,491
email-service 35,443,569 2,517
shipping-service 30,700,372 2,438
translation-service 30,490,570 2,502
auth-service 30,096,860 2,466
chat-service 25,564,111 2,449
recommendation-service 25,366,844 2,483

The query runs a COUNT(DISTINCT) over a high-cardinality field at billion-row scale, and switching languages mid-investigation costs the engineer nothing more than writing SQL instead of PPL. The host counts cluster in the approximately 2,400–2,540 range, so each service runs across a broad slice of the fleet. That confirms the earlier finding: the errors reflect a fleet-wide pattern, not a single node.

The engineer asked five questions and ran five queries, and each answer shaped the next. The optimized engine served every query directly from raw trace data, across both PPL and SQL, without a rollup table or precomputed summary behind any result.

Run these queries where you already work

You don’t need a separate tool to run the queries in this walkthrough.

Figure 1: Investigation queries and results grid in Query Workbench

Query Workbench in OpenSearch Dashboards UI gives you a dedicated editor for PPL and SQL. You write a query, run it, and read the results in a grid, using the same queries shown throughout this post. When you want to move from a written query to interactive exploration, Discover runs the same PPL and SQL against your indexes. In Discover, you can filter, expand fields, and drill into individual documents without leaving the page. The same query language works in both places, so you can start an investigation in Discover and carry it into Query Workbench, or the reverse, without rewriting anything.

Figure 2: PPL query and field list in Discover

Keep all your data and query it as it is

Querying raw data directly only helps if you can afford to keep the raw data. The optimized engine compresses observability data up to 70 percent more efficiently than the default General Purpose engine. That compression turns “keep everything and query it directly” into a practical default. You retain full-fidelity data for the questions you cannot predict in advance. You also pay less to store it than you would to store the raw JSON.

Get started

To try the optimized engine, create an Amazon OpenSearch Service domain running OpenSearch 3.5 or later. Then select the Observability use case during setup, which provisions the domain with the optimized engine.

To learn more about configuring and using the optimized engine, see Optimized for Log Analytics in the Amazon OpenSearch Service documentation. For an overview of the service, visit Amazon OpenSearch Service Log Analytics.

For more information, see the blog post Run log analytics for a fraction of the cost with the new engine for Amazon OpenSearch Service.

Give it a try and send feedback to AWS re:Post for Amazon OpenSearch Service or through your usual AWS Support contacts.


About the authors

Kaushik Krishnan

Kaushik is a Technical Account Manager at Amazon Web Services with a focus on Amazon OpenSearch Service. He is based in the Washington, D.C. area and specializes in troubleshooting critical operational and performance issues as well as conducting architectural reviews of OpenSearch clusters for customers. Outside of work, he enjoys playing soccer and is an avid traveler.

Luis Tiani

Luis is a Sr Solutions Architect at AWS. He specializes in data and analytics topics, with extensive focus on Amazon OpenSearch Service for search, log analytics, and vector environments. Tiani has helped numerous customers across financial services, DNB, SMB, and enterprise segments in their OpenSearch adoption journey, reviewing use cases and providing architecture design and cluster sizing guidance.

Jagadish Kumar

Jagadish is a Senior Solutions Architect at Amazon Web Services, focused on OpenSearch and analytics workloads.

Send rich RCS messages with AWS End User Messaging RCS

Post Syndicated from Rommel Sunga original https://aws.amazon.com/blogs/messaging-and-targeting/send-rich-rcs-messages-with-aws-end-user-messaging-rcs/

When a customer asks where their order is, a plain text reply answers the question. But a rich RCS message with a product photo, a tappable confirmation button, and a calendar chip helps the customer act on it. Rich Communication Services (RCS) messages deliver branded, interactive content, including images, rich cards, carousels, and suggestion chips, to the messaging app already built into the customer’s phone. Unlike Short Message Service (SMS), RCS messages come from a verified sender with your brand name and logo, deliver over a data connection, and support read receipts and structured replies. AWS End User Messaging RCS provides the SendRcsMessage API, a managed way to send RCS messages through a single integration point instead of separate integrations for each carrier.

This post is for developers and solutions architects who want to add RCS messaging to their customer engagement workflows on AWS. It shows how to send every RCS content type (text, files, rich cards, carousels, and suggestions). It also shows how to control delivery with message expiration and SMS fallback, using Python and the AWS End User Messaging RCS API.

The post focuses on the SendRcsMessage API, which is specific to RCS and is the only one of the two that supports rich cards, carousels, and suggestions. The SMS API’s SendTextMessage can also deliver over RCS when you pass an RCS agent as the origination identity, but it is limited to plain text. Every example that follows uses SendRcsMessage.

Prerequisites

Before you run the examples in this post, you need the following:

  • An AWS account with access to AWS End User Messaging.
  • An AWS RCS agent in the Active state. To send to your customers, the agent needs an approved country launch registration for each destination country. To run the examples before launch approval, use an agent with a testing registration and send to a registered test device: an Android phone with RCS enabled, or an iPhone on iOS 18 or later, with a status of VERIFIED.
  • AWS SDK for Python (Boto3) 1.43.37 or later, which includes SendRcsMessage support. Run pip install --upgrade boto3 to get the latest version.
  • Optionally, the AWS Command Line Interface (AWS CLI) version 2.35.12 or later.
  • For the suggestions example, an Amazon Simple Notification Service (Amazon SNS) topic configured for two-way messaging on your RCS agent, so you can receive suggestion tap events.

If you’re new to RCS on AWS, see Getting started with RCS on AWS End User Messaging SMS to create your agent. You pay standard RCS rates for RCS messages, including messages sent to test devices.

IAM permissions

The AWS Identity and Access Management (IAM) principal that runs the examples needs permissions for the following actions:

  • sms-voice:SendRcsMessage, to send RCS message types.
  • sms-voice:SendTextMessage, to send the plain text comparison example and any SMS fallback messages.
  • sms-voice:DescribeRcsAgents, to check that your agent is Active.
  • sms-voice:DescribeVerifiedDestinationNumbers, to confirm a registered test device is VERIFIED, if you send to one.

If you use the SMS fallback example, you also need a phone number or sender ID in your account that can send SMS to the destination country. RCS and SMS are separate origination identities: the RCS agent sends the RCS message, and the fallback needs its own SMS-capable identity.

If you send media from Amazon Simple Storage Service (Amazon S3), the bucket needs a resource policy granting the sms-voice.amazonaws.com service principal s3:GetObject, shown in the “File messages” section. If you use server-side encryption with AWS Key Management Service (AWS KMS) keys for your bucket, your KMS key policy must also grant the service access. For two-way messaging, your SNS topic needs a resource policy allowing the service to publish to it. For details, see Two-way messaging in the AWS End User Messaging SMS User Guide.

Configuration

Create a config.json file in your project directory to store the RCS agent Amazon Resource Name (ARN) that sends the messages and the recipient phone number in E.164 format:

{
  "rcsAgentArn": "arn:aws:sms-voice:us-east-1:111122223333:rcs-agent/rcs-a1b2c3d4",
  "destinationPhoneNumber": "+12065550100"
}

OriginationIdentity accepts the RCS agent ID (RcsAgentId) or ARN (RcsAgentArn), and also a pool ID or pool ARN. The examples use the agent ARN because it stays unambiguous when an account has more than one agent, but the shorter agent ID works the same way.

The config.json file is for local testing only. In production, don’t hardcode phone numbers and identifiers. Use AWS Secrets Manager, AWS Systems Manager Parameter Store, or environment variables instead.

Each example in this post builds a message_content dictionary and sends it with the following code:

import json
import boto3
import os

config_path = os.path.join(os.path.dirname(__file__), 'config.json')
with open(config_path, 'r') as f:
    config = json.load(f)

client = boto3.client('pinpoint-sms-voice-v2')

message_content = { ... }

response = client.send_rcs_message(
    DestinationPhoneNumber=config['destinationPhoneNumber'],
    OriginationIdentity=config['rcsAgentArn'],
    RcsMessageContent=message_content
)
print(f"Message sent. ID: {response['MessageId']}")

For production use, wrap the send call with error handling to manage throttling and validation failures:

try:
    response = client.send_rcs_message(
        DestinationPhoneNumber=config['destinationPhoneNumber'],
        OriginationIdentity=config['rcsAgentArn'],
        RcsMessageContent=message_content
    )
    print(f"Message sent. ID: {response['MessageId']}")
except client.exceptions.ThrottlingException as e:
    print(f"Rate limited. Retry after backoff: {e}")
except client.exceptions.ValidationException as e:
    print(f"Invalid request or media: {e}")
except Exception as e:
    print(f"Failed to send message: {e}")

The following sections show only the message_content for each message type. To send any of these messages, use the shared sending code from this section. The examples follow one scenario: AnyCompany, a fictitious retailer, messaging a customer about an order.

Text messages

Text messages are the most basic RCS content type. You can send plain text two ways. The SendTextMessage API, the same API used for SMS, delivers over RCS when you pass your RCS agent ARN as the origination identity:

response = client.send_text_message(
    DestinationPhoneNumber=config['destinationPhoneNumber'],
    OriginationIdentity=config['rcsAgentArn'],
    MessageBody="Hello from AnyCompany over RCS."
)

The SendRcsMessage API sends the same text as a TextMessage content type, and additionally supports suggestion chips, message expiration, and per-message fallback. An RCS text also arrives as a single message regardless of length, while carriers split SMS over 160 characters into segments that can arrive out of order.

Specifications and requirements

  • Message body: 1–3,072 UTF-8 characters, required.
  • Up to 11 suggestions per message (covered in the “Suggestions” section)
  • Destination phone number must be in E.164 format.
  • Without a FallbackConfiguration, recipients who can’t receive RCS get nothing.

Text message example

RCS text message from the AnyCompany agent confirming that order ORD-2026-001 has shipped

Figure 1: RCS text message confirming that order ORD-2026-001 has shipped

Code example

The following is the message_content for the preceding message:

message_content = {
    "Content": {
        "TextMessage": {
            "Body": "Thanks for reaching out to AnyCompany! Your order ORD-2026-001 has shipped and arrives on Friday, August 7. Reply to this message if you have any questions."
        }
    }
}

File messages

With file messages, you send a single image, video, audio file, or PDF that renders as inline media in the recipient’s messaging app. FileUrl accepts two URL forms, and they fail in different places.

With an S3 URL (s3://amzn-s3-demo-bucket/object-key), the API checks at request time that the object exists, is within the size limit, and is readable with the permissions you granted the service. If any of those checks fail, the call returns a ValidationException describing the problem, so you find out at send time. The service then retrieves the object, rehosts it, and generates a time-limited presigned URL for delivery to the device.

With an HTTPS URL, the URL is passed through to the carrier and isn’t checked the same way at request time. The API accepts the request. Problems such as an unreachable host, a URL that requires authentication, or an unsupported media type surface at delivery instead of in the API response. The URL must be publicly accessible with no authentication. The API doesn’t support plain http:// URLs.

Use S3 URLs when you want bad media to fail loudly at send time. Use HTTPS URLs for media already published on a public CDN, and monitor delivery events for failures.

Specifications and requirements

  • FileUrl: required, S3 or HTTPS URL, up to 2,000 characters.
  • ThumbnailUrl: optional, JPEG or PNG, recommended for video and PDF.
  • Maximum file size: 100 MB at the API layer. Carriers can enforce lower limits (keep video under 5 MB)
  • Supported formats include JPEG, PNG, and GIF images, MP4 and WebM video, MP3 and AAC audio, and PDF documents. Support varies by carrier and device.

To deliver from Amazon S3, add the following bucket policy so the service can read your objects:

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Principal": {
        "Service": "sms-voice.amazonaws.com"
      },
      "Action": "s3:GetObject",
      "Resource": "arn:aws:s3:::amzn-s3-demo-bucket/*"
    }
  ]
}

Replace amzn-s3-demo-bucket with your bucket name. To restrict access to a prefix, replace /* in the Resource ARN with a path such as arn:aws:s3:::YOUR-BUCKET/rcs-media/*.

File message example

RCS file message showing an inline PDF document attachment

Figure 2: RCS file message rendering an inline PDF attachment

Code example

The following is the message_content for the preceding message:

message_content = {
    "Content": {
        "FileMessage": {
            "FileUrl": "https://docs.aws.amazon.com/pdfs/social-messaging/latest/userguide/social-ug.pdf"
        }
    }
}

Rich cards

A rich card combines media, a title, a description, and suggested actions into a single structured message. Rich cards work well for product highlights, booking confirmations, appointment details, and promotional offers.

Specifications and requirements

  • Title: up to 200 characters. Description: up to 2,000 characters.
  • CardContent requires at least one of Media, Title, or Description
  • CardOrientation is required: VERTICAL or HORIZONTAL. Use VERTICAL because horizontal orientation truncates images on iOS.
  • Media Height: SHORT (112 density-independent pixels), MEDIUM (168), or TALL (264). IOS ignores this value.
  • Card-level suggestions: up to 4 per card.
  • URLs in description text are not tappable. Use OpenUrl suggestions for links.

Rich card message example

Vertical rich card with a product image, title, description, and Add to cart and View details buttons

Figure 3: Vertical rich card with a product image, title, description, and action buttons

Code example

The following is the message_content for the preceding message:

message_content = {
    "Content": {
        "RichCard": {
            "CardOrientation": "VERTICAL",
            "CardContent": {
                "Title": "AnyCompany Wireless Headphones",
                "Description": "Noise-cancelling, 30-hour battery life, available in black or silver. Your loyalty discount brings the price to $179.17.",
                "Media": {
                    "FileUrl": "https://example.com/images/headphones.png",
                    "Height": "MEDIUM"
                },
                "Suggestions": [
                    {
                        "Reply": {
                            "Text": "Add to cart",
                            "PostbackData": "cart_add_headphones"
                        }
                    },
                    {
                        "OpenUrl": {
                            "Text": "View details",
                            "PostbackData": "view_headphones",
                            "Url": "https://example.com/products/headphones",
                            "Application": "BROWSER"
                        }
                    }
                ]
            }
        }
    }
}

Carousels

A carousel displays 2–10 rich cards in a horizontally scrollable strip. Carousels fit browse-and-compare experiences such as product catalogs, service menus, plan comparisons, and location listings. Carousel cards use the same content model as standalone rich cards, with two differences: cards always render in a vertical layout, and the TALL media height is not supported.

Specifications and requirements

  • Cards per carousel: minimum 2, maximum 10.
  • CardWidth: SMALL (180 density-independent pixels) or MEDIUM (296). All cards share the same width.
  • Card title: up to 200 characters. Description: up to 2,000 characters.
  • Media Height: SHORT or MEDIUM only.
  • Suggestions: up to four per card, plus message-level chips below the whole carousel.
  • All cards scale to the height of the tallest card.
Carousel showing the first two product cards, Wireless Headphones and Smart Watch, each with a Select button

Figure 4: Carousel showing the Wireless Headphones and Smart Watch cards, each with a Select button

Scrolling right reveals the remaining cards:

Carousel scrolled to show the Portable Speaker card with a Select button

Figure 5: Carousel scrolled to the Portable Speaker card

Code example

The following is the message_content for the preceding message:

message_content = {
    "Content": {
        "Carousel": {
            "CardWidth": "MEDIUM",
            "CardContents": [
                {
                    "Title": "Wireless Headphones",
                    "Description": "Noise-cancelling, 30-hour battery. $179.17 with your discount.",
                    "Media": {
                        "FileUrl": "https://example.com/images/headphones.png",
                        "Height": "SHORT"
                    },
                    "Suggestions": [
                        {
                            "Reply": {
                                "Text": "Select",
                                "PostbackData": "select_headphones"
                            }
                        }
                    ]
                },
                {
                    "Title": "Smart Watch",
                    "Description": "Fitness tracking, 7-day battery, water resistant. $249.00.",
                    "Media": {
                        "FileUrl": "https://example.com/images/watch.png",
                        "Height": "SHORT"
                    },
                    "Suggestions": [
                        {
                            "Reply": {
                                "Text": "Select",
                                "PostbackData": "select_watch"
                            }
                        }
                    ]
                },
                {
                    "Title": "Portable Speaker",
                    "Description": "360-degree sound, 12-hour battery. $89.99.",
                    "Media": {
                        "FileUrl": "https://example.com/images/speaker.png",
                        "Height": "SHORT"
                    },
                    "Suggestions": [
                        {
                            "Reply": {
                                "Text": "Select",
                                "PostbackData": "select_speaker"
                            }
                        }
                    ]
                }
            ]
        }
    }
}

Suggestions

Suggestions are the interactive chips you saw in the earlier examples. They guide recipients through a conversation with predefined replies and actions, without typing. RCS supports six suggestion types: Reply, OpenUrl, DialPhone, ShowLocation, RequestLocation, and CreateCalendarEvent, and you can mix them in one message on any content type. Message-level suggestions live in a Suggestions array that is a sibling of Content, not nested inside it. Card-level suggestions live inside each card’s CardContent.

Every suggestion requires a Text label and PostbackData. The postback data is invisible to the recipient and comes back to your application when the chip is tapped. Encode routing information there (for example, appt_confirm_12345), and route logic on postback data rather than display text.

Specifications and requirements

  • Text label: up to 25 characters; PostbackData: up to 2,048 characters, both required on every suggestion.
  • Message-level suggestions: up to 11. Card-level suggestions: up to four per card.
  • OpenUrl Url must begin with https://. Set Application to WEBVIEW with a WebviewViewMode of FULL, HALF, or TALL to keep the recipient inside the messaging app.
  • DialPhone PhoneNumber must be in E.164 format.
  • CreateCalendarEvent requires Title, StartTime, and EndTime
  • Two-way messaging with an Amazon SNS topic must be configured to receive suggestion taps. Handle the case where a recipient types free text instead of tapping.

Suggestions message example

RCS text message confirming a fitting appointment at AnyCompany Anytown

Figure 6: RCS message confirming a fitting appointment at AnyCompany Anytown

The first suggestion chips shown below the appointment message: Confirm, Reschedule, Manage booking, and Call the store

Figure 7: Suggestion chips below the appointment message: Confirm, Reschedule, Manage booking, and Call the store

Scrolling the chip row reveals the remaining suggestions:

The remaining suggestion chips: View store map, Share my location, and Add to calendar

Figure 8: Remaining suggestion chips: View store map, Share my location, and Add to calendar

Code example

The following message_content combines all six suggestion types on one text message:

message_content = {
    "Content": {
        "TextMessage": {
            "Body": "Your fitting appointment at AnyCompany Anytown is confirmed for Friday, August 7 at 2:00 PM. How would you like to manage your visit?"
        }
    },
    "Suggestions": [
        {
            "Reply": {
                "Text": "Confirm",
                "PostbackData": "appt_confirm_12345"
            }
        },
        {
            "Reply": {
                "Text": "Reschedule",
                "PostbackData": "appt_reschedule_12345"
            }
        },
        {
            "OpenUrl": {
                "Text": "Manage booking",
                "PostbackData": "appt_manage_12345",
                "Url": "https://example.com/bookings/12345",
                "Application": "BROWSER"
            }
        },
        {
            "DialPhone": {
                "Text": "Call the store",
                "PostbackData": "appt_call_12345",
                "PhoneNumber": "+12065550142"
            }
        },
        {
            "ShowLocation": {
                "Text": "View store map",
                "PostbackData": "appt_map_12345",
                "Latitude": 47.6062,
                "Longitude": -122.3321,
                "Label": "AnyCompany Anytown"
            }
        },
        {
            "RequestLocation": {
                "Text": "Share my location",
                "PostbackData": "appt_share_loc_12345"
            }
        },
        {
            "CreateCalendarEvent": {
                "Text": "Add to calendar",
                "PostbackData": "appt_cal_12345",
                "Title": "AnyCompany fitting appointment",
                "StartTime": "2026-08-07T06:00:00Z",
                "EndTime": "2026-08-07T06:30:00Z",
                "Description": "Fitting appointment at AnyCompany Anytown"
            }
        }
    ]
}

When the recipient taps a chip, the messaging app sends the chip text back into the conversation as a reply:

Tapping Confirm sends the chip text back as a reply from the recipient, shown with a read receipt

Figure 9: Tapping Confirm sends the chip text back as a reply, shown with a read receipt

The tap arrives as an inbound event on your two-way SNS topic. The messageBody field contains a JSON string with a type of SUGGESTION, the display text, and the postback data:

{
  "originationNumber": "+12065550101",
  "destinationNumber": "rcs-a1b2c3d4",
  "messageBody": "{\"type\":\"SUGGESTION\",\"text\":\"Confirm\",\"postbackData\":\"appt_confirm_12345\"}",
  "inboundMessageId": "msg-abc123def456"
}

Note the casing difference: request fields use PascalCase (PostbackData), while inbound events use camelCase (postbackData). A RequestLocation tap delivers the recipient’s coordinates in a separate inbound location event.

Message expiration

The TimeToLive parameter sets an expiration window in seconds on a SendRcsMessage request. If the message is not delivered within that window, the service removes it and the recipient never sees it. This matters for time-sensitive content such as one-time passwords (OTPs): a verification code that arrives after the code has expired only confuses the customer.

Specifications and requirements

  • TimeToLive: integer seconds, 1–172,800 (48 hours). Use at least 10 seconds so the carrier can attempt delivery.
  • The countdown starts when the service accepts the request. Omitting TimeToLive means no expiration window.
  • On expiry you receive a TTL_EXPIRATION_REVOKED event (message removed, safe to send a fallback) or TTL_EXPIRATION_REVOKE_FAILED (revoke failed, the message might still deliver, so weigh the duplicate risk)

Message expiration example

RCS verification code message delivered within its five-minute expiration window

Figure 10: RCS verification code delivered within its five-minute expiration window

Code example

The following example sends an OTP that expires after five minutes. TimeToLive is a request parameter, a sibling of RcsMessageContent:

message_content = {
    "Content": {
        "TextMessage": {
            "Body": "Your AnyCompany verification code is 482913. This code expires in 5 minutes."
        }
    }
}

response = client.send_rcs_message(
    DestinationPhoneNumber=config['destinationPhoneNumber'],
    OriginationIdentity=config['rcsAgentArn'],
    RcsMessageContent=message_content,
    TimeToLive=300
)

Per-message fallback

Fallback is optional, and without it a recipient who can’t receive RCS gets nothing. The FallbackConfiguration request parameter routes the message to SMS or Multimedia Messaging Service (MMS). Fallback applies when the device or carrier doesn’t support RCS, when the channel rejects the message, or when the TimeToLive window expires first.

Specifications and requirements

  • Channel: required, SMS or MMS.
  • MessageBody: required for SMS fallback, up to 1,600 characters (compared with 3,072 for the RCS text body); MMS fallback requires at least one of MessageBody or MediaUrls
  • OriginationIdentity for the fallback: a phone number or sender ID registered in your account that can send SMS or MMS to the destination country. Pools and RCS agents are not accepted here.
  • Write the fallback content separately, because suggestion chips and rich cards don’t translate to SMS. Put URLs as plain text in SMS fallback, or use MMS fallback to preserve visual content.

Per-message fallback example

AnyCompany delivery notification delivered over RCS

Figure 11: AnyCompany delivery notification delivered over RCS

On a device without RCS, the SMS fallback version arrives instead from the fallback phone number.

Code example

The following example sends a delivery notification with an SMS fallback from a dedicated phone number:

message_content = {
    "Content": {
        "TextMessage": {
            "Body": "AnyCompany: your delivery arrives today between 2:00 PM and 4:00 PM. Track it at https://example.com/track/1234"
        }
    }
}

response = client.send_rcs_message(
    DestinationPhoneNumber=config['destinationPhoneNumber'],
    OriginationIdentity=config['rcsAgentArn'],
    RcsMessageContent=message_content,
    FallbackConfiguration={
        "Channel": "SMS",
        "MessageBody": "AnyCompany: your delivery arrives today between 2:00 PM and 4:00 PM. Track it at https://example.com/track/1234",
        "OriginationIdentity": "+12065550188"
    }
)

To track outcomes, pass ConfigurationSetName on the send call so delivery, read, expiration, and fallback events route to your configuration set’s event destinations. Set up event destinations before you send, because they don’t retroactively capture events.

Cleaning up

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

  1. Delete the RCS agent if you no longer need it. If you enabled deletion protection when creating it, disable that first. If you registered test devices, remove their verified destination numbers first.
  2. Delete the Amazon SNS topics and configuration set event destinations you created for two-way messaging and status events.
  3. Delete any media objects you uploaded for the examples and the bucket policy from your S3 bucket.
  4. Review Amazon CloudWatch Logs for log groups created by event destinations and delete them if no longer needed.

Conclusion

In this post, you learned how to send every RCS content type with AWS End User Messaging RCS, including text messages, file messages, rich cards, carousels, and suggestions. You also learned how to control delivery with message expiration and per-message SMS fallback. You sent each type from a short Python script, with one shared sending pattern across all content types.

The SendRcsMessage API keeps one pattern across all content types: a Content object for the message body and a sibling Suggestions array for interactivity. Moving from a plain text notification to a full product carousel is a change to one dictionary.

Next steps:

  • Build event-driven replies by subscribing an AWS Lambda function to your two-way Amazon SNS topic and routing on postback data.
  • Design a fallback strategy that pairs TimeToLive values with per-message SMS or MMS fallback for each use case.
  • If you started with a testing registration, submit a country launch registration when you’re ready to send to your customers.

Create your first RCS agent in the AWS End User Messaging SMS & RCS console and send a test message today. Tell us about your experience: share your use cases and questions in the comments.

Additional resources


About the authors

Consistency is the new latency: AI at the data layer

Post Syndicated from Suman Chatterjee original https://aws.amazon.com/blogs/architecture/consistency-is-the-new-latency-ai-at-the-data-layer/

As AI applications scale from reactive bots to autonomous agents, their reliability is bound to the speed and accuracy of the data layer beneath them.

The integrity crisis nobody is talking about

There’s a quiet assumption baked into most AI architectures today regarding data layer consistency, and it’s costing companies more than they realize. The assumption is that the data your AI agent reads is the current state of reality.

In a world of distributed systems, cross-region replication, and autonomous agents making millisecond decisions, this assumption breaks down.

I’ve spent extensive time working with enterprise teams building agentic AI, and a recurring failure pattern emerges.

The breakdown isn’t in the model or the prompts. It’s in how we manage replication consistency when an agent performs the reading.

The context window is the new database row

In a modern agentic Retrieval-Augmented Generation (RAG) architecture, the database is the active memory of your AI. When an agent performs a task, it retrieves data to build its context window, forming the foundation of the large language model’s (LLM) reasoning.

If that data is even slightly out of date, the agent’s entire reasoning chain is invalidated. We must shift from simply managing data availability to strictly verifying contextual integrity.

The silent poison of asynchronous lag

In traditional web applications, asynchronous replication scales global reads with minimal write impact. If a user sees a post 500ms late, nobody notices.

For an autonomous AI agent, a 500ms delay is silent poison. If an agent writes a decision to a primary node and immediately reads from a lagging replica, it treats stale data as ground truth. It then executes a logically coherent, multi-step plan based on factually incorrect inputs.

In the age of AI, a fast answer that is wrong is more expensive than a slightly slower answer that is right.

The anatomy of a stale-read failure: When memory betrays logic

Consider an autonomous Inventory Reconciliation Agent managing a flash sale:

  1. The write: The agent updates available_stock to 500 units on the primary database in us-east-1.
  2. The lag: Network congestion causes a 2-second replication lag to the ap-south-1 (Mumbai) replica.
  3. The read: A secondary agent instance in Mumbai queries the replica and retrieves the old value: 0 units.
  4. The failure: The agent triggers a “Sold Out” notification and halts the sale, despite having 500 units in the warehouse.

The agent didn’t make a reasoning error. It performed logical operations on poisoned context.

Diagram of the stale read cascade, showing how replication lag feeds outdated data into an AI agent’s context

Figure 1: The stale read cascade, showing how replication lag poisons an AI agent’s context

The hallucination debt problem

When an agent writes an incorrect conclusion back to the database, that error becomes long-term memory. Future retrievals pull this poisoned history, creating a self-reinforcing cycle of “Hallucination Debt.”

LLMs amplify this because they lack a temporal compass. They cooperatively treat retrieved database results as current facts without hesitation. The burden of verifying contextual integrity falls entirely on the architecture.

The replication trinity: Choosing your truth

Not all AI tasks have the same consistency requirements. You must match your replication model to the specific “truth requirement” of the task.

Here are three architectural patterns I’ve found most effective.

Pattern A: Precision through global consistency

When an agent manages high-stakes data (user permissions, security policies, financial records, core system instructions), the cost of a stale read is unacceptable. You need strong consistency.

For many workloads, Amazon Aurora Global Database provides the necessary foundation. While its cross-region storage replication is asynchronous by default, you can close the consistency gap by turning on Global Write Forwarding with a GLOBAL consistency level.

To verify Read-Your-Own-Writes integrity, you configure the SESSION consistency level, which makes an agent wait for its own forwarded writes to replicate back before reading.

For the strongest consistency, the GLOBAL level makes a read query wait for replication to catch up to the exact point in time when the read started.

For the next generation of globally distributed AI, Amazon Aurora DSQL addresses this need. Aurora DSQL offers native synchronous strong consistency across multiple regions, so multi-agent systems can scale globally without compromising accuracy.

Every agent, regardless of location, operates on the exact same ground truth.

Best for: Identity metadata, financial ledgers, immutable system prompts.

Why it matters: Eliminates “mid-thought” state changes that cause contradictory behavior between agent instances.

Pattern B: Global availability at scale

For global AI agents that need ultra-low latency at massive scale, Amazon DynamoDB Global Tables offer a multi-leader architecture where data replicates across regions. For replication details, refer to the DynamoDB documentation.

The key technique here is Conditional Writes. By using a ConditionExpression that checks a version timestamp or whether an attribute exists, an agent updates a record only if the data hasn’t changed since it was last retrieved.

If the condition fails, DynamoDB returns a ConditionalCheckFailedException. This is a critical signal: it tells the agent to re-read the current state and reconsider its decision, rather than blindly overwriting another agent’s work.

This pattern prevents the “Lost Update” anomaly (where two agents running in parallel overwrite each other’s reasoning) without requiring synchronous global coordination.

Best for: Conversational history, user session state, personalized agent memory.

Why it matters: Handles concurrent updates from distributed agents while maintaining a shared memory that’s resilient to race conditions.

Pattern C: High-velocity intake

Some AI agents perform real-time anomaly detection or trend analysis on massive streams of telemetry data. In these cases, you need unthrottled ingestion above all else.

A leaderless architecture like Amazon Keyspaces (for Apache Cassandra) is designed for this workload.

Keyspaces provides highly available, predictable performance by automatically replicating data across three Availability Zones.

Every write is durably committed using LOCAL_QUORUM.

To make sure your AI agent doesn’t miss a critical spike in telemetry, you enforce strong consistency by setting its read operations to LOCAL_QUORUM rather than the eventually consistent LOCAL_ONE.

This quorum overlap means the agent retrieves the latest data without slowing down the high-speed ingestion pipeline.

It transforms a noisy, high-frequency data stream into a reliable foundation for real-time AI decision-making.

Best for: Internet of Things (IoT) telemetry, real-time log analysis, high-frequency sensor data.

Why it matters: Throughput is the priority, but you still need a safety valve to confirm the agent doesn’t miss critical spike data.

Conclusion: Becoming a context architect

Our role as architects has evolved.

We can no longer treat database replication as a background infrastructure concern, something to configure once and forget. In the era of autonomous agents, the stability of the data layer is the direct prerequisite for the trustworthiness of the AI. The two are inseparable.

By matching your replication model to your agent’s reasoning requirements, you move beyond simply managing data. You become a Context Architect, someone who works to confirm that every decision your AI makes is grounded in a synchronized version of the truth.

Because in the end, an AI is only as good as the context it operates in. And context is only as good as the data it’s built on.

Get the database layer right, and everything else follows.

References:


About the author

How Zepto powers sub-second search using OpenSearch Service OR2 instances

Post Syndicated from Kayalvizhi Kandasamy original https://aws.amazon.com/blogs/big-data/how-zepto-powers-sub-second-search-using-opensearch-service-or2-instances/

Sub-second search is the starting point of every order on Zepto, a fast-growing quick-commerce platform in India, founded in 2021 with endeavor to provide delivery in minutes. Powering the search experience is Amazon OpenSearch Service, a managed retrieval engine built on OpenSearch for agentic AI, search, and analytics.

Zepto operates hundreds of delivery hubs (dark stores) across Indian cities where it provides logistics services to sellers operating on Zepto Platform. Each hub maintains its own inventory levels, pricing, and assortment spanning thousands of Stock Keeping Units (SKUs). As the company scaled to hundreds of hubs, driving linear increases in indexing volume and maintaining sub-second product search latency while controlling costs became increasingly challenging.

To address this, Zepto migrated its OpenSearch Service data nodes from memory-optimized instances to OpenSearch Optimized instances. This instance family is purpose-built for high indexing throughput and cost efficiency. It uses local storage as the primary data tier, with Apache Lucene segments copied synchronously to Amazon Simple Storage Service (Amazon S3) for durability. With this migration, Zepto now serves the same workload with two-thirds of their previous data node count, achieving over 100% higher indexing throughput and 30% cost savings.

In this post, we explore the architecture decisions along with the load testing outcomes that led Zepto to select OpenSearch Optimized instances for latency-sensitive product search. We also discuss the key lessons learned during the production migration.

Zepto’s search platform

Zepto’s search platform is built around a localized delivery hub model. Each hub maintains its own inventory, capacity, and fulfillment priority. When a customer searches for a product, the query is not resolved against a global catalog. Instead, it is resolved in the context of the specific delivery hub or hubs serving that customer’s delivery address. This distinction is critical: Every customer journey on Zepto’s application begins with product discovery through search, browse, and promotional surfaces. All these must reflect hub-specific availability in real time to fulfill orders in minutes.

An event-driven architecture powers this experience, keeping results fresh as products, prices, offers, and inventory change across hundreds of delivery hubs. The following architecture diagram illustrates Zepto’s end-to-end indexing and search pipeline, from event production through stream processing to the search indices on OpenSearch Service.

Figure 1: Zepto’s end-to-end indexing and search pipeline architecture

Event producers and consumers: Zepto’s application microservices are deployed on Amazon Elastic Kubernetes Service (Amazon EKS), a fully managed service for running Kubernetes workloads on AWS. These microservices serve as both event producers and consumers. Sellers and Zepto Admin users interact with the Zepto Partner and Admin application.

Key microservices: The Catalog Management Service emits events when product metadata changes like new product additions, attribute updates, and category reclassifications. The Inventory Management Service publishes stock-level changes across delivery hubs in real time as warehouse teams pick, pack, and replenish inventory. The Pricing Management Service generates events whenever sellers update pricing. The Offers Management Service broadcasts events when promotional offers are created, activated, modified, or expired. Together, these microservices capture every relevant update for downstream indexing, producing events into the streaming layer whenever business state changes.

Search events streaming: All domain events flow through Amazon Managed Streaming for Apache Kafka (Amazon MSK), a managed streaming data service that manages Apache Kafka infrastructure and operations.

The system organizes events into dedicated Kafka topics by business domain. These include Catalog for product metadata changes, Inventory for hub-level stock updates, Pricing for price changes across stores, and Offers for promotional offer lifecycle events and more. This topic-based partitioning provides independent scaling per domain, ensures ordered delivery within each topic and consumer isolation, so that a surge in inventory events does not disrupt catalog indexing.

Stream processing and routing: Events from MSK topics are consumed and routed into two priority-based indexing pipelines through dedicated Apache Flink OpenSearch Connector jobs deployed on an Amazon EKS cluster:

  • Job #1: P0 indexing events (Pipeline #1): Processes high-priority events requiring near real-time index freshness, such as inventory changes, catalog enrichment, and pricing updates.
  • Job #2: P1 indexing events (Pipeline #2): Handles lower-priority but higher-volume events, such as tag updates, semantic embedding generation, offer activations, and nightly revenue per impression (RPI) score recomputation. These updates improve search quality but can tolerate slightly higher latency.

With this dual-job approach, Zepto maintains sub-second freshness for critical signals like stock availability and current pricing. Compute-heavy enrichment updates are processed separately without creating backpressure on real-time updates.

Search indices: Zepto hosts the search index on OpenSearch Service, structured at the city-product level. Delivery hub-specific metadata, such as stock status and hub-level demand signals, is stored as nested documents within each record. The following example depicts a typical document in the search index.

{
    "city": "mumbai",
    "product_id": "SKU-29401",
    "product_name": "Amul Butter 500g",
    "category": "Dairy",
    "offer_ids": ["OFFER-201", "OFFER-305"],
    "tags: [weekend, liquidation],
    "hubs": [
        {
            "hub_id": "MUM-HUB-01",
            "rpi_score": 0.0142,
            "stock_status": "in_stock",
            "hub_signals": {"demand": "peak"},
            "active": "true"
        },
        {
            "hub_id": "MUM-HUB-02",
            "rpi_score": 0.0147,
            "stock_status": "low_stock",
            "hub_signals": {"demand": "low"},
            "active": "false"
        },
        ..
    ]
}

The document structure supports store-level personalization while organizing the index by city-product pairs.

Search pipeline: Zepto’s search platform decouples the search request flow from the indexing pipeline at the application layer. When a customer initiates a search, the request passes through the Zepto application to the Search Service and Orchestration layer, which queries the OpenSearch index and assembles the response.

The Search Service and Orchestration layer handles the complete query lifecycle. This includes query understanding, candidate retrieval, machine learning (ML) ranking, ad slotting, and response assembly. For a detailed overview of Zepto’s full search architecture, refer to Building Search for a 10-Minute World on the Zepto engineering blog.

Scaling challenge

Zepto’s search platform started with a single use case, basic product search. As the business expanded, the platform introduced increasingly sophisticated experiences and each new experience added indexing signals to the pipeline like offer events, liquidation tags, pricing changes, ranking scores and more. All needed to be ingested and reflected in the index. Simultaneously, growing user traffic and the expansion of browse surfaces increased read throughput demands on the cluster.

The challenge was most acute during festive events like Diwali and New Year, when traffic surges required scaling to 1.4× the data node count. Although the cluster handled the node additions, the team needed to monitor shard relocation progress and validate that search latencies remained within service level agreements (SLAs) at each step. This operational overhead grew with each scaling event.

Adding more nodes to the cluster would address the immediate throughput constraints, but at the cost of proportionally higher infrastructure spend. To find a solution, Zepto set a clear goal: “Improve throughput without increasing the data node count.”

Solution overview

With the goal of keeping the node count intact, Zepto experimented with multiple configurations. One approach was resharding, adjusting the number of primary shards to better distribute the workload across existing nodes. However, load testing under production-representative traffic revealed that each resharding configuration degraded search latencies. The resharding operations themselves were also operationally expensive, requiring full index recreation, data migration, and extended validation windows.

The team needed a fundamentally different approach. The approach needed to improve throughput without adding nodes or resharding the index.

Evaluating OpenSearch Optimized instances

OpenSearch Optimized instances are an instance family purpose-built for workloads that require high indexing throughput with cost efficiency. They are commonly used for log analytics and time series use cases. These instances store data on Amazon Elastic Block Store (Amazon EBS) volumes for fast local access. Apache Lucene segments are synchronously replicated to Amazon S3, providing 11 nines of data durability.

Despite the common use case association with log analytics, we recommended evaluating OpenSearch Optimized instances type OR2 for Zepto’s product search workload. The team assessed two key criteria to determine viability:

Criterion 1: Does segment replication address the throughput bottleneck?

With document replication (the default on memory-optimized instances), every write is indexed on the primary shard and then re-indexed independently on each replica. This duplicates CPU work across the cluster. With segment replication on OpenSearch Optimized instances, segments are built once on the primary shard. They are then copied as complete files to replicas. This eliminates the duplicate indexing pipeline on replicas and frees their compute for serving search queries. Zepto’s workload involved continuous indexing from multiple pipelines that competed with search traffic. This separation was the key architectural advantage.

Criterion 2: Can the search platform tolerate the 10-second refresh interval?

OpenSearch Optimized instances use a 10-second segment replication refresh interval that is longer than the default one-second refresh on memory-optimized instances. This means newly indexed documents become searchable with up to 10 seconds of additional delay. The team evaluated whether this trade-off was acceptable for their search use cases.

Rahul Pradeep, Senior Architect at Zepto, explains:

“Out-of-stock or in-stock is not a primary parameter for retrieval. It is more like a tiebreaker. Relevance is our primary parameter. We retrieve hundreds of products in one query and then do a last-minute validation against our real-time inventory service. That is why we may not need one-second refresh.”

Zepto’s existing architecture where the Product Enrichment Service validates inventory after retrieval indicated that the 10-second refresh interval would not impact customer experience; see how Zepto built Product Enrichment at scale for further details. The migration was viable without any application-level changes.

Based on this evaluation, the solution involved migrating from memory-optimized Graviton-based data nodes to OpenSearch Optimized instances. This shift changed how indexing work is distributed across the cluster. Instead of a model where every node duplicates the full indexing pipeline, only the primary shard performs indexing, and replicas receive pre-built segments.

Load testing

To validate the hypothesis before committing to a migration, we designed a proof of concept, a load testing setup in their lower environment that mirrored production characteristics:

  • Baseline cluster with r7g.12xlarge instances and a parallel testing cluster with or2.12xlarge instances, having four nodes per cluster.
  • Identical shard configuration (X primary shards, Y replica, Z shard copies per node).
  • Simultaneous indexing and read load simulation.

Document structure improvements

In addition to validating the infrastructure change, the team identified an opportunity to optimize the document structure itself to further improve search latency. They added an active_hubs attribute to the base document, a flat array listing only the hubs where the product is currently stocked and active as shown in the following updated document structure.

{
    <City and product metadata>,
    "active_hubs: [MUM-HUB-01, MUM-HUB-03],
    "hubs": [
        {
            "hub_id": "MUM-HUB-01",
            "rpi_score": 0.0142,
            "stock_status": "in_stock",
            "hub_signals": {"demand": "peak"},
            "active": "true"
        },
        ..
    ]
}

The following table summarizes the key metrics from the load test comparing the r7g.12xlarge baseline cluster against the or2.12xlarge test cluster under identical conditions.

Metric r7g.12xlarge or2.12xlarge Change
Peak indexing lag ~12M docs ~6M docs 2X Faster
Indexing throughput Baseline 2× higher 100% Improvement
Search latency (p90) 187 ms 89.1 ms 52% Improvement
Search latency (p99) 244 ms 175 ms 28% Improvement

The following graph depicts the P90 search latency comparison between the two clusters.

Line graph comparing P90 search latency for the r7g.12xlarge and or2.12xlarge clusters over time

Figure 2: P90 search latency comparison between the r7g and OR2 clusters

The following graph depicts the P99 search latency comparison between the two clusters.

Line graph comparing P99 search latency for the r7g.12xlarge and or2.12xlarge clusters over time

Figure 3: P99 search latency comparison between the r7g and OR2 clusters

The following graphs depict the indexing latency comparison between the two clusters.

Figure 4: Indexing latency comparison between the r7g (left) and OR2 (right) clusters

Key insights

  • Improvement in indexing throughput: The higher indexing throughput of OR2 is most visible during nightly batch operations when events from RPI score recomputation, tag updates, and catalog enrichment flood the indexing pipeline simultaneously. On the r7g cluster, the P1 indexing lag peaked at over 12M docs. On OR2, with approximately 2× the indexing throughput, the same event volume produced a peak lag of only 6M docs. Higher throughput translates directly to lower lag and fresher search results. It is attributed to the segment replication approach of OR2 that eliminates redundant indexing work on replicas. Each document is indexed once on the primary shard rather than being replayed on each replica.
  • Reduction in search latency: P90 search latency dropped from 187 ms to 89.1 ms (52% improvement) and P99 from 244 ms to 175 ms (28% improvement). These gains are primarily attributable to the active_hubs document structure change rather than the instance type migration alone. By pre-computing a flat list of active hubs at index time, the query no longer needs to traverse nested hub documents to determine availability. This creates a lightweight pre-filter that eliminates unnecessary computation at search time.

Production planning and rollout

The load test results gave Zepto the confidence to make a key architectural decision: reduce the overall data node count. Higher per-node indexing throughput meant the same workload could be served with fewer nodes with OR2, and the cost savings compounded. Each eliminated node removed compute, storage, and operational overhead from the cluster. Zepto carried this forward into production, provisioning the OR2 cluster at two-thirds of the original node count. The following table summarizes the before-and-after comparison.

Metric r7g.12xlarge or2.12xlarge Change
Data nodes required 3X Nodes 2X Nodes -33.3%
Cost savings Baseline 2/3 of Baseline +30%

Rather than a complete cutover, Zepto adopted a phased rollout strategy using bucket-based traffic routing, completing the migration over approximately two months with zero downtime:

  1. Provisioned a new OpenSearch Service domain on OR2 instances with segment replication turned on.
  2. Executed parallel indexing pipelines to populate the OR2 cluster while the existing r7g cluster continued serving production traffic.
  3. Routed internal users to the OR2 cluster first to validate search quality, relevance, and latency characteristics under real query patterns.
  4. Gradually increased external user traffic in buckets, monitoring comparison dashboards at each increment for latency regressions or relevance drift.
  5. Maintained parallel dashboards throughout the migration to compare the OR2 and r7g clusters in real time. Key metrics monitored included p50 and p99 search latency, indexing throughput, replica lag, Java Virtual Machine (JVM) heap utilization, circuit breaker events, I/O operations per second (IOPS) utilization, and disk throughput.

Challenges and lessons learned

During the migration, the team encountered one notable challenge: latency spikes during segment merges. This observation offers practical guidance for teams evaluating OpenSearch Optimized instances for search workloads.

Symptom: After shifting significant traffic to OR2, Zepto observed intermittent p99 latency spikes correlating with segment merge operations.

Root cause: Large segment merges consumed significant I/O bandwidth, temporarily impacting concurrent search query performance. The original 256 GB EBS volumes did not provide sufficient IOPS buffer for concurrent merge and search operations.

Resolution: Implemented the following two changes:

  1. Increased EBS volume size to 1 TB per node. For gp3 volumes, baseline IOPS increase with volume size. This provided buffer for concurrent operations.
  2. Tuned the segment merge policy. Reduced max_merged_segment (see OpenSearch: Force Merge API for more details) from 5 GB to 2 GB and segments_per_tier (see OpenSearch: Index Settings for more details) from 10 to 5. This produces smaller, more frequent merges that distribute I/O load more evenly rather than infrequent large merges that spike latency.

After increasing EBS volume size and tuning the segment merge policy, latency spikes decreased. Transient spikes still occurred during merges but settled quickly within acceptable bounds.

Production cutover

Finally, Zepto shifted from partial to 100% traffic over four weeks and decommissioned the previous cluster after confirming stable performance across multiple peak traffic cycles. The following table summarizes the cluster configuration before and after migration.

Parameter Previous Cluster Current Cluster
Instance type r7g.12xlarge or2.12xlarge
Data nodes 3X Nodes 2X Nodes
RAM per node 384 GiB 384 GiB
Replication strategy Document replication Segment replication
Default refresh interval 1 Second 10 Seconds
Durability Cross-Availability Zone replicas S3 synchronous replication

Conclusion

In this post, we described Zepto’s evaluation of OpenSearch Optimized instances for latency-sensitive product search and the results of their production migration. By moving from memory-optimized data nodes to OpenSearch Optimized instances with segment replication enabled, Zepto achieved over 100% higher indexing throughput and 30% cost savings while reducing their cluster to two-thirds of the previous data node count.

Zepto’s migration demonstrates that OpenSearch Optimized instances are a viable choice for latency-sensitive product search and not just log analytics. Workloads where the retrieval layer can tolerate seconds-level staleness because real-time consistency is resolved at a different layer are candidates for adopting OR2 instances. For ecommerce and quick-commerce platforms that separate candidate generation from availability validation, this pattern can deliver significant infrastructure cost reduction.

If your workload has high indexing volume, and can tolerate a 10-second refresh interval, consider evaluating OpenSearch Optimized instances for your cluster. To get started:

  1. Assess your workload fit: review your current indexing throughput, replica count, and refresh interval requirements. Prioritize this approach if your workload has a high write-to-read ratio.
  2. Execute a proof of concept: provision a small OpenSearch Optimized cluster in a lower environment with identical shard configuration and restore a production index snapshot. Execute simultaneous indexing and search load to validate throughput and latency.
  3. Plan a phased rollout: use parallel indexing and bucket-based traffic routing to migrate incrementally with zero downtime, monitoring indexing lag and search latency at each step.

To explore the architecture behind OpenSearch Optimized instances, see Under the Hood: OpenSearch Optimized Instances. For practical configuration guidance, see Improve performance with OpenSearch Optimized Instances. We welcome your questions and feedback in the comments section below.


About the authors

Mayank Agarwal

Mayank Agarwal

Mayank is a Principal Architect at Zepto, focused on ML platforms, search, and supply-chain systems at scale. He writes about technology at agarwalknayam.com.

Kayalvizhi Kandasamy

Kayalvizhi Kandasamy

Kayalvizhi is a Principal Solutions Architect at AWS. She specializes in helping customers turn ideas into production-ready solutions using AI/ML, analytics, serverless, and microservices on AWS. A FIDE-rated chess player, Kayalvizhi enjoys passing on her love for the game by coaching her daughters.

Rahul Pradeep

Rahul Pradeep

Rahul is a Senior Architect at Zepto working across systems in the Shopping Journey of a user. His focus areas are to build high throughput systems in search, discovery and checkout domains that can stand the test of growing scale and changing business needs. He writes about tech at raahulpradeep.com

Bhagwati Malav

Bhagwati Malav

Bhagwati is an Engineering Leader at Zepto, building and scaling Search and Discovery systems. His work focuses on information retrieval, search relevance, semantic search, ranking systems, and distributed systems.

Pawananjani Kumar

Pawananjani Kumar

Pawananjani is a Senior Engineer on Zepto’s Search team. He works on the core mechanics of search retrieval, result ranking, and data indexing.

Rugved Sawarkar

Rugved Sawarkar

Rugved is a Senior Engineer within the Search team at Zepto. His primary responsibilities center on building and refining search retrieval, ranking algorithms, and indexing pipelines.

Harpreet Singh

Harpreet Singh

Harpreet is a Senior Technical Account Manager at AWS based in Bangalore, specializing in cloud optimization, resilience, and GenAI-driven operations. He develops innovative strategies to solve complex customer challenges across high-growth industries. He aims to drive cloud adoption and operational excellence at scale. Outside work, Harpreet enjoys playing badminton and exploring new technology.

Aashi Agarwal

Aashi Agarwal

Aashi is a Solutions Architect at AWS, where she specializes in the analytics domain. She guides customers through the transformative process of migration and modernization. With a blend of visionary architecture and robust security, she crafts resilient systems and seamlessly integrates cutting-edge AI/ML services, including the marvels of generative AI, into their technological tapestry. Outside of work, she loves to explore new things and discovers music.

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.

OAuth 2.0, LDAP, and HTTP auth for Amazon MQ for RabbitMQ

Post Syndicated from Vinodh Kannan Sadayamuthu original https://aws.amazon.com/blogs/big-data/oauth-2-0-ldap-and-http-auth-for-amazon-mq-for-rabbitmq/

This is Part 2 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 AWS Identity and Access Management (IAM) authentication, see Part 3.

When you deploy Amazon MQ for RabbitMQ in an enterprise environment, authentication quickly becomes more complex than a single broker configuration. Your organization might already have an Active Directory managing thousands of users, or a cloud identity provider handling application access, or workloads that require short-lived, token-based credentials. Maintaining a separate set of static RabbitMQ credentials alongside these systems creates operational overhead and introduces security gaps. This is especially true when users change roles, leave the organization, or when credentials need to be rotated across multiple brokers.

Amazon MQ for RabbitMQ supports OAuth 2.0, LDAP, and HTTP-based authentication backends, so you can connect your broker directly to the identity infrastructure you already use. This post explains how each approach works, highlights the key configurations, and helps you decide which one fits your use case.

Overview

This post covers three authentication and authorization integrations for Amazon MQ for RabbitMQ:

  1. OAuth 2.0: Token-based authentication where clients obtain short-lived tokens from an identity provider and present them to the broker as bearer credentials. The broker validates tokens using JSON Web Key Sets (JWKS) and derives permissions from token scopes.
  2. LDAP: Directory-based authentication where the broker delegates credential verification to an LDAP directory such as Active Directory. Users authenticate with their directory credentials, and RabbitMQ permissions map to LDAP group memberships.
  3. HTTP authentication backend: A flexible approach where the broker delegates authentication and authorization decisions to an external HTTP service, so you can implement custom logic or integrate with identity systems that don’t support OAuth 2.0 or LDAP natively.

All three approaches eliminate the need to manage broker-local credentials. They provide centralized user management, fine-grained access control, and audit capabilities through your existing identity infrastructure.

How OAuth 2.0 authentication works

OAuth 2.0 authentication eliminates static broker credentials by using short-lived tokens issued by an external identity provider. Instead of storing usernames and passwords in the broker, clients obtain access tokens and present them as credentials when connecting.

When a client connects to a broker configured with OAuth 2.0 authentication:

  1. The client requests an access token from the OAuth 2.0 identity provider, specifying the required scopes.
  2. The identity provider validates the client credentials and issues a signed JWT (JSON Web Token) containing the granted scopes.
  3. The client connects to the Amazon MQ broker and presents the JWT as the password.
  4. The broker retrieves the identity provider’s public keys through the JWKS endpoint.
  5. The broker validates the token signature, expiration, and audience claim.
  6. The broker extracts RabbitMQ permissions from the token scopes and grants access accordingly.

The following diagram shows the OAuth 2.0 authentication flow.

OAuth 2.0 authentication flow between a client, an identity provider, and the Amazon MQ for RabbitMQ broker

Figure 1: OAuth 2.0 authentication flow for Amazon MQ for RabbitMQ

Scope-to-permission mapping

The broker maps OAuth 2.0 scopes to RabbitMQ permissions using a configurable prefix. For example, with the resource server ID rabbitmq, the following scopes grant specific access:

OAuth 2.0 scope RabbitMQ permission
rabbitmq.read:*/* Read access to all resources in all vhosts
rabbitmq.write:*/* Write access to all resources in all vhosts
rabbitmq.configure:*/* Configure access to all resources in all vhosts
rabbitmq.read:orders/* Read access to all resources in the orders vhost
rabbitmq.tag:management Management UI access
rabbitmq.tag:administrator Administrator access

Key configuration

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

# Enable OAuth 2.0 authentication (with internal fallback for the monitoring user)
auth_backends.1 = oauth2
auth_backends.2 = internal

# OAuth 2.0 resource server configuration
auth_oauth2.resource_server_id = rabbitmq
auth_oauth2.preferred_username_claims.1 = sub

# JWKS endpoint for token validation
auth_oauth2.jwks_uri = https://your-idp.example.com/.well-known/jwks.json

# Additional token validation
auth_oauth2.issuer = https://your-idp.example.com
auth_oauth2.scope_prefix = rabbitmq.

# Skip audience validation for IdPs that do not emit an aud claim matching resource_server_id
auth_oauth2.verify_aud = false

The following table describes each configuration setting.

Setting Purpose
auth_backends.1 = oauth2 Enables the OAuth 2.0 authentication backend (use auth_backends.2 = internal for the monitoring user fallback)
auth_oauth2.resource_server_id Identifies this broker as a resource server. Used as the scope prefix
auth_oauth2.preferred_username_claims.1 JWT claim used to extract the username for display and logging
auth_oauth2.jwks_uri URL of the identity provider’s JWKS endpoint for token signature validation (named jwks_url on RabbitMQ 3.x, jwks_uri on 4.x)
auth_oauth2.issuer Expected token issuer. Tokens from other issuers are rejected
auth_oauth2.verify_aud Whether the broker validates the token’s aud claim against resource_server_id. Set to false for IdPs that do not emit a matching aud
auth_oauth2.scope_prefix Prefix applied to scopes when mapping to RabbitMQ permissions

Important considerations

  1. By default the broker validates the token’s aud (audience) claim against the resource_server_id and rejects tokens without a match. Some identity providers (for example, Amazon Cognito) don’t emit an aud claim matching the resource server. For those, set auth_oauth2.verify_aud = false.
  2. If your identity provider cannot issue scopes in the native RabbitMQ form (for example, it disallows the * wildcard), use auth_oauth2.scope_aliases entries to translate the provider’s scope names to RabbitMQ scopes such as rabbitmq.read:*/*.
  3. Configure short-lived tokens (one hour or less) and implement token refresh logic in your client applications.
  4. The JWKS endpoint must be reachable from the broker’s network. For private identity providers, verify network connectivity and DNS resolution.
  5. On RabbitMQ 3.x the JWKS endpoint setting is auth_oauth2.jwks_url. On RabbitMQ 4.x it is auth_oauth2.jwks_uri. Use the setting name that matches your broker engine version.
  6. Amazon MQ automatically creates a system user named monitoring-AWS-OWNED-DO-NOT-DELETE with monitoring-only permissions. This user uses the internal RabbitMQ authentication system even on OAuth 2.0-enabled brokers.

How LDAP authentication works

LDAP authentication connects your RabbitMQ broker to an existing directory service such as Active Directory. Instead of managing users locally in the broker, the broker delegates authentication to the LDAP server and derives permissions from directory group memberships. This centralizes user management and lets you apply your existing password policies, account lockout rules, and audit trails to broker access.

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

  1. The client connects to the Amazon MQ broker with a username and password.
  2. The broker constructs a Distinguished Name (DN) from the username using the configured user_dn_pattern.
  3. The broker performs an LDAP bind operation against the directory server using the constructed DN and the client’s password.
  4. If the bind succeeds, the broker queries the directory for the user’s group memberships.
  5. The broker maps group memberships to RabbitMQ permissions (vhost access, resource permissions, and management tags).
  6. The client is authenticated and authorized based on the LDAP query results.

The following diagram shows the LDAP authentication flow.

LDAP authentication flow showing the Amazon MQ broker binding to a directory server and mapping group memberships to permissions

Figure 2: LDAP authentication flow for Amazon MQ for RabbitMQ

LDAP directory structure

This implementation uses a group-centric LDAP model where RabbitMQ concepts (vhosts, exchanges, queues, and tags) are represented as sub-OUs under a single groups hierarchy:

OU=rabbitmq
├── OU=users
│   ├── CN=app-orders-producer
│   └── CN=app-orders-consumer
│
└── OU=groups
    ├── OU=vhosts
    │   ├── CN=vhost-orders
    │   └── CN=vhost-payments
    │
    ├── OU=exchanges
    │   ├── CN=orders-publisher
    │   └── CN=payments-publisher
    │
    ├── OU=queues
    │   ├── CN=orders-consumer
    │   └── CN=payments-consumer
    │
    └── OU=tags
        ├── CN=rmq-admin
        └── CN=rmq-monitor

Users are assigned to groups based on their required access. For example, app-orders-producer would be a member of vhost-orders and orders-publisher, granting it access to the orders vhost and write permissions on the orders exchange.

Key configuration

The following rabbitmq.conf snippet shows the essential settings for LDAP authentication:

# Enable LDAP as primary backend with internal as fallback
auth_backends.1 = ldap
auth_backends.2 = internal

# LDAP server connection (LDAPS on port 636)
auth_ldap.servers.1 = your-active-directory-server.example.com
auth_ldap.port = 636
auth_ldap.user_dn_pattern = CN=${username},OU=users,OU=rabbitmq,DC=example,DC=com
auth_ldap.use_ssl = true
auth_ldap.ssl_options.verify = verify_peer
auth_ldap.log = true

# AWS integration: assume an IAM role to retrieve the CA certificate for LDAPS
aws.arns.assume_role_arn = arn:aws:iam::111122223333:role/AmazonMqLdapRole
aws.arns.auth_ldap.ssl_options.cacertfile = arn:aws:s3:::your-ca-cert-bucket/ca-cert.pem

# Management console tags
auth_ldap.queries.tags = '''
[{administrator, {in_group, "CN=rmq-admin,OU=tags,OU=groups,OU=rabbitmq,DC=example,DC=com"}},
{management, {in_group, "CN=rmq-monitor,OU=tags,OU=groups,OU=rabbitmq,DC=example,DC=com"}}]
'''

# Vhost access control
auth_ldap.queries.vhost_access = '''
{in_group, "CN=vhost-${vhost},OU=vhosts,OU=groups,OU=rabbitmq,DC=example,DC=com"}
'''

# Resource access control
auth_ldap.queries.resource_access = '''
{for, [{permission, configure,
{in_group, "CN=rmq-admin,OU=tags,OU=groups,OU=rabbitmq,DC=example,DC=com"}},
{permission, write,
{for, [{resource, exchange,
{in_group, "CN=orders-publisher,OU=exchanges,OU=groups,OU=rabbitmq,DC=example,DC=com"}}]}},
{permission, read,
{for, [{resource, queue,
{in_group, "CN=orders-consumer,OU=queues,OU=groups,OU=rabbitmq,DC=example,DC=com"}}]}}]}
'''

The following table describes each configuration setting.

Setting Purpose
auth_backends.1 = ldap Sets LDAP as the primary authentication backend
auth_backends.2 = internal Falls back to internal authentication if LDAP is unavailable
auth_ldap.servers.1 LDAP server hostname or IP address
auth_ldap.user_dn_pattern Template for constructing the user DN from the provided username
auth_ldap.port LDAP server port; 636 for LDAPS
auth_ldap.use_ssl Enables an encrypted LDAPS connection to the directory server. Amazon MQ requires that you explicitly set either auth_ldap.use_ssl = true or auth_ldap.use_starttls = true. The broker fails configuration validation if neither is set.
auth_ldap.ssl_options.verify Certificate verification mode for the LDAPS connection. Verify_peer validates the server certificate
aws.arns.assume_role_arn ARN of the IAM role the broker assumes to retrieve the CA certificate
aws.arns.auth_ldap.ssl_options.cacertfile ARN of the CA certificate (in S3) used to validate the LDAP server’s TLS certificate
auth_ldap.queries.tags Maps directory group membership to the administrator and management console tags
auth_ldap.queries.vhost_access LDAP query that determines which vhosts a user can access based on group membership
auth_ldap.queries.resource_access LDAP query that determines resource-level permissions (configure, write, read) based on group membership

Important considerations

  1. Amazon MQ requires an encrypted LDAP connection: you must explicitly set either auth_ldap.use_ssl = true (LDAPS on port 636) or auth_ldap.use_starttls = true (StartTLS on port 389). The broker rejects the configuration if neither is set. Unencrypted LDAP transmits credentials in plaintext, so always use one of these options to protect credentials in transit between the broker and your directory server.
  2. The user_dn_pattern must match your directory’s organizational structure exactly. Verify the pattern with an LDAP browser before applying it to the broker.
  3. With Active Directory, user DNs are usually based on the display name rather than the sign-in name, so a fixed user_dn_pattern often will not match. In that case, configure DN lookup (auth_ldap.dn_lookup_bind, auth_ldap.dn_lookup_base, and auth_ldap.dn_lookup_attribute = sAMAccountName) so the broker resolves each username to its full DN before binding.
  4. LDAP configuration changes require a broker reboot to take effect. However, user permission changes in the directory (group membership additions or removals) take effect immediately for new connections.
  5. Configure the internal backend as a fallback to maintain access if the LDAP server becomes temporarily unavailable.

How HTTP authentication works

The HTTP authentication backend delegates all authentication and authorization decisions to an external HTTP service. When a client connects, the broker sends requests over HTTPS to your service, which responds with allow or deny decisions. Amazon MQ requires encrypted connections and rejects any configuration that uses a plain http endpoint. This approach provides maximum flexibility for integrating with identity systems that don’t support OAuth 2.0 or LDAP natively, or when you need custom authentication logic. The HTTP authentication backend is available on Amazon MQ for RabbitMQ version 4 and above.

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

  1. The client connects to the Amazon MQ broker with a username and password.
  2. The broker sends an HTTPS POST request to the configured authentication endpoint with the username and password.
  3. The external authentication service validates the credentials against its identity store and responds with allow or deny.
  4. For each authorization check (vhost access, resource permissions, topic permissions), the broker sends additional HTTPS requests to the corresponding endpoints.
  5. The authentication service evaluates the authorization request and responds with allow, deny, or allow with tags.
  6. The client is authenticated and authorized based on the authentication service responses.

The following diagram shows the HTTP authentication flow.

HTTP authentication flow showing the Amazon MQ broker sending credential and authorization checks to an external HTTP service

Figure 3: HTTP authentication flow for Amazon MQ for RabbitMQ

The broker sends HTTPS POST requests to four endpoints. Each endpoint must return a plain-text response:

Endpoint Request parameters Expected response
/auth/user username, password allow [tag1, tag2] or deny
/auth/vhost username, vhost, ip allow or deny
/auth/resource username, vhost, resource, name, permission allow or deny
/auth/topic username, vhost, resource, name, permission, routing_key allow or deny

Key configuration

The following rabbitmq.conf snippet shows the essential settings for HTTP authentication:

# Enable the HTTP backend with caching to reduce load on the auth service
auth_backends.1 = cache
auth_backends.2 = http
auth_cache.cached_backend = http

# HTTP authentication endpoints (HTTPS required)
auth_http.http_method = post
auth_http.user_path = https://your-auth-service.example.com/auth/user
auth_http.vhost_path = https://your-auth-service.example.com/auth/vhost
auth_http.resource_path = https://your-auth-service.example.com/auth/resource
auth_http.topic_path = https://your-auth-service.example.com/auth/topic

# TLS configuration for the HTTPS connection to the auth service
auth_http.ssl_options.verify = verify_peer
auth_http.ssl_options.sni = your-auth-service.example.com

# AWS integration: IAM role and CA certificate for secure credential retrieval
aws.arns.assume_role_arn = <your-assume-role-arn>
aws.arns.auth_http.ssl_options.cacertfile = <your-ca-cert-arn>

The following table describes each configuration setting.

Setting Purpose
auth_backends.1 = cache auth_backends.2 = http Enables the HTTP authentication backend with a cache layer in front, which reduces the number of calls to your authentication service
auth_http.user_path URL the broker calls to authenticate users
auth_http.vhost_path URL the broker calls to check vhost access
auth_http.resource_path URL the broker calls to check resource permissions (queues, exchanges)
auth_http.topic_path URL the broker calls to check topic-level permissions
auth_http.http_method HTTP method the broker uses to call the endpoints. Set to post
auth_http.ssl_options.verify Certificate verification mode for the HTTPS connection to the auth service. Verify_peer validates the server certificate
auth_http.ssl_options.sni Server Name Indication hostname sent during the TLS handshake with the auth service
aws.arns.assume_role_arn ARN of the IAM role the broker assumes to securely retrieve the CA certificate
aws.arns.auth_http.ssl_options.cacertfile ARN of the CA certificate the broker uses to validate the auth service’s TLS certificate

Important considerations

  1. The HTTP authentication service must be highly available. If the service is unreachable, all authentication attempts fail. Consider deploying it behind a load balancer with health checks.
  2. HTTPS is mandatory for all authentication endpoints. The broker rejects any endpoint configured with a plain http URL, ensuring credentials are always protected in transit.
  3. Front the HTTP backend with the cache backend (auth_backends.1 = cache) to reduce the number of calls to your authentication service and improve connection latency. Also keep your service’s response times low to avoid connection timeouts and degraded broker performance.
  4. The authentication service receives plaintext passwords. Make sure the service handles credentials securely and doesn’t log them.
  5. The broker connects to your authentication service over TLS. Configure certificate validation with auth_http.ssl_options.verify = verify_peer, and provide the CA certificate and the IAM role for retrieving it through the aws.arns.auth_http.ssl_options.cacertfile and aws.arns.assume_role_arn settings.

Implementation guides

For step-by-step deployment and validation instructions, see the following resources:

  1. Amazon MQ for RabbitMQ OAuth 2.0 authentication – Configure OAuth 2.0 token-based authentication for Amazon MQ.
  2. Amazon MQ for RabbitMQ LDAP integration – Configure LDAP directory integration for Amazon MQ.
  3. Amazon MQ for RabbitMQ HTTP authentication backend – Configure the HTTP authentication backend for Amazon MQ.
  4. Amazon MQ samples repository – AWS Cloud Development Kit (AWS CDK) stacks and sample code for LDAP and OAuth 2.0 integrations.

Conclusion

This post explained how OAuth 2.0, LDAP, and HTTP authentication backends work for Amazon MQ for RabbitMQ, and when to use each one. OAuth 2.0 provides token-based, passwordless authentication with automatic credential expiration. LDAP connects your broker to existing directory infrastructure for centralized user and group management. The HTTP backend offers maximum flexibility for custom identity integrations. Used individually or in combination, these approaches eliminate broker-local credential management and provide centralized access control through your existing identity infrastructure.

In the next post in this series, we cover IAM authentication and OAuth 2.0 authorization for Amazon MQ for RabbitMQ.

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

  1. Amazon MQ Developer Guide: Security
  2. RabbitMQ OAuth 2.0 plugin documentation
  3. RabbitMQ LDAP plugin documentation
  4. Amazon MQ samples repository

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.


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.

Sarath Kumar Kallayil Sreedharan

Sarath Kumar Kallayil Sreedharan

Sarath Kumar K.S. is a Senior Technical Account Manager/Enterprise Support lead at Amazon Web Services. Sarath works with enterprise customers to help them architect and build highly reliable and cost-effective solutions on AWS. He specializes in serverless, messaging technologies, and AI services, and has a background in application development and architecture. In his spare time, he enjoys reading, traveling, playing cricket, and spending time with his family

Mutual TLS and SSL certificate authentication for Amazon MQ for RabbitMQ

Post Syndicated from Harshith Mithamar original https://aws.amazon.com/blogs/big-data/mutual-tls-and-ssl-certificate-authentication-for-amazon-mq-for-rabbitmq/

This is Part 1 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 OAuth 2.0, LDAP, and HTTP authentication, see Part 2. For IAM authentication, see Part 3.

When you use Amazon MQ for RabbitMQ to handle sensitive data, standard TLS encryption alone might not meet your compliance requirements. Compliance frameworks like SOX, HIPAA, and PCI DSS often require verification of the identity of both parties in a connection. Features like mutual TLS (mTLS) and SSL certificate authentication can help support those requirements by adding certificate-based identity verification to your messaging infrastructure.

Amazon MQ for RabbitMQ version 4 or later supports two certificate-based security features that address these needs: SSL certificate authentication for passwordless certificate-only login, and mTLS for certificate-based peer verification with username and password authentication. This post explains how each approach works, highlights the key configuration options, and helps you decide which one fits your use case.

Overview

This post covers two certificate-based security features for Amazon MQ for RabbitMQ:

  1. SSL certificate authentication: Passwordless authentication where clients authenticate solely using X.509 client certificates through the EXTERNAL SASL mechanism. The broker extracts the username directly from the certificate, eliminating the need for passwords.
  2. Mutual TLS (mTLS): Certificate-based peer verification where both the client and broker prove their identities using certificates, while clients still authenticate with a username and password. This secures AMQP connections and the RabbitMQ management interface.

Both features are available for Amazon MQ for RabbitMQ version 4 and above, and both use AWS ARNs for certificate and credential references, integrating with AWS Certificate Manager (ACM), and AWS Identity and Access Management (IAM).

How SSL certificate authentication works

SSL certificate authentication eliminates the need to transmit credentials during connection. The broker extracts the client’s identity from the certificate, though the corresponding user must exist in RabbitMQ’s internal store for authorization. Instead of using certificates only for transport-layer verification, the broker uses the EXTERNAL SASL mechanism to extract the client’s identity directly from the X.509 certificate.

When a client connects to a broker configured with SSL certificate authentication:

  1. The client initiates a TLS connection and presents its client certificate.
  2. The Amazon MQ broker assumes an IAM role to retrieve the CA certificate from ACM.
  3. The broker validates the client certificate against the configured CA certificate.
  4. The broker extracts the username from the client certificate using the configured field (Common Name, Distinguished Name, or Subject Alternative Name).
  5. The broker authenticates the client using the extracted username. No password required.

The following diagram shows the SSL certificate authentication flow. On the left, the client application holds only an X.509 client certificate with no credentials. In the center, the arrows show the TLS handshake carrying the client certificate to the broker, and the return path confirming authentication with no password needed. On the right, the Amazon MQ for RabbitMQ broker performs certificate validation, assuming an IAM role to retrieve the CA certificate from ACM. It then uses the EXTERNAL SASL mechanism to extract the username from the certificate’s CN, DN, or SAN field and establishes the authenticated session.

Client authenticates to the Amazon MQ for RabbitMQ broker using only an X.509 certificate through the EXTERNAL SASL mechanism, with no password

Figure 1: SSL certificate authentication flow

Username extraction options

The broker can extract the client identity from different fields of the X.509 certificate:

ssl_cert_login_from value Certificate field used Example
common_name Common Name (CN) CN=myapp → username myapp
distinguished_name Full Distinguished Name CN=myapp,O=MyOrg → username CN=myapp,O=MyOrg
subject_alternative_name Subject Alternative Name (SAN) entry SAN dns:myapp.example.com → username myapp.example.com

When you use subject_alternative_name, you also configure ssl_cert_login_san_type (dns, ip, email, uri, or other_name) and ssl_cert_login_san_index to specify which SAN entry to use.

Note: The username extraction options for ssl_cert_login_from apply only to SSL certificate authentication. mTLS doesn’t extract identity from the client certificate.

Key configuration

The following rabbitmq.conf snippet shows the essential settings for SSL certificate authentication:

# Enable certificate-only authentication
auth_mechanisms.1 = EXTERNAL
ssl_cert_login_from = common_name
auth_backends.1 = internal
# Require client certificates
ssl_options.verify = verify_peer
ssl_options.fail_if_no_peer_cert = true
# AWS integration for certificate retrieval
aws.arns.assume_role_arn = ${AmazonMqAssumeRoleArn}
aws.arns.ssl_options.cacertfile = ${CaCertArn}

The following table describes what each setting controls:

Setting Purpose
auth_mechanisms.1 = EXTERNAL Enables the EXTERNAL SASL mechanism, authenticating clients using their X.509 certificate instead of a username and password
ssl_cert_login_from = common_name Tells the broker which certificate field to extract the username from
ssl_options.verify = verify_peer Enables client certificate verification
ssl_options.fail_if_no_peer_cert = true Rejects connections from clients that do not present a certificate
aws.arns.assume_role_arn IAM role ARN the broker assumes to retrieve certificates from ACM
aws.arns.ssl_options.cacertfile ARN of the CA certificate in ACM used to validate client certificates

Note: EXTERNAL and internal serve different purposes. EXTERNAL is the authentication mechanism that verifies client identity using the X.509 certificate. internal is the authorization backend that resolves permissions for the authenticated user from RabbitMQ’s built-in user store.

Important considerations

  1. Client certificates must be signed by a trusted Certificate Authority (CA). The broker validates the certificate chain during authentication.
  2. Amazon MQ enforces the use of AWS ARNs for certificate-related settings. Use aws.arns.ssl_options.cacertfile instead of ssl_options.cacertfile.
  3. 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 SSL certificate-enabled brokers and is restricted to loopback interface access only.
  4. If any setting requires the use of an AWS ARN, you must also provide aws.arns.assume_role_arn.
  5. Amazon MQ doesn’t currently support CRL or OCSP for certificate revocation. To revoke a client certificate that’s no longer trusted, replace the CA certificate on AWS Private Certificate Authority (AWS Private CA), re-issue valid client certificates, and apply a configuration update to the broker.
  6. To rotate certificates, update the CA certificate on AWS Private CA and update the broker configuration. Configuration changes don’t take effect immediately. To apply your changes, wait for the next maintenance window or reboot the broker.

How mutual TLS (mTLS) works

Standard TLS works like visiting a secure website: only the server proves its identity to your browser using a certificate. With mTLS, both your client application and the message broker must prove their identities using certificates. This two-way authentication helps verify that only authorized clients can connect to your broker. Unlike SSL certificate authentication, mTLS still requires a username and password at the application layer.

When your client connects to an Amazon MQ broker with mTLS enabled, the following authentication process occurs:

  1. The client initiates a TLS connection and presents its client certificate.
  2. The Amazon MQ broker assumes an IAM role to retrieve the CA certificate from ACM.
  3. The broker validates the client certificate against the CA certificate.
  4. The client authenticates with a username and password in the application layer.
  5. Authentication succeeds, and the broker establishes a secure, encrypted connection with the client.

Note: Unlike SSL certificate authentication, mTLS doesn’t extract the username from the certificate. The client certificate proves transport-layer trust only. The broker validates it against the CA certificate but does not use any certificate fields for application-level authentication. The username provided at login doesn’t need to match the client certificate’s CN.

The following diagram illustrates this two-layer flow. On the left, the client application holds both a client certificate and a username and password. In the center, the arrows show the TLS handshake carrying the client certificate to the broker, followed by the credentials. On the right, the Amazon MQ for RabbitMQ broker performs certificate validation at the transport layer, assuming an IAM role to retrieve the CA certificate from ACM. It then authenticates the username and password at the application layer before establishing the secure connection to the client.

Mutual TLS flow in which the broker validates the client certificate, then authenticates the username and password at the application layer

Figure 2: Mutual TLS authentication flow

With mTLS, you can secure:

  • Client connections to the AMQP endpoint.
  • The RabbitMQ management interface.
  • Connections to OAuth 2.0 identity providers.
  • HTTPS authentication server connections.
  • Lightweight Directory Access Protocol (LDAP) server communications.

Key configuration

The following rabbitmq.conf snippet shows the essential settings for mTLS:

auth_backends.1 = internal
# Require client certificates for AMQP and management
ssl_options.verify = verify_peer
ssl_options.fail_if_no_peer_cert = true
management.ssl.verify = verify_peer
# AWS integration for certificate retrieval
aws.arns.assume_role_arn = ${AmazonMqAssumeRoleArn}
aws.arns.ssl_options.cacertfile = ${CaCertArn}
aws.arns.management.ssl.cacertfile = ${CaCertArn}

The following table describes the mTLS-specific settings and their purpose:

Setting Purpose
ssl_options.verify = verify_peer Enables client certificate verification for AMQP connections
ssl_options.fail_if_no_peer_cert = true Rejects connections from clients that do not present a certificate
management.ssl.verify = verify_peer Enables client certificate verification for the RabbitMQ management interface
aws.arns.ssl_options.cacertfile ARN of the CA certificate in ACM used to validate client certificates for AMQP
aws.arns.management.ssl.cacertfile ARN of the CA certificate in ACM used to validate client certificates for the management interface

Certificate requirements

Both SSL certificate authentication and mTLS require three types of certificates:

  • Server certificate: Authenticates the broker to clients. Obtain from AWS Private Certificate Authority (AWS Private CA) and reference using an AWS ARN.
  • Client certificates: Authenticate each client application to the broker. Issue from your organization’s CA or AWS Private CA.
  • CA certificate: Validates client certificates on the broker side. Store in ACM and reference in the broker’s SSL configuration.

Comparing SSL certificate authentication and mTLS

Use the following table to decide which method fits your security requirements:

Aspect SSL certificate authentication Mutual TLS (mTLS)
Authentication mechanism EXTERNAL SASL — certificate is the sole credential Transport-layer cert verification + username/password at application layer
Password required No Yes
Username source Extracted from certificate (CN, DN, or SAN) Provided by client at login
SASL mechanism EXTERNAL PLAIN (default)
Management interface cert verification Not included by default Supported through management.ssl.verify
Key config directive auth_mechanisms.1 = EXTERNAL ssl_options.verify = verify_peer
Use case Passwordless environments, PKI-managed identities Adding cert verification to existing credential-based auth
Compliance fit Environments requiring no passwords on the wire Frameworks requiring two-factor (something you have + something you know)

Choose SSL certificate authentication when eliminating passwords entirely from your messaging layer, or when your PKI infrastructure already manages client identities. Choose mTLS when adding transport-layer certificate verification to an existing deployment that relies on username/password authentication, or when compliance frameworks mandate two-factor authentication.

Additional SSL options

Both methods support the following additional configuration options:

Configuration Description
ssl_options.depth Maximum certificate chain depth for verification
ssl_options.hostname_verification Hostname verification mode: wildcard or none
ssl_cert_login_san_type SAN type when using Subject Alternative Name: dns, ip, email, uri, or other_name
ssl_cert_login_san_index Zero-based index of the SAN entry to use

Implementation guides

For step-by-step deployment and validation instructions, see the following resources:

Both tutorials use AWS CDK for infrastructure deployment and include validation scripts to test connectivity.

Conclusion

SSL certificate authentication and mTLS each address different security requirements for Amazon MQ for RabbitMQ. SSL certificate authentication uses the X.509 certificate as the sole credential through the EXTERNAL SASL mechanism, eliminating passwords entirely. mTLS adds transport-layer certificate verification on top of existing username/password authentication, giving you two-factor security. If you are building a regulated environment, SSL certificate authentication removes passwords from the wire entirely, which might help support security requirements in frameworks that address credential management. If you’re incrementally hardening an existing deployment, mTLS lets you add transport-layer verification without changing how clients authenticate. In the next post in this series, we cover OAuth 2.0, LDAP, and HTTP authentication for Amazon MQ for RabbitMQ.

To get started with Amazon MQ for RabbitMQ, see the Amazon MQ service page.

Additional resources

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


About the authors

Harshith Mithamar

Harshith Mithamar

Harshith is a Technical Account Manager at AWS. He works with enterprise customers to help them build secure, scalable messaging solutions on AWS.

Vinodh Kannan Sadayamuthu

Vinodh Kannan Sadayamuthu

Vinodh Kannan 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.

Authentication and authorization options for Amazon MQ for RabbitMQ

Post Syndicated from Vinodh Kannan Sadayamuthu original https://aws.amazon.com/blogs/big-data/authentication-and-authorization-options-for-amazon-mq-for-rabbitmq/

Managing authentication for message brokers at scale is complex: credentials sprawl, audit requirements, and integration with existing identity providers create operational overhead. The default approach of creating RabbitMQ users with static usernames and passwords works for getting started, but it quickly becomes a liability at scale. Credentials must be distributed securely, rotated regularly, and revoked promptly when team members change roles or leave the organization. For regulated industries, auditors want to see that your messaging infrastructure enforces the same identity and access controls as the rest of your environment.

Different organizations have different identity infrastructures. Some manage users through Active Directory. Others have standardized on OAuth 2.0. Platform teams building on AWS want to use AWS Identity and Access Management (IAM) roles and policies they understand. Security-conscious environments might require certificate-based authentication where no passwords are transmitted over the network at all.

Amazon MQ for RabbitMQ supports multiple authentication and authorization methods, so you can connect your broker to the identity infrastructure you already use. This post introduces the available options and helps you choose the right one for your use case.

Authentication methods at a glance

Amazon MQ for RabbitMQ supports the following authentication and authorization methods:

Method Credential type User management Recommended for
Simple credentials Username / password Broker-local Getting started, development environments
OAuth 2.0 Bearer tokens from external identity provider External identity provider Workloads that need short-lived tokens from a third-party identity provider
IAM authentication Short-lived JSON Web Tokens (JWTs) from AWS Security Token Service (AWS STS) IAM AWS-native workloads, multi-tenant isolation, credential-free authentication
LDAP Directory credentials Active Directory or LDAP server Organizations with existing directory services
HTTP-based auth backend Username / password validated by external server External HTTP server Custom auth logic, centralized user management across brokers
SSL certificate authentication Certificate only (passwordless) Broker-local (username extracted from cert) Eliminating passwords entirely with certificate-only identity
Mutual TLS (mTLS) Certificate and username/password Broker-local Adding transport-layer certificate verification to existing credential-based auth

Choosing the right method

The right choice depends on your existing identity infrastructure, security requirements, and operational preferences.

Simple credentials

The default method. You create RabbitMQ users with usernames and passwords directly on the broker. This is a straightforward way to get started, but it requires you to manage credentials manually. Choose this for development, testing, or small-scale deployments where credential management overhead is acceptable.

OAuth 2.0

Clients obtain short-lived tokens from any OAuth 2.0-compatible identity provider and present them to the broker as bearer tokens. Choose this when you have an existing identity provider (other than IAM) that issues tokens for your applications, and you want automatic token expiration without managing broker-local credentials.

IAM authentication

IAM serves as an identity provider. Client applications use their IAM credentials to obtain a short-lived JWT from AWS Security Token Service (AWS STS) and present it as a bearer token. IAM policies control which roles can obtain tokens. RabbitMQ scope aliases on the broker map each role’s Amazon Resource Name (ARN) to specific resource permissions (read, write, configure, and administrator). AWS CloudTrail logs every token issuance for auditing. Choose this when your workloads run on AWS compute services with IAM roles, and you want credential-free, IAM-native authentication with broker-level authorization.

LDAP

Connect your broker to an existing directory service such as Active Directory. Users authenticate with their directory credentials, and RabbitMQ permissions map to LDAP group memberships. Choose this when your organization already manages users and groups through a directory service, and you want to apply existing password policies and group-based access control to broker access.

HTTP-based auth backend

Delegates authentication and authorization decisions to a custom HTTPS server. The broker sends HTTP requests to your server for user validation, virtual host access, resource permissions, and topic permissions. Choose this when you need custom authentication logic, want to centralize user management across multiple brokers, or need to integrate with an identity system that doesn’t support OAuth 2.0 or LDAP natively.

SSL certificate authentication

Removes passwords entirely. The broker uses the EXTERNAL Simple Authentication and Security Layer (SASL) mechanism to extract the client’s identity directly from the X.509 certificate (for example, from the Common Name field) and uses it as the RabbitMQ username. With this method, your application doesn’t transmit credentials over the network. Choose this when your security policy requires passwordless authentication, and you manage client identities through a public key infrastructure (PKI).

Mutual TLS (mTLS)

Adds certificate verification on top of existing username/password authentication. During the TLS handshake, the client validates the broker’s certificate and the broker validates the client’s certificate, then the client provides a username and password at the application layer. This gives you two-factor security: something you have (the certificate) plus something you know (the password). Choose this when compliance frameworks require mutual authentication, but you want to retain your existing username/password authentication flow.

Conclusion

Amazon MQ for RabbitMQ version 4 supports seven authentication and authorization methods. With these methods, you can align your message broker security with your existing identity infrastructure. Your organization might standardize on IAM, manage identities through Active Directory, federate access through a third-party identity providers like Okta or Microsoft Entra ID, or rely on PKI for certificate-based trust. In each case, you can eliminate the operational overhead of managing static credentials at scale.

Choose your implementation path:

For sample code and infrastructure templates, clone the
Amazon MQ samples repository and deploy the CDK stack for your chosen authentication method.

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.

Vignesh Selvam

Vignesh Selvam

Vignesh is the Principal Product Manager for Amazon MQ at AWS. He works with customers to solve their messaging needs and with the open-source communities for innovating with message brokers. Prior to joining AWS, he built products for security and analytics.

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.

Trace cascading decision failures with a blame graph on Amazon OpenSearch Service

Post Syndicated from Jon Handler original https://aws.amazon.com/blogs/big-data/trace-cascading-decision-failures-with-a-blame-graph-on-amazon-opensearch-service/

Multi-agent systems are straightforward to build but hard to debug. You chain a few agents together, each one does its part, and most of the time it works. When it doesn’t, you’re left with a large volume of logs. They tell you what every agent said, but nothing about which agent caused the bad outcome.

Working with AWS customers building multi-agent systems, we kept seeing the same problem. A pipeline of agents decides, the decision turns out wrong, and no one can say which agent caused it. The logs are complete, but they don’t answer that question. So we, two AWS Solutions Architects, built a stock-research pipeline to reproduce it and show a solution approach.

Five agents work in sequence, and the last one makes a BUY, SELL, or HOLD call. In our test cases, the agent kept recommending BUY, and the positions kept losing money. Every step was logged. The logs still didn’t tell us who broke the pipeline.

In this post, we show you how to build a blame graph that traces which agent caused a failure in a multi-agent pipeline, using Amazon OpenSearch Service for graph storage and Amazon Bedrock for embeddings and reasoning.

Prerequisites

You must have the following prerequisites to follow along with this post.

  • Download the source code from the GitHub repository: It includes everything needed to set up and run the demo end to end:
    • The five-agent pipeline.
    • The instrumentation layer.
    • AWS CloudFormation template.
    • OpenSearch UI dashboard export
    • Sample data.
    • Step-by-step setup instructions (README.md, DEPLOYMENT.md).
  • An AWS account with access to Amazon Bedrock (Anthropic Claude Sonnet 4.5 and Amazon Titan Text Embeddings V2 enabled in us-west-2).
  • An OpenSearch Service domain.
  • Python 3.11+.
  • AWS Command Line Interface (AWS CLI) v2 configured with valid credentials.

The challenge

The pipeline is a chain of five agents. A Researcher gathers the facts, a Risk Analyst weighs the downside, a Valuation Analyst runs the numbers, and a Macro Economist sets up the market backdrop. Each one builds on the output of the agents before it. The Strategist (AI agent) sits at the end and turns all of it into a single call: BUY, SELL, or HOLD.

We set up three failures, each one a pattern common in production agent deployments (hallucinated facts from retrieval, suppressed minority signals, stale data from delayed ingestion):

  • A hallucination. The Researcher invents a company partnership that doesn’t exist.
  • A buried warning. The Risk Analyst flags a regulatory risk and gets outvoted.
  • Stale data. The Researcher misses a filing published three days earlier.

We engineered each failure deterministically, so the demo is reproducible and has a known answer. For each scenario, we hand-authored the five agents’ outputs as fixed JavaScript Object Notation (JSON). The pipeline replays these outputs while the instrumentation computes embeddings, influence, and blame live. We recorded a ground-truth root cause (for example, researcher for hallucination).

In every case, the pipeline recommends BUY, and the position drops. Standard logging records each agent’s output, but it can’t tell you which claim drove the final decision. Closing the gap between logging and root-cause attribution is what we set out to do.

Solution

We treat agent reasoning as a graph and measure influence between agents, then walk that graph backward from the failed decision to find the root cause.

Three services make up the stack:

  • Strands Agents runs the five-agent pipeline.
  • Amazon Bedrock provides the models: Amazon Titan Text Embeddings V2 to embed each claim, and Anthropic Claude Sonnet 4.5 for agent reasoning and the incident write-up.
  • An OpenSearch UI application, an analytics interface hosted in the AWS Cloud with a single endpoint, connects to the domain as a data source and serves the dashboard, Discover, and the Dev Tools console we use to investigate. 

Here is how blame attribution works. Every claim an agent makes becomes a document with an Amazon Titan embedding. When a downstream agent cites something, we measure the cosine similarity between that citation and each upstream claim. Cosine similarity becomes the influence one agent had on another.

We store these as edges. To find the root cause, we start at the failed decision and walk backward through the edges. Whoever contributed the most gets the most blame.

Alongside the graph we record three things per run: an explainability score for how much of the decision traces back to evidence, the confidence of the attribution, and whether a dissenting agent was overruled.

A note on method: there is no industry standard yet for root-cause attribution in multi-agent large language model (LLM) pipelines. Our approach combines two established ideas: a credit assignment (attributing an outcome to the steps that produced it) and embedding similarity for tracing how claims propagate, with an LLM-as-a-judge style check. The metrics here (influence, explainability) are pragmatic, reproducible measures we define in this post, not standardized benchmarks.

Architecture

Five parts make up the flow:

  • Agents run on the Strands Agents, with reasoning on Claude Sonnet 4.5.
  • An instrumentation layer extracts each claim, embeds it with Amazon Titan Text Embeddings V2, scores influence with cosine similarity, runs the backward traversal, and generates an incident report.
  • Amazon OpenSearch Service holds seven indices, including the claims index with k-nearest neighbor (kNN) vectors and the blame, metrics, and incident indices.
  • Analysts review the results in the OpenSearch UI application (the dashboard, Discover, and the Dev Tools console), launched from the Amazon OpenSearch Service console.
  • We use OpenSearch UI rather than the domain’s built-in dashboards. Because OpenSearch UI is hosted in the AWS Cloud, the application stays available during domain maintenance and can bring multiple data sources into one view. The pipeline still writes to the domain, and OpenSearch UI reads it as a registered data source. 
Five-agent pipeline: Researcher, Risk Analyst, Valuation, Macro Economist, Strategist in sequence, ending at BUY decision.

Figure 1a: The five-agent runtime pipeline

Instrumentation layer sending embeddings and blame edges to Amazon OpenSearch Service, with Amazon Bedrock providing Amazon Titan and Claude models.

Figure 1b: The instrumentation and OpenSearch Service data plane

Walking through a failure

We ran the pipeline nine times, three runs per scenario, on an Amazon OpenSearch Service domain running OpenSearch 2.17. The decision under investigation is the final BUY. We know it failed because each scenario carries a ground-truth outcome: the position lost money. The failure is the known bad outcome we trace backward from, not something the system infers.  Everything the pipeline produces is a document you can query, so the investigation is a series of queries we run from the Dev Tools console in the OpenSearch UI application. 

To follow along, launch the OpenSearch UI application from the Amazon OpenSearch Service console, open your workspace, and choose Dev Tools (near the bottom of the left navigation panel). Paste each query below into the left pane and choose the run button. Every query in this section is in the repository at devtools_queries.md, in the same order as the walkthrough, so you can copy them from there instead of retyping. The equivalent queries as Python are in queries.py. 

Start with the outcome

Every run is a BUY, and every loss is negative, down to 72 percent. Standard logging stops here. You know it failed, but you don’t know who to fix.

Dev Tools query results showing nine pipeline runs, all recommending BUY with losses from 58% to 72%.

Figure 2: Pipeline run results: all nine runs recommend BUY with losses up to 72%

Next, look at who influenced whom

Among all agents, the Researcher sources the most edges. Nearly every node downstream gets its data from the Researcher, making it the first place to look. A lead, not a verdict.

Dev Tools aggregation showing influence edges by source agent; Researcher has the most edges.

Figure 3: Influence edges aggregated by source agent

Query the blame metrics for each scenario

Blame lands on the Researcher, with a score around 0.45, and the attribution is correct on all three runs. A fabricated partnership flowed straight into the final BUY. Stale-data scenario behaves the same way: the Researcher again, at 0.46, correct.

Dev Tools query showing root cause attribution: Researcher at 0.45 for hallucination scenario.

Figure 4: Root cause attribution for hallucination runs

Here are the raw edges in Discover, sorted from highest influence to lowest

In the OpenSearch UI application, choose Discover and select the agent-blame index pattern, then set the time range to Last 30 days and sort by influence_score descending. Each row is one edge- a claim passed from a source agent (source_agent.agent_id) to a downstream agent (target_agent.agent_id), scored by how strongly it shaped that agent’s output. The top rows are the highest-influence edges: the ones that most shaped the final BUY.

Discover view of blame edges sorted by influence score, highest to lowest.

Figure 5: Blame edges sorted by influence score

When attribution is hard

It’s the buried-warning scenario that the graph gets wrong, and it’s the most useful result in the post.

The Risk Analyst was right. It flagged the regulatory risk. The Strategist saw the warning, weighted it at 0.15, and bought it anyway. Who actually failed? The Strategist.

But the blame graph points at the Risk Analyst, with the highest score in that run at 0.37. Why? Our method measures influence, and the dissent is a distinct claim that the method traces directly, so it scores high. Influence is not the same as responsibility.

Why did the Strategist ignore it? In the scenario, the Strategist acknowledged the dissent but reasoned that the strength of the clinical data made the compound “differentiated” from past failures. It weighted that bullish evidence at 0.85 against the Risk Analyst’s 0.15. The Strategist rationalized the warning away instead of treating high-confidence, time-bound regulatory risk as a hard stop. The model recorded that reasoning, which is exactly why we can see how the dissent was discounted.

This gap between influence and responsibility is why we track dissent on our own.

Dissent was present, acknowledged, and weighted at 0.15. A flag catches what the graph misses: a valid warning was heard and then ignored. One signal is not enough. Influence tells you what is propagated. Dissent flags tell you what was wrongly dismissed. You need both.

Dev Tools query showing suppressed dissent: dissent_weight_given 0.15, dissent_suppressed true.

Figure 6: Suppressed dissent detection

Reviewing the metrics dashboard

OpenSearch UI rolls up all nine runs. To open it, launch the OpenSearch UI application, open your workspace, and choose Dashboards in the left navigation, then open the Multi-Agent Blame Game — Observability dashboard. Set the time range to Last 30 days to see all nine runs. If you haven’t imported it yet, go to Manage Workspace and choose Import under Assets. Upload blame-game-dashboard.ndjson from the repository, mapping the index patterns to your domain’s data source.

Full OpenSearch metrics dashboard with panels for root cause, explainability, loss, influence, and propagation.

Figure 7: Full metrics dashboard

Each panel earns its place. A few are worth calling out. Root cause distribution flags the Researcher six times and the Risk Analyst three times. That Risk Analyst slice is the dissent misattribution from earlier, not a real culprit.

Root cause distribution: Researcher in 6 runs, Risk Analyst in 3 (misattribution).

Figure 8: Root cause distribution

Explainability averages 0.826, a metric we define, not a standard score.

Explainability score averaging 0.826 across nine runs.

Figure 9: Explainability score

Preventable loss versus realized loss splits the damage attribution can pin on one agent from the damage it can’t. And average influence clusters rather than spikes, showing no single cause. That is the whole reason attribution sums influence instead of trusting one edge.

Preventable loss panel showing dollar amounts attributed to root-cause agent per scenario.

Figure 10: Preventable loss

Realized loss panel showing total financial damage across all runs before attribution.

Figure 11: Realized loss

Blame and loss comparison table: hallucination and stale-data rows show small errors. Dissent row shows largest gap.

Figure 12: Blame and loss table

Average influence by source agent: scores cluster between 0.29 and 0.42, no single spike.

Figure 13: Average influence by source agent

Propagation type breakdown: most edges weak or independent, few amplified.

Figure 14: Propagation type breakdown

Exploring it interactively

For demos we wrapped the same pipeline in a small Streamlit app. To run it, from the repository root install the dependencies and start the app:   

source .env 
streamlit run src/app.py --server.address localhost

It opens in your browser at http://localhost:8501. It runs two ways: pick a prepared scenario and replay it, or type in a company of your own and have the five agents run live on Amazon Bedrock against it. Either way you watch the agents execute, and the blame graph form, with the verdict and the incident narrative on one screen. A History tab reads the metrics index, so you can review past runs without leaving the app. 

A live run has no ground truth, so the app doesn’t claim the attribution is right or wrong. You just see where the influence landed. The prepared scenarios are still the way to demonstrate a specific, known failure. 

Streamlit demo app showing a pipeline run with agent panels, claims, and blame verdict.

Figure 15: Streamlit demo app

Explaining every decision: The evidence each agent weighed

Blame attribution is only useful if you can see the evidence behind it. Every claim an agent makes is stored with the confidence the agent assigned and the source it came from. Sources include an SEC filing, a clinical trial registry, an FDA page, or an earnings call. A blame score is never a bare number. You can open any agent and read the exact claims and sources it weighed before it spoke.

Consider the final decision as the clearest example. The Strategist doesn’t only emit a BUY. The Strategist records which upstream claim it relied on and how much weight it gave each one. Recording those weights turns the last step from a black box into a list of citations you can audit.

Explainability captures exactly that. A high score means most of the recommendation traces back to specific, sourced claims rather than to unexplained reasoning. It is the difference between the model said BUY and the model said BUY because of these claims, from these sources, weighted this way.

Streamlit app detail: Financial Researcher claims expanded with confidence scores and sources.

Figure 16: Per-agent evidence and reasoning for the BioGenX run

Performance and results

Across nine runs, the system identified the correct root cause six times, or 67 percent. The three misses are all the buried-warning scenarios, where influence and responsibility diverge. We would rather report the real number and explain the miss than round it up.

A full run takes about 25 seconds from end to end. Almost all of that is the Bedrock calls: about 74 embeddings per run plus one Claude write-up.

Attribution alone, the part that walks the graph and assigns blame, runs in about 74 milliseconds. That is cheap enough to run on every pipeline execution, not only after something goes wrong.

End-to-end latency chart: full run about 25 seconds, attribution step about 74 milliseconds.

Figure 17: End-to-end latency scenario

What this means for building agent pipelines

Our data points at three concrete changes:

  • Make the Researcher cross-check any major claim against a second source.
  • Give the Strategist a hard rule so a high-confidence dissent near a binary event can’t be overridden silently.
  • Add a freshness check so old data can’t drive a decision.

More broadly, treat influence and responsibility as separate questions. Measure both. A blame graph is a strong default for tracing propagation, but you need side signals like dissent suppression to catch up on the cases it can’t see.

From detection to prevention: Guardrails that stop the loss

Attribution tells you who broke a run after the fact. The same signals can stop the break before anyone acts on it. We added a guardrail layer that sits between the pipeline’s decision and the action, and overrides the call when a known failure pattern appears. The demo implements this layer (run with --guardrails, or toggle it in the app). It answers the question of whether the fixes are in the code: they are.

Guardrail gate diagram: blame signals feed three checks (dissent-override, source cross-check, freshness) before decision passes or is held.

Figure 18: The guardrail gate between the decision and the action

Each guardrail targets one of the three failure modes:

  • Dissent-override (Strategist): When the Risk Analyst raises a high-confidence dissent near a binary event and the Strategist under-weights it, the decision is forced to the safe action (HOLD).
  • Source cross-check (Researcher): A material claim resting on a single self-reported source cannot drive a BUY. It must be corroborated, or the call is held.
  • Freshness (Researcher): If material information was published just before the analysis and was not reflected in the inputs, the call is held.

With all three enabled, every scenario that previously issued a losing BUY is caught and held. Across the three runs that is about USD $3.79 million of illustrative loss prevented.

Guardrails effect: approximately $3.79M illustrative loss converted from realized to prevented.

Figure 19: Guardrails convert realized loss into prevented loss

The three guardrails are demo-grade heuristics, and we want to be explicit about that. The single-source and freshness checks work here only because the scenario data is engineered with known sources and dates. They are proxies, not real controls. A well-formed hallucination with a plausible citation would pass without detection of the cross-check, and the freshness rule only knows about data it is handed.

To make this production-grade, replace each proxy with real control.

For hallucinations, don’t count sources. Verify the material claims a decision rests on against a trusted source such as a knowledge base in Amazon OpenSearch Service or an authoritative filings and market-data API. Use an entailment or LLM-as-a-judge check to confirm the evidence actually supports the claim, requiring corroboration from independent sources before a claim can drive a BUY.

For freshness, wire in a live data feed and a scheduled-catalyst calendar. Hold whenever a decision rests on inputs that predate a material update or sits too close to a binary event. For dissent, keep the override but calibrate its threshold on historical outcomes and route borderline, high-value calls to a human rather than auto-deciding.

Underneath all of it, store the rules and thresholds as versioned policy in OpenSearch Service. Keep the blame graph running so you can confirm the guardrails fire for the right reasons. Log every override for audit, and evaluate the whole layer on real outcomes. Watch the false-positive rate as closely as the catches, because a guardrail that blocks good trades is only a new failure mode. Stay conservative: Prefer holding a good trade to taking a bad one, and make every block explainable.

Responsible AI considerations

This solution uses Amazon Titan Text Embeddings V2 and Anthropic Claude Sonnet 4.5 for agent reasoning and incident narrative generation. LLM-generated blame attributions and incident reports are informational aids, not authoritative verdicts. Always pair automated attribution with human review before making operational decisions. The influence score measures semantic similarity between claims, not true causation. The buried-warning scenario in this post demonstrates exactly where that distinction matters.

All company names, financial figures, and scenarios are fictional. No real market data or customer information is used. The stock-research pipeline is an illustrative vehicle for demonstrating blame attribution and observability. It isn’t investment advice, and the BUY/SELL/HOLD outputs are not stock recommendations. Don’t use this system, as built, to make financial or investment decisions.

Before adapting this approach to production pipelines, validate attribution accuracy against your own ground-truth data and implement safeguards appropriate to your risk level. The guardrails module in this repo is a starting point, not a complete solution. For more information, see Responsible AI with AWS.

Clean up

To avoid ongoing charges, delete the Amazon OpenSearch Service domain when you are done. The demo uses a single CloudFormation stack, so one command removes everything.

aws cloudformation delete-stack --region us-west-2 --stack-name blame-game-demo

Amazon Bedrock is billed per request, so there is nothing to tear down there.

Conclusion

Multi-agent pipelines fail in ways logs can’t explain. By embedding each claim with Amazon Titan Text Embeddings V2, scoring influence in Amazon OpenSearch Service, and walking the graph backward from the failed decision, we turned “something broke” into “here is the agent that broke it, and here is the evidence.” We also showed where that approach falls short, and the extra signal that covers it.

Code, queries, and deployment steps are in the repository. The hard part isn’t the infrastructure. It’s deciding to measure influence and responsibility as two different things.

Learn more

To dive deeper, get the full source in the GitHub repository, and see the Amazon OpenSearch Service and Amazon Bedrock documentation to adapt this to your own pipelines.


About the authors

Jon Handler

Jon Handler

Jon is a Senior Principal Solutions Architect for Search Services at Amazon Web Services. Jon works closely with OpenSearch and Amazon OpenSearch Service, providing help and guidance to a broad range of customers who have search and log analytics workloads. Prior to joining AWS, Jon’s career ranged across distributed systems and search at startups and large organizations. His career as a software developer included four years of coding a large-scale, eCommerce search engine.

Smita Singh

Smita Singh

Smita is a Senior Solutions Architect at AWS. She comes with 20 years of experience in the industry. She focuses on defining technical strategic vision and works on architecture, design, and implementation of modern, scalable platforms for large-scale global enterprises and SaaS providers. She specializes in architecture and implementation of large-scale platform solutions for global enterprises and SaaS providers, with a focus on data, analytics, and generative AI workloads.

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.