diff --git a/Directory.Packages.props b/Directory.Packages.props
index f1fa135a5c4..43cad480473 100644
--- a/Directory.Packages.props
+++ b/Directory.Packages.props
@@ -4,10 +4,12 @@
true
+
+
@@ -112,7 +114,7 @@
-
+
diff --git a/Orleans.slnx b/Orleans.slnx
index 002a00aff92..158099d0792 100644
--- a/Orleans.slnx
+++ b/Orleans.slnx
@@ -82,6 +82,7 @@
+
diff --git a/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHost.csproj b/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHost.csproj
index 6f74f623a39..344177e41c5 100644
--- a/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHost.csproj
+++ b/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHost.csproj
@@ -21,6 +21,11 @@
+
+
+
+
diff --git a/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHostExamples.cs b/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHostExamples.cs
index fe82b979ee1..0dede7ce1f4 100644
--- a/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHostExamples.cs
+++ b/docs/site/src/content/docs/host/snippets/aspire/AppHost/AppHostExamples.cs
@@ -1,5 +1,7 @@
+using Amazon;
using Aspire.Hosting;
using Aspire.Hosting.Azure;
+using Aspire.Hosting.Orleans;
namespace Orleans.Docs.Snippets.Aspire;
@@ -8,6 +10,41 @@ namespace Orleans.Docs.Snippets.Aspire;
public static class AppHostExamples
{
+ //
+ public static void SqsStreaming(string[] args)
+ {
+ var builder = DistributedApplication.CreateBuilder(args);
+
+ var aws = builder.AddAWSSDKConfig()
+ .WithRegion(RegionEndpoint.USEast1);
+
+ var orleans = builder.AddOrleans("cluster")
+ .WithDevelopmentClustering()
+ .WithMemoryGrainStorage("PubSubStore")
+ .WithSqsStreaming(
+ "Orders",
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = "orders-service",
+ PartitionCount = 16,
+ FifoQueue = true,
+ ReceiveWaitTimeSeconds = 20,
+ VisibilityTimeoutSeconds = 60,
+ });
+
+ var silo = builder.AddProject("silo")
+ .WithReference(orleans)
+ .WithReplicas(3);
+
+ builder.AddProject("client")
+ .WithReference(orleans.AsClient())
+ .WaitFor(silo);
+
+ builder.Build().Run();
+ }
+ //
+
//
public static void BasicOrleansCluster(string[] args)
{
@@ -263,4 +300,5 @@ public static void ExplicitClusterIds(string[] args)
builder.Build().Run();
}
//
+
}
diff --git a/docs/site/src/content/docs/host/snippets/aspire/Client/Client.csproj b/docs/site/src/content/docs/host/snippets/aspire/Client/Client.csproj
index 2f0f625952c..19a2a367d79 100644
--- a/docs/site/src/content/docs/host/snippets/aspire/Client/Client.csproj
+++ b/docs/site/src/content/docs/host/snippets/aspire/Client/Client.csproj
@@ -16,7 +16,6 @@
-
diff --git a/docs/site/src/content/docs/host/snippets/aspire/Silo/Silo.csproj b/docs/site/src/content/docs/host/snippets/aspire/Silo/Silo.csproj
index 8b5ee5962a0..05b85281b42 100644
--- a/docs/site/src/content/docs/host/snippets/aspire/Silo/Silo.csproj
+++ b/docs/site/src/content/docs/host/snippets/aspire/Silo/Silo.csproj
@@ -27,7 +27,6 @@
-
diff --git a/docs/site/src/content/docs/resources/nuget-packages.md b/docs/site/src/content/docs/resources/nuget-packages.md
index 97780fdc996..03160f13e85 100644
--- a/docs/site/src/content/docs/resources/nuget-packages.md
+++ b/docs/site/src/content/docs/resources/nuget-packages.md
@@ -92,6 +92,7 @@ Use reminders for recurring durable callbacks and durable jobs for scheduled one
| [Microsoft.Orleans.Streaming.NATS](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.NATS) | NATS JetStream streams. |
| [Microsoft.Orleans.Streaming.Redis](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.Redis) | Redis Streams integration. |
| [Microsoft.Orleans.Streaming.SQS](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS) | Amazon SQS streams. |
+| [Microsoft.Orleans.Streaming.SQS.Aspire](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS.Aspire) | .NET Aspire and AWS CDK integration for Orleans SQS streams. |
| [Microsoft.Orleans.BroadcastChannel](https://www.nuget.org/packages/Microsoft.Orleans.BroadcastChannel) | Lightweight broadcast channels. |
## Grain directories
diff --git a/docs/site/src/content/docs/snippets/compiled/Streaming/SqsSnippets.cs b/docs/site/src/content/docs/snippets/compiled/Streaming/SqsSnippets.cs
index d7789ae8b5d..3b127bb576b 100644
--- a/docs/site/src/content/docs/snippets/compiled/Streaming/SqsSnippets.cs
+++ b/docs/site/src/content/docs/snippets/compiled/Streaming/SqsSnippets.cs
@@ -1,4 +1,5 @@
using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Hosting;
using Orleans.Hosting;
using Orleans.Serialization;
using Orleans.Streaming.SQS.Streams;
@@ -7,6 +8,24 @@ namespace Documentation.Streaming;
internal static class SqsSnippets
{
+ internal static void ConfigureAspireSilo(string[] args)
+ {
+ //
+ var builder = Host.CreateApplicationBuilder(args);
+ builder.UseOrleans();
+ builder.Build().Run();
+ //
+ }
+
+ internal static void ConfigureAspireClient(string[] args)
+ {
+ //
+ var builder = Host.CreateApplicationBuilder(args);
+ builder.UseOrleansClient();
+ builder.Build().Run();
+ //
+ }
+
internal static void ConfigureSilo(ISiloBuilder siloBuilder)
{
//
diff --git a/docs/site/src/content/docs/streaming/sqs-streaming.md b/docs/site/src/content/docs/streaming/sqs-streaming.md
index 96ff7249f28..6ffd09d6cf0 100644
--- a/docs/site/src/content/docs/streaming/sqs-streaming.md
+++ b/docs/site/src/content/docs/streaming/sqs-streaming.md
@@ -9,7 +9,7 @@ ms.topic: how-to
The [`Microsoft.Orleans.Streaming.SQS`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS) package connects Orleans persistent streams to [Amazon Simple Queue Service](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/welcome.html). Orleans maps streams across a configurable set of SQS queues, creates a queue when its mapped partition is first used, receives batches through persistent-stream pulling agents, and deletes messages after successful delivery.
-SQS streams provide [at-least-once delivery](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/standard-queues-at-least-once-delivery.html). A message becomes visible again when its [visibility timeout](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-visibility-timeout.html) expires before Orleans acknowledges it, so consumers must handle duplicates. The provider assigns receiver-local sequence tokens as messages arrive and doesn't support rewind to an earlier token.
+SQS streams provide [at-least-once delivery](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/standard-queues-at-least-once-delivery.html). A message becomes visible again when its [visibility timeout](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/sqs-visibility-timeout.html) expires before Orleans acknowledges it, so consumers must handle duplicates. The provider assigns receiver-local sequence tokens as messages arrive, and subscriptions consume current SQS deliveries rather than seeking to an earlier token.
## Configure a standard queue provider
@@ -25,6 +25,30 @@ Configure each Orleans client which publishes through the provider with the same
The `Service` connection value accepts an AWS region such as `us-east-1` or an SQS-compatible endpoint such as `http://localhost:4566`. With a region, the provider uses the [AWS SDK for .NET credential resolution chain](https://docs.aws.amazon.com/sdk-for-net/v4/developer-guide/creds-assign.html) when the connection string contains no explicit credentials. Prefer workload credentials such as an IAM role. A protected connection string can supply `AccessKey`, `SecretKey`, and `SessionToken` when the deployment requires explicit temporary credentials.
+## Configure SQS streams with Aspire
+
+Install [`Microsoft.Orleans.Streaming.SQS.Aspire`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS.Aspire) in the AppHost. Install [`Microsoft.Orleans.Streaming.SQS`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS) in every silo and Orleans client project that uses the provider. The AppHost package's `WithSqsStreaming` extension uses the AWS-supported [`Aspire.Hosting.AWS`](https://www.nuget.org/packages/Aspire.Hosting.AWS) integration to configure AWS SDK for .NET v4 and provision the complete Orleans queue topology through AWS CDK.
+
+Configure the AWS SDK region and one `SqsStreamingOptions` object:
+
+:::code language="csharp" source="../host/snippets/aspire/AppHost/AppHostExamples.cs" id="sqs_streaming_apphost":::
+
+The integration uses those options for both the CDK queue resources and the `Orleans:Streaming:Orders` configuration emitted to every referenced silo and client. The provisioned queue names therefore match the runtime topology exactly. The AWS SDK reference supplies `AWS_PROFILE` and `AWS_REGION`; workload identity and shared AWS configuration remain available through the SDK credential chain.
+
+Provider names are unique ignoring case within each Orleans service. Registration rejects names such as `orders_primary` and `Orders_primary` before changing provider configuration or creating resources. It also checks physical queue ownership across SQS stacks in the same AppHost and region: the service ID, lowercased provider name, partition number, and optional `.fifo` suffix identify each queue. Use distinct service IDs or provider names for separate topologies in the same region, including when AWS profiles select different accounts. Hyphens and underscores remain distinct in physical queue names, so `orders-primary` and `orders_primary` can coexist.
+
+The extension applies the stable service ID and makes each referenced silo and client wait for the CDK stack automatically. The stack uses CloudFormation in the account and region selected by the AWS SDK configuration. Grant the AppHost credentials permission to create and update the stack; bootstrap the environment when added constructs use CDK assets. See [Provisioning application resources with AWS CDK](https://github.com/aws/integrations-on-dotnet-aspire-for-aws/blob/main/src/Aspire.Hosting.AWS/README.md#provisioning-application-resources-with-aws-cdk) for the integration contract.
+
+The silo activates generated `Orleans:Streaming:Orders` configuration through :
+
+:::code language="csharp" source="../snippets/compiled/Streaming/SqsSnippets.cs" id="sqs_streaming_silo":::
+
+The client activates the matching publishing provider through :
+
+:::code language="csharp" source="../snippets/compiled/Streaming/SqsSnippets.cs" id="sqs_streaming_client":::
+
+For SQS-compatible local services, emit `ServiceEndpoint=http://localhost:9324` in the provider configuration. Provider configuration also accepts `Region`, `ConnectionString`, `ReceiveWaitTimeSeconds`, `VisibilityTimeoutSeconds`, indexed `ReceiveMessageAttributes` and `ReceiveMessageSystemAttributes`, and a keyed `ISQSDataAdapter` through `DataAdapterKey`.
+
## Preserve per-stream order with FIFO queues
Set to use SQS FIFO queues:
@@ -60,7 +84,7 @@ List every application-defined attribute required by the decoder in | Requests the application-defined attributes consumed by a custom data adapter. |
| | Requests SQS system attributes consumed by the provider or application adapter. |
-Queue-creation settings apply when the queue is absent. Manage changes to retention, visibility, encryption, access policy, and dead-letter redrive policy through SQS administration for existing queues.
+With the Aspire integration, the CDK stack owns queue creation and updates. Define retention, visibility, encryption, access policy, and dead-letter redrive policy in the queue constructs so CloudFormation applies the declared configuration consistently. In non-Aspire hosts, Orleans creates a missing mapped queue when its partition is first used; operators manage subsequent queue-property changes through SQS administration.
Serialized Orleans batches must fit within the [SQS message quotas](https://docs.aws.amazon.com/AWSSimpleQueueService/latest/SQSDeveloperGuide/quotas-messages.html). Bound batch size before messages approach the service limit, and include custom envelope and message-attribute overhead in that calculation.
diff --git a/docs/site/src/content/docs/streaming/stream-providers.md b/docs/site/src/content/docs/streaming/stream-providers.md
index b74a1624f04..0d589bcd9c7 100644
--- a/docs/site/src/content/docs/streaming/stream-providers.md
+++ b/docs/site/src/content/docs/streaming/stream-providers.md
@@ -1,7 +1,7 @@
---
title: Orleans stream providers
description: Compare built-in Orleans stream providers by durability, rewindability, status, and prerequisites.
-ms.date: 08/18/2026
+ms.date: 08/26/2026
ms.topic: concept-article
---
@@ -18,6 +18,7 @@ A stream provider connects the Orleans streaming API to a transport and defines
| Azure Event Hubs | [`Microsoft.Orleans.Streaming.EventHubs`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.EventHubs) | Stable | Yes, within Event Hubs retention | Yes | Event Hubs namespace, hub, consumer group, and checkpoint storage |
| Amazon Kinesis | [`Microsoft.Orleans.Streaming.Kinesis`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.Kinesis) | Stable | Yes, within Kinesis retention | Yes | Kinesis data stream, AWS credentials, region, and durable checkpoint storage |
| Amazon SQS | [`Microsoft.Orleans.Streaming.SQS`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS) | Stable | Yes, within SQS retention | No | AWS account, queue permissions, region/endpoint configuration |
+| Amazon SQS Aspire AppHost | [`Microsoft.Orleans.Streaming.SQS.Aspire`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS.Aspire) | Stable | Uses the SQS provider's queues | No | .NET Aspire AppHost, AWS CDK, AWS credentials, and region configuration |
| ADO.NET | [`Microsoft.Orleans.Streaming.AdoNet`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.AdoNet) | **Alpha** | Yes, in relational tables until expiry/dead-letter eviction | No | Supported database, ADO.NET driver, and Orleans streaming SQL schema |
| NATS JetStream | [`Microsoft.Orleans.Streaming.NATS`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.NATS) | **Alpha** | Configurable; file storage is the default | No | NATS server with JetStream and sufficient storage; subject/stream administration |
| Redis Streams | [`Microsoft.Orleans.Streaming.Redis`](https://www.nuget.org/packages/Microsoft.Orleans.Streaming.Redis) | **Alpha** | Configurable through Redis persistence and stream retention | Yes, while entries remain | Redis deployment, persistence/HA policy, and retention sizing |
diff --git a/docs/site/src/data/external-link-allowlist.json b/docs/site/src/data/external-link-allowlist.json
index a8932606fa7..9e450e20221 100644
--- a/docs/site/src/data/external-link-allowlist.json
+++ b/docs/site/src/data/external-link-allowlist.json
@@ -13,6 +13,7 @@
"https://docs.aws.amazon.com/streams/latest/dev/introduction.html": "AWS serves this public Kinesis overview to browsers but rejects the bounded automated HEAD/GET probe with HTTP 403.",
"https://www.nuget.org/packages/Microsoft.Orleans.Journaling.S3": "The new S3 journaling package is documented but not yet published; remove this entry and its unpublished API-package entry after publication.",
"https://en.wikipedia.org/wiki/Kalman_filter": "Wikipedia serves this public article to browsers but rejects the bounded automated HEAD/GET probe with HTTP 403.",
+ "https://www.nuget.org/packages/Microsoft.Orleans.Streaming.SQS.Aspire": "This package is introduced by the current source and will be available when the corresponding Orleans release is published.",
"https://www.f5.com/company/blog/nginx/nginx-power-of-two-choices-load-balancing-algorithm": "The F5 site serves this canonical NGINX article to browsers but rejects the bounded automated HEAD/GET probe with HTTP 403."
}
}
diff --git a/docs/site/src/data/unpublished-api-packages.json b/docs/site/src/data/unpublished-api-packages.json
index b523ea70c0f..d00c1801080 100644
--- a/docs/site/src/data/unpublished-api-packages.json
+++ b/docs/site/src/data/unpublished-api-packages.json
@@ -1,6 +1,7 @@
{
"description": "Generated API assemblies which are not currently published as standalone NuGet packages.",
"packages": {
- "Microsoft.Orleans.Journaling.S3": "The new provider is awaiting its first NuGet publication."
+ "Microsoft.Orleans.Journaling.S3": "The new provider is awaiting its first NuGet publication.",
+ "Microsoft.Orleans.Streaming.SQS.Aspire": "This package is introduced by the current source and will be published with the corresponding Orleans release."
}
}
diff --git a/src/AWS/Orleans.Streaming.SQS.Aspire/Orleans.Streaming.SQS.Aspire.csproj b/src/AWS/Orleans.Streaming.SQS.Aspire/Orleans.Streaming.SQS.Aspire.csproj
new file mode 100644
index 00000000000..caa8e6a5b4d
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS.Aspire/Orleans.Streaming.SQS.Aspire.csproj
@@ -0,0 +1,26 @@
+
+
+ README.md
+ Microsoft.Orleans.Streaming.SQS.Aspire
+ Microsoft Orleans Aspire integration for AWS SQS Streaming
+ Provisions and configures Orleans SQS stream providers in .NET Aspire AppHost projects.
+ $(PackageTags) Aspire AWS SQS CDK
+ $(DefaultTargetFrameworks)
+ false
+ false
+ false
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/src/AWS/Orleans.Streaming.SQS.Aspire/OrleansSqsStreamingExtensions.cs b/src/AWS/Orleans.Streaming.SQS.Aspire/OrleansSqsStreamingExtensions.cs
new file mode 100644
index 00000000000..edf7481a1fc
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS.Aspire/OrleansSqsStreamingExtensions.cs
@@ -0,0 +1,352 @@
+using System.Runtime.CompilerServices;
+using System.Security.Cryptography;
+using System.Text;
+using System.Text.RegularExpressions;
+using Amazon.CDK;
+using Amazon.CDK.AWS.SQS;
+using Aspire.Hosting.ApplicationModel;
+using Aspire.Hosting.AWS;
+using Aspire.Hosting.AWS.CDK;
+using Aspire.Hosting.Orleans;
+using OrleansAWSUtils.Storage;
+
+namespace Aspire.Hosting;
+
+///
+/// Extension methods for configuring Orleans SQS streaming in .NET Aspire.
+///
+public static partial class OrleansSqsStreamingExtensions
+{
+ ///
+ /// Adds an AWS CDK-provisioned SQS stream provider to an Orleans service.
+ ///
+ /// The Orleans service.
+ /// The stream provider name.
+ /// The AWS SDK profile and region configuration.
+ /// The SQS topology and runtime options.
+ /// The Orleans service.
+ ///
+ /// A stream provider name conflicts case-insensitively, a physical queue is already owned in the
+ /// same AppHost and region, a generated resource name is already registered, or the service
+ /// identifier conflicts with the Orleans service.
+ ///
+ public static OrleansService WithSqsStreaming(
+ this OrleansService orleansService,
+ string name,
+ IAWSSDKConfig awsSdkConfig,
+ SqsStreamingOptions options)
+ {
+ AddSqsStreaming(orleansService, name, awsSdkConfig, options);
+ return orleansService;
+ }
+
+ ///
+ /// Adds an AWS CDK-provisioned SQS stream provider to an Orleans service and returns its resource model.
+ ///
+ /// The Orleans service.
+ /// The stream provider name.
+ /// The AWS SDK profile and region configuration.
+ /// The SQS topology and runtime options.
+ /// The SQS stream provider resource.
+ ///
+ /// A stream provider name conflicts case-insensitively, a physical queue is already owned in the
+ /// same AppHost and region, a generated resource name is already registered, or the service
+ /// identifier conflicts with the Orleans service.
+ ///
+ public static SqsStreamingResource AddSqsStreaming(
+ this OrleansService orleansService,
+ string name,
+ IAWSSDKConfig awsSdkConfig,
+ SqsStreamingOptions options)
+ {
+ ArgumentNullException.ThrowIfNull(orleansService);
+ ArgumentException.ThrowIfNullOrWhiteSpace(name);
+ ArgumentNullException.ThrowIfNull(awsSdkConfig);
+ ArgumentNullException.ThrowIfNull(options);
+ if (name.Contains("__", StringComparison.Ordinal))
+ {
+ throw new ArgumentException(
+ "SQS stream provider names must form one configuration path segment; use a separator other than '__'.",
+ nameof(name));
+ }
+
+ var validatedOptions = ValidateAndCopy(name, awsSdkConfig, options);
+ var queueNames = GetPhysicalQueueNames(name, validatedOptions);
+ var region = awsSdkConfig.Region!.SystemName;
+ ValidateIdentities(orleansService, name, region, queueNames);
+ ValidateServiceId(orleansService, validatedOptions.ServiceId, allowUnset: true);
+
+ var resourceName = CreateResourceName($"{orleansService.Name}-{name}-sqs");
+ var queueResourceNames = Enumerable.Range(0, validatedOptions.PartitionCount)
+ .Select(partition => $"{resourceName}-{partition}")
+ .ToArray();
+ ValidateResourceNames(orleansService.Builder, resourceName, queueResourceNames);
+ var stack = orleansService.Builder.AddAWSCDKStack(resourceName).WithReference(awsSdkConfig);
+ var queues = CreateQueues(stack, validatedOptions, queueNames, queueResourceNames);
+ stack.WithAnnotation(new SqsTopologyAnnotation(name, region, queueNames));
+ var resource = new SqsStreamingResource(
+ orleansService,
+ name,
+ validatedOptions,
+ awsSdkConfig,
+ stack,
+ queues);
+
+ orleansService
+ .WithServiceId(validatedOptions.ServiceId)
+ .WithStreaming(name, resource);
+ return resource;
+ }
+
+ private static void ValidateIdentities(
+ OrleansService orleansService,
+ string name,
+ string region,
+ IReadOnlyList queueNames)
+ {
+ foreach (var existingName in orleansService.Streaming.Keys)
+ {
+ if (string.Equals(existingName, name, StringComparison.OrdinalIgnoreCase))
+ {
+ throw new InvalidOperationException(
+ $"SQS stream provider '{name}' conflicts with existing stream provider '{existingName}'. "
+ + "Stream provider names must be unique ignoring case within an Orleans service.");
+ }
+ }
+
+ var physicalNames = new HashSet(queueNames, StringComparer.Ordinal);
+ foreach (var resource in orleansService.Builder.Resources)
+ {
+ foreach (var topology in resource.Annotations.OfType())
+ {
+ if (!string.Equals(topology.Region, region, StringComparison.Ordinal))
+ {
+ continue;
+ }
+
+ foreach (var queueName in topology.QueueNames)
+ {
+ if (physicalNames.Contains(queueName))
+ {
+ throw new InvalidOperationException(
+ $"SQS stream provider '{name}' generates physical queue '{queueName}' in region '{region}', "
+ + $"which is already owned by stream provider '{topology.ProviderName}' in stack '{resource.Name}'. "
+ + "Use distinct service identifiers or provider names for separate queue topologies.");
+ }
+ }
+ }
+ }
+ }
+
+ private sealed record SqsTopologyAnnotation(
+ string ProviderName,
+ string Region,
+ IReadOnlyList QueueNames) : IResourceAnnotation;
+
+ private static void ValidateResourceNames(
+ IDistributedApplicationBuilder builder,
+ string stackName,
+ IReadOnlyList queueResourceNames)
+ {
+ var names = new HashSet(queueResourceNames, StringComparer.OrdinalIgnoreCase) { stackName };
+ foreach (var resource in builder.Resources)
+ {
+ if (names.Contains(resource.Name))
+ {
+ throw new InvalidOperationException(
+ $"SQS streaming resource name '{resource.Name}' is already registered in the AppHost. "
+ + "Use distinct Orleans service names and provider names for separate stacks.");
+ }
+ }
+ }
+
+ private static SqsStreamingOptions ValidateAndCopy(
+ string providerName,
+ IAWSSDKConfig awsSdkConfig,
+ SqsStreamingOptions options)
+ {
+ ArgumentException.ThrowIfNullOrWhiteSpace(options.ServiceId);
+ if (awsSdkConfig.Region is null)
+ {
+ throw new ArgumentException("SQS streaming requires a concrete AWS region.", nameof(awsSdkConfig));
+ }
+
+ if (options.PartitionCount <= 0)
+ {
+ throw new ArgumentOutOfRangeException(nameof(options), "PartitionCount must be greater than zero.");
+ }
+
+ ValidateRange(options.ReceiveWaitTimeSeconds, 0, 20, nameof(options.ReceiveWaitTimeSeconds));
+ ValidateRange(options.VisibilityTimeoutSeconds, 0, 43_200, nameof(options.VisibilityTimeoutSeconds));
+ if (options.CacheSize is <= 0)
+ {
+ throw new ArgumentOutOfRangeException(nameof(options), "CacheSize must be greater than zero.");
+ }
+
+ if (string.Equals(options.DataAdapterKey, providerName, StringComparison.Ordinal))
+ {
+ throw new ArgumentException(
+ "DataAdapterKey must differ from the stream provider name.",
+ nameof(options));
+ }
+
+ var result = new SqsStreamingOptions
+ {
+ ServiceId = options.ServiceId,
+ PartitionCount = options.PartitionCount,
+ FifoQueue = options.FifoQueue,
+ ReceiveWaitTimeSeconds = options.ReceiveWaitTimeSeconds,
+ VisibilityTimeoutSeconds = options.VisibilityTimeoutSeconds,
+ CacheSize = options.CacheSize,
+ DataAdapterKey = options.DataAdapterKey,
+ ReceiveMessageAttributes = CopyValues(options.ReceiveMessageAttributes, nameof(options.ReceiveMessageAttributes)),
+ ReceiveMessageSystemAttributes = CopyValues(options.ReceiveMessageSystemAttributes, nameof(options.ReceiveMessageSystemAttributes)),
+ };
+
+ foreach (var queueName in GetPhysicalQueueNames(providerName, result))
+ {
+ if (queueName.Length > 80 || !SqsQueueNamePattern().IsMatch(queueName))
+ {
+ throw new ArgumentException(
+ $"The generated SQS queue name '{queueName}' must be at most 80 characters and contain only letters, numbers, hyphens, underscores, and an optional .fifo suffix.",
+ nameof(options));
+ }
+ }
+
+ return result;
+ }
+
+ private static IReadOnlyList>> CreateQueues(
+ IResourceBuilder stack,
+ SqsStreamingOptions options,
+ IReadOnlyList queueNames,
+ IReadOnlyList queueResourceNames)
+ {
+ var result = new List>>(queueNames.Count);
+ for (var index = 0; index < queueNames.Count; index++)
+ {
+ var queueName = queueNames[index];
+ result.Add(
+ stack.AddSQSQueue(
+ queueResourceNames[index],
+ new QueueProps
+ {
+ QueueName = queueName,
+ Fifo = options.FifoQueue,
+ ContentBasedDeduplication = options.FifoQueue,
+ DeduplicationScope = options.FifoQueue ? DeduplicationScope.MESSAGE_GROUP : null,
+ FifoThroughputLimit = options.FifoQueue ? FifoThroughputLimit.PER_MESSAGE_GROUP_ID : null,
+ ReceiveMessageWaitTime = options.ReceiveWaitTimeSeconds is { } receiveWaitTime
+ ? Duration.Seconds(receiveWaitTime)
+ : null,
+ VisibilityTimeout = options.VisibilityTimeoutSeconds is { } visibilityTimeout
+ ? Duration.Seconds(visibilityTimeout)
+ : null,
+ }));
+ }
+
+ return result.AsReadOnly();
+ }
+
+ private static IReadOnlyList GetPhysicalQueueNames(
+ string providerName,
+ SqsStreamingOptions options)
+ => Enumerable.Range(0, options.PartitionCount)
+ .Select(partition => SqsQueueName.Create(providerName, partition, options.FifoQueue, options.ServiceId))
+ .ToArray();
+
+ private static IReadOnlyList CopyValues(IReadOnlyList? values, string name)
+ {
+ if (values is null)
+ {
+ throw new ArgumentNullException(name);
+ }
+
+ var result = new string[values.Count];
+ for (var index = 0; index < values.Count; index++)
+ {
+ if (string.IsNullOrWhiteSpace(values[index]))
+ {
+ throw new ArgumentException($"{name} cannot contain empty values.", name);
+ }
+
+ result[index] = values[index];
+ }
+
+ return Array.AsReadOnly(result);
+ }
+
+ private static void ValidateRange(int? value, int minimum, int maximum, string name)
+ {
+ if (value is not null && (value < minimum || value > maximum))
+ {
+ throw new ArgumentOutOfRangeException(name, $"{name} must be between {minimum} and {maximum}.");
+ }
+ }
+
+ private static string NormalizeResourceName(string value)
+ {
+ var result = new StringBuilder(value.Length);
+ var previousWasHyphen = false;
+ foreach (var character in value)
+ {
+ if (char.IsLetterOrDigit(character))
+ {
+ result.Append(char.ToLowerInvariant(character));
+ previousWasHyphen = false;
+ }
+ else if (!previousWasHyphen)
+ {
+ result.Append('-');
+ previousWasHyphen = true;
+ }
+ }
+
+ var normalized = result.ToString().Trim('-');
+ return normalized.Length == 0 ? "sqs" : normalized;
+ }
+
+ private static string CreateResourceName(string value)
+ {
+ var normalized = NormalizeResourceName(value);
+ if (string.Equals(normalized, value.ToLowerInvariant(), StringComparison.Ordinal))
+ {
+ return normalized;
+ }
+
+ var hash = SHA256.HashData(Encoding.UTF8.GetBytes(value));
+ return $"{normalized}-{Convert.ToHexString(hash.AsSpan(0, 4)).ToLowerInvariant()}";
+ }
+
+ internal static void ValidateServiceId(
+ OrleansService orleansService,
+ string expectedServiceId,
+ bool allowUnset)
+ {
+ var configuredServiceId = GetServiceId(orleansService);
+ if (allowUnset
+ && configuredServiceId is ParameterResource generatedServiceId
+ && !orleansService.Builder.Resources.Contains(generatedServiceId))
+ {
+ return;
+ }
+
+ if (configuredServiceId is not string serviceId)
+ {
+ throw new InvalidOperationException(
+ "SQS streaming requires Orleans ServiceId to be a concrete string.");
+ }
+
+ if (!string.Equals(serviceId, expectedServiceId, StringComparison.Ordinal))
+ {
+ throw new InvalidOperationException(
+ $"SQS streaming ServiceId '{expectedServiceId}' conflicts with Orleans ServiceId '{serviceId}'.");
+ }
+ }
+
+ [UnsafeAccessor(UnsafeAccessorKind.Method, Name = "get_ServiceId")]
+ private static extern object? GetServiceId(OrleansService orleansService);
+
+ [GeneratedRegex(@"^[A-Za-z0-9_-]+(?:\.fifo)?$", RegexOptions.CultureInvariant)]
+ private static partial Regex SqsQueueNamePattern();
+}
diff --git a/src/AWS/Orleans.Streaming.SQS.Aspire/README.md b/src/AWS/Orleans.Streaming.SQS.Aspire/README.md
new file mode 100644
index 00000000000..fc0a6da0666
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS.Aspire/README.md
@@ -0,0 +1,53 @@
+# Microsoft Orleans Aspire integration for Amazon SQS
+
+The `Microsoft.Orleans.Streaming.SQS.Aspire` package configures an Orleans SQS stream provider and provisions its complete partition queue topology through the AWS CDK integration for .NET Aspire.
+
+## Install
+
+```shell
+dotnet add package Microsoft.Orleans.Streaming.SQS.Aspire
+```
+
+Install `Microsoft.Orleans.Streaming.SQS` in every silo and Orleans client project that uses the provider.
+
+## Configure
+
+```csharp
+using Amazon;
+using Aspire.Hosting;
+
+var builder = DistributedApplication.CreateBuilder(args);
+var aws = builder.AddAWSSDKConfig()
+ .WithRegion(RegionEndpoint.USEast1);
+
+var orleans = builder.AddOrleans("cluster")
+ .WithDevelopmentClustering()
+ .WithMemoryGrainStorage("PubSubStore")
+ .WithSqsStreaming(
+ "Orders",
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = "orders-service",
+ PartitionCount = 16,
+ FifoQueue = true,
+ ReceiveWaitTimeSeconds = 20,
+ VisibilityTimeoutSeconds = 60,
+ });
+
+builder.AddProject("silo")
+ .WithReference(orleans);
+
+builder.AddProject("client")
+ .WithReference(orleans.AsClient());
+```
+
+The options define both the AWS CDK queue resources and the Orleans configuration emitted to every referenced silo and client. The integration applies the stable service ID, attaches the AWS SDK profile and region, and makes each referenced resource wait for the CDK stack automatically.
+
+Provider names are unique ignoring case within an Orleans service, including providers registered through other streaming extensions. Each physical queue has one owning SQS stack per AppHost and region. Registration validates both identities before adding resources or changing Orleans configuration, and reports the conflicting provider or queue. Use distinct service IDs or provider names for separate topologies in the same region, including when AWS profiles select different accounts.
+
+Physical queue names preserve the service ID and lowercase the provider name: `orders-service-orders_primary-0`. Hyphens and underscores remain distinct in physical names; the CDK resource identifiers distinguish names such as `orders-primary` and `orders_primary`.
+
+## Documentation
+
+See [Stream with Amazon SQS](https://dotnet.github.io/orleans/docs/streaming/sqs-streaming/) for runtime behavior, permissions, delivery semantics, and operations.
diff --git a/src/AWS/Orleans.Streaming.SQS.Aspire/SqsStreamingOptions.cs b/src/AWS/Orleans.Streaming.SQS.Aspire/SqsStreamingOptions.cs
new file mode 100644
index 00000000000..3d98234851a
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS.Aspire/SqsStreamingOptions.cs
@@ -0,0 +1,52 @@
+namespace Aspire.Hosting;
+
+///
+/// Configures an Orleans SQS stream provider and its AWS CDK queue resources.
+///
+public sealed class SqsStreamingOptions
+{
+ ///
+ /// Gets or sets the stable Orleans service identifier used to prefix every physical queue name.
+ ///
+ public required string ServiceId { get; init; }
+
+ ///
+ /// Gets or sets the number of SQS queues in the stream provider topology.
+ ///
+ public int PartitionCount { get; init; } = 8;
+
+ ///
+ /// Gets or sets a value indicating whether the provider uses FIFO queues.
+ ///
+ public bool FifoQueue { get; init; }
+
+ ///
+ /// Gets or sets the SQS long-poll duration, in seconds.
+ ///
+ public int? ReceiveWaitTimeSeconds { get; init; }
+
+ ///
+ /// Gets or sets the SQS visibility timeout, in seconds.
+ ///
+ public int? VisibilityTimeoutSeconds { get; init; }
+
+ ///
+ /// Gets or sets the silo queue-cache size.
+ ///
+ public int? CacheSize { get; init; }
+
+ ///
+ /// Gets or sets the keyed dependency injection key for the SQS data adapter.
+ ///
+ public string? DataAdapterKey { get; init; }
+
+ ///
+ /// Gets or sets application-defined message attributes requested when receiving messages.
+ ///
+ public IReadOnlyList ReceiveMessageAttributes { get; init; } = [];
+
+ ///
+ /// Gets or sets SQS system attributes requested when receiving messages.
+ ///
+ public IReadOnlyList ReceiveMessageSystemAttributes { get; init; } = [];
+}
diff --git a/src/AWS/Orleans.Streaming.SQS.Aspire/SqsStreamingResource.cs b/src/AWS/Orleans.Streaming.SQS.Aspire/SqsStreamingResource.cs
new file mode 100644
index 00000000000..366a8ae4079
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS.Aspire/SqsStreamingResource.cs
@@ -0,0 +1,131 @@
+using System.Globalization;
+using Amazon.CDK.AWS.SQS;
+using Aspire.Hosting.ApplicationModel;
+using Aspire.Hosting.AWS;
+using Aspire.Hosting.AWS.CDK;
+using Aspire.Hosting.Orleans;
+
+namespace Aspire.Hosting;
+
+///
+/// Represents an Orleans SQS stream provider and its AWS CDK resources.
+///
+public sealed class SqsStreamingResource : IProviderConfiguration
+{
+ private readonly OrleansService _orleansService;
+ private readonly SqsStreamingOptions _options;
+
+ internal SqsStreamingResource(
+ OrleansService orleansService,
+ string name,
+ SqsStreamingOptions options,
+ IAWSSDKConfig awsSdkConfig,
+ IResourceBuilder stack,
+ IReadOnlyList>> queues)
+ {
+ _orleansService = orleansService;
+ Name = name;
+ _options = options;
+ AwsSdkConfig = awsSdkConfig;
+ Stack = stack;
+ Queues = queues;
+ }
+
+ ///
+ /// Gets the Orleans stream provider name.
+ ///
+ public string Name { get; }
+
+ ///
+ /// Gets the validated configuration used for Orleans and AWS CDK.
+ ///
+ public SqsStreamingOptions Options => _options;
+
+ ///
+ /// Gets the AWS SDK configuration associated with the provider.
+ ///
+ public IAWSSDKConfig AwsSdkConfig { get; }
+
+ ///
+ /// Gets the AWS CDK stack which owns the partition queues.
+ ///
+ public IResourceBuilder Stack { get; }
+
+ ///
+ /// Gets the AWS CDK queue resources in the provider topology.
+ ///
+ public IReadOnlyList>> Queues { get; }
+
+ ///
+ public void ConfigureResource(IResourceBuilder resourceBuilder, string configSectionPath)
+ where T : IResourceWithEnvironment
+ {
+ var prefix = $"Orleans__{configSectionPath.Replace(":", "__", StringComparison.Ordinal)}";
+ resourceBuilder
+ .WithReference(AwsSdkConfig)
+ .WithEnvironment($"{prefix}__ProviderType", "SQS")
+ .WithEnvironment($"{prefix}__Region", AwsSdkConfig.Region!.SystemName)
+ .WithEnvironment($"{prefix}__PartitionCount", _options.PartitionCount.ToString(CultureInfo.InvariantCulture))
+ .WithEnvironment($"{prefix}__FifoQueue", _options.FifoQueue.ToString());
+ resourceBuilder.WithEnvironment(context =>
+ {
+ OrleansSqsStreamingExtensions.ValidateServiceId(
+ _orleansService,
+ _options.ServiceId,
+ allowUnset: false);
+ });
+ if (resourceBuilder.Resource is IResourceWithWaitSupport)
+ {
+ resourceBuilder.WithAnnotation(
+ new WaitAnnotation(Stack.Resource, WaitType.WaitUntilHealthy, exitCode: 0));
+ }
+
+ AddOptionalValue(resourceBuilder, prefix, nameof(_options.ReceiveWaitTimeSeconds), _options.ReceiveWaitTimeSeconds);
+ AddOptionalValue(resourceBuilder, prefix, nameof(_options.VisibilityTimeoutSeconds), _options.VisibilityTimeoutSeconds);
+ AddOptionalValue(resourceBuilder, prefix, nameof(_options.CacheSize), _options.CacheSize);
+ AddOptionalValue(resourceBuilder, prefix, nameof(_options.DataAdapterKey), _options.DataAdapterKey);
+ AddValues(resourceBuilder, prefix, nameof(_options.ReceiveMessageAttributes), _options.ReceiveMessageAttributes);
+ AddValues(resourceBuilder, prefix, nameof(_options.ReceiveMessageSystemAttributes), _options.ReceiveMessageSystemAttributes);
+ }
+
+ private static void AddOptionalValue(
+ IResourceBuilder resourceBuilder,
+ string prefix,
+ string name,
+ int? value)
+ where T : IResourceWithEnvironment
+ {
+ if (value is { } configuredValue)
+ {
+ resourceBuilder.WithEnvironment(
+ $"{prefix}__{name}",
+ configuredValue.ToString(CultureInfo.InvariantCulture));
+ }
+ }
+
+ private static void AddOptionalValue(
+ IResourceBuilder resourceBuilder,
+ string prefix,
+ string name,
+ string? value)
+ where T : IResourceWithEnvironment
+ {
+ if (!string.IsNullOrWhiteSpace(value))
+ {
+ resourceBuilder.WithEnvironment($"{prefix}__{name}", value);
+ }
+ }
+
+ private static void AddValues(
+ IResourceBuilder resourceBuilder,
+ string prefix,
+ string name,
+ IReadOnlyList values)
+ where T : IResourceWithEnvironment
+ {
+ for (var index = 0; index < values.Count; index++)
+ {
+ resourceBuilder.WithEnvironment($"{prefix}__{name}__{index}", values[index]);
+ }
+ }
+}
diff --git a/src/AWS/Orleans.Streaming.SQS/Hosting/SqsStreamProviderBuilder.cs b/src/AWS/Orleans.Streaming.SQS/Hosting/SqsStreamProviderBuilder.cs
new file mode 100644
index 00000000000..42f8da18aea
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS/Hosting/SqsStreamProviderBuilder.cs
@@ -0,0 +1,366 @@
+using System;
+using System.Collections.Generic;
+using System.Globalization;
+using System.Linq;
+using Amazon;
+using Microsoft.Extensions.Configuration;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Options;
+using Orleans;
+using Orleans.Configuration;
+using Orleans.Hosting;
+using Orleans.Providers;
+using Orleans.Streaming.SQS;
+using Orleans.Streaming.SQS.Streams;
+
+[assembly: RegisterProvider("SQS", "Streaming", "Silo", typeof(SqsStreamProviderBuilder))]
+[assembly: RegisterProvider("SQS", "Streaming", "Client", typeof(SqsStreamProviderBuilder))]
+[assembly: RegisterProvider("AmazonSQS", "Streaming", "Silo", typeof(SqsStreamProviderBuilder))]
+[assembly: RegisterProvider("AmazonSQS", "Streaming", "Client", typeof(SqsStreamProviderBuilder))]
+
+namespace Orleans.Hosting;
+
+///
+/// Configures Amazon SQS stream providers from Orleans configuration.
+///
+public sealed class SqsStreamProviderBuilder : IProviderBuilder, IProviderBuilder
+{
+ ///
+ public void Configure(ISiloBuilder builder, string? name, IConfigurationSection configurationSection)
+ {
+ ArgumentNullException.ThrowIfNull(name);
+
+ builder.AddSqsStreams(name, streams =>
+ {
+ ConfigureCommon(streams, name, configurationSection);
+
+ var cacheSize = GetPositiveInt(configurationSection, "CacheSize");
+ if (cacheSize.HasValue)
+ {
+ streams.ConfigureCache(cacheSize.Value);
+ }
+ });
+ }
+
+ ///
+ public void Configure(IClientBuilder builder, string? name, IConfigurationSection configurationSection)
+ {
+ ArgumentNullException.ThrowIfNull(name);
+
+ builder.AddSqsStreams(name, streams => ConfigureCommon(streams, name, configurationSection));
+ }
+
+ private static void ConfigureCommon(
+ SiloSqsStreamConfigurator streams,
+ string providerName,
+ IConfigurationSection configurationSection)
+ {
+ streams.ConfigureSqs(GetSqsOptionsBuilder(configurationSection));
+ ConfigurePartitioning(streams.ConfigurePartitioning, configurationSection);
+ ConfigureDataAdapter(streams.UseDataAdapter, providerName, configurationSection);
+ }
+
+ private static void ConfigureCommon(
+ ClusterClientSqsStreamConfigurator streams,
+ string providerName,
+ IConfigurationSection configurationSection)
+ {
+ streams.ConfigureSqs(GetSqsOptionsBuilder(configurationSection));
+ ConfigurePartitioning(streams.ConfigurePartitioning, configurationSection);
+ ConfigureDataAdapter(streams.UseDataAdapter, providerName, configurationSection);
+ }
+
+ private static Action> GetSqsOptionsBuilder(IConfigurationSection configurationSection)
+ => optionsBuilder => optionsBuilder.Configure((options, services) =>
+ {
+ options.ConnectionString = ResolveConnectionString(configurationSection, services);
+ options.FifoQueue = GetBoolean(configurationSection, nameof(options.FifoQueue)) ?? options.FifoQueue;
+ options.ReceiveWaitTimeSeconds = GetIntInRange(
+ configurationSection,
+ nameof(options.ReceiveWaitTimeSeconds),
+ minimum: 0,
+ maximum: 20)
+ ?? options.ReceiveWaitTimeSeconds;
+ options.VisibilityTimeoutSeconds = GetIntInRange(
+ configurationSection,
+ nameof(options.VisibilityTimeoutSeconds),
+ minimum: 0,
+ maximum: 43_200)
+ ?? options.VisibilityTimeoutSeconds;
+
+ ConfigureList(
+ configurationSection.GetSection(nameof(options.ReceiveMessageAttributes)),
+ options.ReceiveMessageAttributes);
+ ConfigureList(
+ configurationSection.GetSection(nameof(options.ReceiveMessageSystemAttributes)),
+ options.ReceiveMessageSystemAttributes);
+ });
+
+ private static string ResolveConnectionString(
+ IConfigurationSection configurationSection,
+ IServiceProvider services)
+ {
+ var connectionString = configurationSection["ConnectionString"];
+ var serviceKey = configurationSection["ServiceKey"];
+ var connectionName = configurationSection["ConnectionName"];
+ if (!string.IsNullOrWhiteSpace(serviceKey)
+ && !string.IsNullOrWhiteSpace(connectionName)
+ && !string.Equals(serviceKey, connectionName, StringComparison.OrdinalIgnoreCase))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration specifies different ServiceKey and ConnectionName values. Configure one referenced SQS service.");
+ }
+
+ var configuration = services.GetRequiredService();
+ var referenceName = !string.IsNullOrWhiteSpace(serviceKey) ? serviceKey : connectionName;
+ if (!string.IsNullOrWhiteSpace(connectionString) && !string.IsNullOrWhiteSpace(referenceName))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration specifies both a connection reference and ConnectionString. Configure one SQS connection source.");
+ }
+
+ if (!string.IsNullOrWhiteSpace(referenceName))
+ {
+ connectionString = configuration.GetConnectionString(referenceName);
+ if (string.IsNullOrWhiteSpace(connectionString))
+ {
+ throw new OrleansConfigurationException(
+ $"SQS streaming connection reference '{referenceName}' did not resolve to a connection string.");
+ }
+ }
+
+ var service = configurationSection["Service"];
+ var region = configurationSection["Region"];
+ var configuredServiceEndpoint = configurationSection["ServiceEndpoint"];
+ var configuredEndpoint = configurationSection["Endpoint"];
+ if (!string.IsNullOrWhiteSpace(configuredServiceEndpoint)
+ && !string.IsNullOrWhiteSpace(configuredEndpoint)
+ && !string.Equals(
+ configuredServiceEndpoint.Trim(),
+ configuredEndpoint.Trim(),
+ StringComparison.Ordinal))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration specifies different ServiceEndpoint and Endpoint values. Configure one SQS service endpoint.");
+ }
+
+ var serviceEndpoint = GetFirstNonWhiteSpace(
+ configuredServiceEndpoint,
+ configuredEndpoint);
+ var configuredLocations = new[] { service, region, serviceEndpoint }
+ .Count(value => !string.IsNullOrWhiteSpace(value));
+ if (configuredLocations > 1)
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration specifies multiple values among Service, Region, and ServiceEndpoint. Configure one SQS service location.");
+ }
+
+ if (!string.IsNullOrWhiteSpace(serviceEndpoint))
+ {
+ ValidateServiceEndpoint(serviceEndpoint);
+ }
+
+ var configuredService = GetFirstNonWhiteSpace(service, region, serviceEndpoint);
+ if (!string.IsNullOrWhiteSpace(connectionString))
+ {
+ if (!string.IsNullOrWhiteSpace(configuredService))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration specifies both a connection string and a service location. Configure one SQS connection source.");
+ }
+
+ ValidateConnectionString(connectionString);
+ return connectionString;
+ }
+
+ if (string.IsNullOrWhiteSpace(configuredService))
+ {
+ var awsRegion = configuration["AWS_REGION"];
+ var awsDefaultRegion = configuration["AWS_DEFAULT_REGION"];
+ if (!string.IsNullOrWhiteSpace(awsRegion)
+ && !string.IsNullOrWhiteSpace(awsDefaultRegion)
+ && !string.Equals(awsRegion, awsDefaultRegion, StringComparison.OrdinalIgnoreCase))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration found different AWS_REGION and AWS_DEFAULT_REGION values. Configure one AWS region.");
+ }
+
+ configuredService = !string.IsNullOrWhiteSpace(awsRegion) ? awsRegion : awsDefaultRegion;
+ }
+
+ if (string.IsNullOrWhiteSpace(configuredService))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming requires a service location. Configure ServiceKey, ConnectionName, ConnectionString, Region, ServiceEndpoint, AWS_REGION, or AWS_DEFAULT_REGION.");
+ }
+
+ ValidateServiceLocation(configuredService);
+ return $"Service={configuredService}";
+ }
+
+ private static void ValidateConnectionString(string connectionString)
+ {
+ var properties = SqsConnectionString.Parse(connectionString);
+ SqsConnectionString.ValidateCredentials(properties);
+ if (!properties.TryGetValue("Service", out var service))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming connection strings require a non-empty Service value containing an AWS region or SQS-compatible endpoint.");
+ }
+
+ ValidateServiceLocation(service);
+ }
+
+ private static void ValidateServiceLocation(string service)
+ {
+ if (service.Contains("://", StringComparison.Ordinal))
+ {
+ ValidateServiceEndpoint(service);
+ }
+ else if (service.Contains(':'))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming service endpoint values must be absolute HTTP or HTTPS URIs.");
+ }
+ else if (!RegionEndpoint.EnumerableAllRegions.Any(
+ region => string.Equals(region.SystemName, service, StringComparison.OrdinalIgnoreCase)))
+ {
+ throw new OrleansConfigurationException(
+ $"SQS streaming region '{service}' is not recognized by the AWS SDK.");
+ }
+ }
+
+ private static void ValidateServiceEndpoint(string serviceEndpoint)
+ {
+ if (!Uri.TryCreate(serviceEndpoint, UriKind.Absolute, out var endpoint)
+ || endpoint.Scheme is not ("http" or "https"))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming ServiceEndpoint values must be absolute HTTP or HTTPS URIs.");
+ }
+ }
+
+ private static void ConfigurePartitioning(
+ Func configurePartitioning,
+ IConfigurationSection configurationSection)
+ {
+ var partitionCount = GetPositiveInt(configurationSection, "PartitionCount");
+ if (partitionCount.HasValue)
+ {
+ configurePartitioning(partitionCount.Value);
+ }
+ }
+
+ private static void ConfigureDataAdapter(
+ Func, object> useDataAdapter,
+ string providerName,
+ IConfigurationSection configurationSection)
+ {
+ var dataAdapterKey = configurationSection["DataAdapterKey"];
+ var dataAdapterServiceKey = configurationSection["DataAdapterServiceKey"];
+ if (!string.IsNullOrWhiteSpace(dataAdapterKey)
+ && !string.IsNullOrWhiteSpace(dataAdapterServiceKey)
+ && !string.Equals(dataAdapterKey, dataAdapterServiceKey, StringComparison.Ordinal))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming configuration specifies different DataAdapterKey and DataAdapterServiceKey values. Configure one keyed data adapter.");
+ }
+
+ var key = !string.IsNullOrWhiteSpace(dataAdapterKey) ? dataAdapterKey : dataAdapterServiceKey;
+ if (!string.IsNullOrWhiteSpace(key))
+ {
+ if (string.Equals(key, providerName, StringComparison.Ordinal))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming DataAdapterKey must differ from the stream provider name.");
+ }
+
+ useDataAdapter((services, _) => services.GetRequiredKeyedService(key));
+ }
+ }
+
+ private static void ConfigureList(IConfigurationSection section, List values)
+ {
+ var configuredValues = section.GetChildren()
+ .OrderBy(child => GetConfigurationIndex(child.Key))
+ .Select(child => child.Value)
+ .Where(value => !string.IsNullOrWhiteSpace(value))
+ .Select(value => value!)
+ .ToList();
+ if (configuredValues.Count > 0)
+ {
+ values.Clear();
+ values.AddRange(configuredValues);
+ }
+ }
+
+ private static int GetConfigurationIndex(string key)
+ => int.TryParse(key, NumberStyles.None, CultureInfo.InvariantCulture, out var index)
+ ? index
+ : int.MaxValue;
+
+ private static string? GetFirstNonWhiteSpace(params string?[] values)
+ => values.FirstOrDefault(value => !string.IsNullOrWhiteSpace(value));
+
+ private static bool? GetBoolean(IConfigurationSection section, string key)
+ {
+ var value = section[key];
+ if (string.IsNullOrWhiteSpace(value))
+ {
+ return null;
+ }
+
+ if (bool.TryParse(value, out var result))
+ {
+ return result;
+ }
+
+ throw new OrleansConfigurationException(
+ $"SQS streaming configuration value '{key}' must be true or false.");
+ }
+
+ private static int? GetPositiveInt(IConfigurationSection section, string key)
+ {
+ var result = GetInt(section, key);
+ if (result is <= 0)
+ {
+ throw new OrleansConfigurationException(
+ $"SQS streaming configuration value '{key}' must be greater than zero.");
+ }
+
+ return result;
+ }
+
+ private static int? GetIntInRange(
+ IConfigurationSection section,
+ string key,
+ int minimum,
+ int maximum)
+ {
+ var result = GetInt(section, key);
+ if (result is not null && (result < minimum || result > maximum))
+ {
+ throw new OrleansConfigurationException(
+ $"SQS streaming configuration value '{key}' must be between {minimum} and {maximum}.");
+ }
+
+ return result;
+ }
+
+ private static int? GetInt(IConfigurationSection section, string key)
+ {
+ var value = section[key];
+ if (string.IsNullOrWhiteSpace(value))
+ {
+ return null;
+ }
+
+ if (int.TryParse(value, out var result))
+ {
+ return result;
+ }
+
+ throw new OrleansConfigurationException(
+ $"SQS streaming configuration value '{key}' must be an integer.");
+ }
+}
diff --git a/src/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.csproj b/src/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.csproj
index 663fb6dbe80..824326558ab 100644
--- a/src/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.csproj
+++ b/src/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.csproj
@@ -12,6 +12,7 @@
+
@@ -28,4 +29,3 @@
-
diff --git a/src/AWS/Orleans.Streaming.SQS/README.md b/src/AWS/Orleans.Streaming.SQS/README.md
index b5bd68f9280..8f9171fd665 100644
--- a/src/AWS/Orleans.Streaming.SQS/README.md
+++ b/src/AWS/Orleans.Streaming.SQS/README.md
@@ -28,6 +28,8 @@ clientBuilder.AddSqsStreams("Orders", options =>
When explicit credentials aren't present in the connection string, the provider uses the AWS SDK credential resolution chain. Prefer workload credentials such as an IAM role.
+For .NET Aspire AppHost projects, install `Microsoft.Orleans.Streaming.SQS.Aspire` to provision the partition queues through AWS CDK and emit the matching silo/client configuration from one options object.
+
## Documentation
See [Stream with Amazon SQS](https://dotnet.github.io/orleans/docs/streaming/sqs-streaming/) for FIFO queues, custom data adapters, delivery semantics, permissions, tuning, and operational guidance.
diff --git a/src/AWS/Orleans.Streaming.SQS/SqsConnectionString.cs b/src/AWS/Orleans.Streaming.SQS/SqsConnectionString.cs
new file mode 100644
index 00000000000..e79c1d6f083
--- /dev/null
+++ b/src/AWS/Orleans.Streaming.SQS/SqsConnectionString.cs
@@ -0,0 +1,64 @@
+using System;
+using System.Collections.Generic;
+using Orleans.Configuration;
+
+namespace Orleans.Streaming.SQS;
+
+internal static class SqsConnectionString
+{
+ public static IReadOnlyDictionary Parse(string connectionString)
+ {
+ if (string.IsNullOrWhiteSpace(connectionString))
+ {
+ throw new OrleansConfigurationException("SQS streaming requires a non-empty connection string.");
+ }
+
+ var result = new Dictionary(StringComparer.OrdinalIgnoreCase);
+ foreach (var segment in connectionString.Split(
+ ';',
+ StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries))
+ {
+ var separator = segment.IndexOf('=');
+ if (separator <= 0 || string.IsNullOrWhiteSpace(segment[(separator + 1)..]))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming connection strings use non-empty key=value segments separated by semicolons.");
+ }
+
+ var key = segment[..separator].Trim();
+ var value = segment[(separator + 1)..].Trim();
+ if (key.Length == 0)
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming connection string property names must be non-empty.");
+ }
+
+ if (!result.TryAdd(key, value))
+ {
+ throw new OrleansConfigurationException(
+ $"SQS streaming connection string property '{key}' is configured more than once.");
+ }
+ }
+
+ return result;
+ }
+
+ public static void ValidateCredentials(IReadOnlyDictionary properties)
+ {
+ var hasAccessKey = properties.ContainsKey("AccessKey");
+ var hasSecretKey = properties.ContainsKey("SecretKey");
+ var hasSessionToken = properties.ContainsKey("SessionToken");
+
+ if (hasSessionToken && (!hasAccessKey || !hasSecretKey))
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming connection string property 'SessionToken' requires both AccessKey and SecretKey.");
+ }
+
+ if (hasAccessKey != hasSecretKey)
+ {
+ throw new OrleansConfigurationException(
+ "SQS streaming connection strings must configure AccessKey and SecretKey together.");
+ }
+ }
+}
diff --git a/src/AWS/Orleans.Streaming.SQS/Storage/SQSStorage.cs b/src/AWS/Orleans.Streaming.SQS/Storage/SQSStorage.cs
index 615f77c195b..cecc177b427 100644
--- a/src/AWS/Orleans.Streaming.SQS/Storage/SQSStorage.cs
+++ b/src/AWS/Orleans.Streaming.SQS/Storage/SQSStorage.cs
@@ -3,7 +3,6 @@
using System.Globalization;
using System.Linq;
using System.Net;
-using System.Text;
using System.Threading.Tasks;
using Amazon.Runtime;
using Amazon.SQS;
@@ -58,7 +57,7 @@ public SQSStorage(ILoggerFactory loggerFactory, string queueName, SqsOptions sqs
{
if (sqsOptions is null) throw new ArgumentNullException(nameof(sqsOptions));
this.sqsOptions = sqsOptions;
- QueueName = ConstructQueueName(queueName, sqsOptions, serviceId);
+ QueueName = SqsQueueName.Create(queueName, sqsOptions.FifoQueue, serviceId);
ParseDataConnectionString(sqsOptions.ConnectionString);
Logger = loggerFactory.CreateLogger();
CreateClient();
@@ -75,68 +74,46 @@ public SQSStorage(ILoggerFactory loggerFactory, string queueName, SqsOptions sqs
private void ParseDataConnectionString(string dataConnectionString)
{
- if (string.IsNullOrEmpty(dataConnectionString)) throw new ArgumentNullException(nameof(dataConnectionString));
-
- var parameters = dataConnectionString.Split(new[] { ';' }, StringSplitOptions.RemoveEmptyEntries);
-
- var serviceConfig = Array.Find(parameters, p => p.Contains(ServicePropertyName));
- if (!string.IsNullOrWhiteSpace(serviceConfig))
- {
- var value = serviceConfig.Split('=', StringSplitOptions.RemoveEmptyEntries);
- if (value.Length == 2 && !string.IsNullOrWhiteSpace(value[1]))
- service = value[1];
- }
-
- var secretKeyConfig = Array.Find(parameters, p => p.Contains(SecretKeyPropertyName));
- if (!string.IsNullOrWhiteSpace(secretKeyConfig))
- {
- var value = secretKeyConfig.Split('=', StringSplitOptions.RemoveEmptyEntries);
- if (value.Length == 2 && !string.IsNullOrWhiteSpace(value[1]))
- secretKey = value[1];
- }
-
- var accessKeyConfig = Array.Find(parameters, p => p.Contains(AccessKeyPropertyName));
- if (!string.IsNullOrWhiteSpace(accessKeyConfig))
+ var parameters = SqsConnectionString.Parse(dataConnectionString);
+ SqsConnectionString.ValidateCredentials(parameters);
+ if (!parameters.TryGetValue(ServicePropertyName, out var serviceValue))
{
- var value = accessKeyConfig.Split('=', StringSplitOptions.RemoveEmptyEntries);
- if (value.Length == 2 && !string.IsNullOrWhiteSpace(value[1]))
- accessKey = value[1];
+ throw new OrleansConfigurationException(
+ "SQS streaming connection strings require a non-empty Service value containing an AWS region or SQS-compatible endpoint.");
}
- var sessionTokenConfig = parameters.Where(p => p.Contains(SessionTokenPropertyName)).FirstOrDefault();
- if (!string.IsNullOrWhiteSpace(sessionTokenConfig))
- {
- var value = sessionTokenConfig.Split('=', 2, StringSplitOptions.RemoveEmptyEntries);
- if (value.Length == 2 && !string.IsNullOrWhiteSpace(value[1]))
- sessionToken = value[1];
- }
+ service = serviceValue;
+ parameters.TryGetValue(SecretKeyPropertyName, out secretKey);
+ parameters.TryGetValue(AccessKeyPropertyName, out accessKey);
+ parameters.TryGetValue(SessionTokenPropertyName, out sessionToken);
}
private void CreateClient()
{
- if (service.StartsWith("http://", StringComparison.OrdinalIgnoreCase) ||
- service.StartsWith("https://", StringComparison.OrdinalIgnoreCase))
+ var isServiceEndpoint = service.StartsWith("http://", StringComparison.OrdinalIgnoreCase)
+ || service.StartsWith("https://", StringComparison.OrdinalIgnoreCase);
+ var config = isServiceEndpoint
+ ? new AmazonSQSConfig { ServiceURL = service }
+ : new AmazonSQSConfig { RegionEndpoint = AWSUtils.GetRegionEndpoint(service) };
+
+ if (!string.IsNullOrEmpty(accessKey) && !string.IsNullOrEmpty(secretKey) && !string.IsNullOrEmpty(sessionToken))
{
- // Local SQS instance (for testing)
- var credentials = new BasicAWSCredentials("dummy", "dummyKey");
- sqsClient = new AmazonSQSClient(credentials, new AmazonSQSConfig { ServiceURL = service });
+ sqsClient = new AmazonSQSClient(
+ new SessionAWSCredentials(accessKey, secretKey, sessionToken),
+ config);
}
- else if (!string.IsNullOrEmpty(accessKey) && !string.IsNullOrEmpty(secretKey) && !string.IsNullOrEmpty(sessionToken))
+ else if (!string.IsNullOrEmpty(accessKey) && !string.IsNullOrEmpty(secretKey))
{
- // AWS SQS instance (auth via explicit credentials)
- var credentials = new SessionAWSCredentials(accessKey, secretKey, sessionToken);
- sqsClient = new AmazonSQSClient(credentials, new AmazonSQSConfig { RegionEndpoint = AWSUtils.GetRegionEndpoint(service) });
+ sqsClient = new AmazonSQSClient(new BasicAWSCredentials(accessKey, secretKey), config);
}
- else if (!string.IsNullOrEmpty(accessKey) && !string.IsNullOrEmpty(secretKey))
+ else if (service.StartsWith("http://", StringComparison.OrdinalIgnoreCase))
{
- // AWS SQS instance (auth via explicit credentials)
- var credentials = new BasicAWSCredentials(accessKey, secretKey);
- sqsClient = new AmazonSQSClient(credentials, new AmazonSQSConfig { RegionEndpoint = AWSUtils.GetRegionEndpoint(service) });
+ // Local SQS-compatible services typically require a signed request but do not validate credentials.
+ sqsClient = new AmazonSQSClient(new BasicAWSCredentials("dummy", "dummyKey"), config);
}
else
{
- // AWS SQS instance (implicit auth - EC2 IAM Roles etc)
- sqsClient = new AmazonSQSClient(new AmazonSQSConfig { RegionEndpoint = AWSUtils.GetRegionEndpoint(service) });
+ sqsClient = new AmazonSQSClient(config);
}
}
@@ -440,24 +417,6 @@ private void ReportErrorAndRethrow(Exception exc, string operation)
throw new AggregateException($"Error doing {operation} for SQS queue {QueueName}", exc);
}
- private static string ConstructQueueName(string queueName, SqsOptions sqsOptions, string serviceId)
- {
- var queueNameBuilder = new StringBuilder();
- if (!string.IsNullOrWhiteSpace(serviceId))
- {
- queueNameBuilder.Append(serviceId);
- queueNameBuilder.Append('-');
- }
-
- queueNameBuilder.Append(queueName);
- if (sqsOptions.FifoQueue)
- {
- queueNameBuilder.Append(".fifo");
- }
-
- return queueNameBuilder.ToString();
- }
-
[LoggerMessage(
EventId = (int)ErrorCode.StreamProviderManagerBase,
Level = LogLevel.Error,
diff --git a/src/AWS/Shared/Storage/SqsQueueName.cs b/src/AWS/Shared/Storage/SqsQueueName.cs
new file mode 100644
index 00000000000..96d6dc7c452
--- /dev/null
+++ b/src/AWS/Shared/Storage/SqsQueueName.cs
@@ -0,0 +1,10 @@
+namespace OrleansAWSUtils.Storage;
+
+internal static class SqsQueueName
+{
+ public static string Create(string queueName, bool fifoQueue, string serviceId)
+ => $"{(string.IsNullOrWhiteSpace(serviceId) ? string.Empty : $"{serviceId}-")}{queueName}{(fifoQueue ? ".fifo" : string.Empty)}";
+
+ public static string Create(string providerName, int partition, bool fifoQueue, string serviceId)
+ => Create($"{providerName.ToLowerInvariant()}-{partition}", fifoQueue, serviceId);
+}
diff --git a/src/api/AWS/Orleans.Streaming.SQS.Aspire/Orleans.Streaming.SQS.Aspire.cs b/src/api/AWS/Orleans.Streaming.SQS.Aspire/Orleans.Streaming.SQS.Aspire.cs
new file mode 100644
index 00000000000..e2eba9e4446
--- /dev/null
+++ b/src/api/AWS/Orleans.Streaming.SQS.Aspire/Orleans.Streaming.SQS.Aspire.cs
@@ -0,0 +1,56 @@
+//------------------------------------------------------------------------------
+//
+// This code was generated by a tool.
+//
+// Changes to this file may cause incorrect behavior and will be lost if
+// the code is regenerated.
+//
+//------------------------------------------------------------------------------
+namespace Aspire.Hosting
+{
+ public static partial class OrleansSqsStreamingExtensions
+ {
+ public static SqsStreamingResource AddSqsStreaming(this Orleans.OrleansService orleansService, string name, AWS.IAWSSDKConfig awsSdkConfig, SqsStreamingOptions options) { throw null; }
+
+ public static Orleans.OrleansService WithSqsStreaming(this Orleans.OrleansService orleansService, string name, AWS.IAWSSDKConfig awsSdkConfig, SqsStreamingOptions options) { throw null; }
+ }
+
+ public sealed partial class SqsStreamingOptions
+ {
+ public int? CacheSize { get { throw null; } init { } }
+
+ public string? DataAdapterKey { get { throw null; } init { } }
+
+ public bool FifoQueue { get { throw null; } init { } }
+
+ public int PartitionCount { get { throw null; } init { } }
+
+ public System.Collections.Generic.IReadOnlyList ReceiveMessageAttributes { get { throw null; } init { } }
+
+ public System.Collections.Generic.IReadOnlyList ReceiveMessageSystemAttributes { get { throw null; } init { } }
+
+ public int? ReceiveWaitTimeSeconds { get { throw null; } init { } }
+
+ public required string ServiceId { get { throw null; } init { } }
+
+ public int? VisibilityTimeoutSeconds { get { throw null; } init { } }
+ }
+
+ public sealed partial class SqsStreamingResource : Orleans.IProviderConfiguration
+ {
+ internal SqsStreamingResource() { }
+
+ public AWS.IAWSSDKConfig AwsSdkConfig { get { throw null; } }
+
+ public string Name { get { throw null; } }
+
+ public SqsStreamingOptions Options { get { throw null; } }
+
+ public System.Collections.Generic.IReadOnlyList>> Queues { get { throw null; } }
+
+ public ApplicationModel.IResourceBuilder Stack { get { throw null; } }
+
+ public void ConfigureResource(ApplicationModel.IResourceBuilder resourceBuilder, string configSectionPath)
+ where T : ApplicationModel.IResourceWithEnvironment { }
+ }
+}
\ No newline at end of file
diff --git a/src/api/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.cs b/src/api/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.cs
index e8650b7b73e..2e22a42c0e4 100644
--- a/src/api/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.cs
+++ b/src/api/AWS/Orleans.Streaming.SQS/Orleans.Streaming.SQS.cs
@@ -52,6 +52,15 @@ public static partial class SiloBuilderExtensions
public static ISiloBuilder AddSqsStreams(this ISiloBuilder builder, string name, System.Action configure) { throw null; }
}
+ public sealed partial class SqsStreamProviderBuilder : Providers.IProviderBuilder, Providers.IProviderBuilder
+ {
+ public SqsStreamProviderBuilder() { }
+
+ public void Configure(IClientBuilder builder, string? name, Microsoft.Extensions.Configuration.IConfigurationSection configurationSection) { }
+
+ public void Configure(ISiloBuilder builder, string? name, Microsoft.Extensions.Configuration.IConfigurationSection configurationSection) { }
+ }
+
public partial class SiloSqsStreamConfigurator : SiloPersistentStreamConfigurator
{
public SiloSqsStreamConfigurator(string name, System.Action> configureServicesDelegate) : base(default!, default!, default!) { }
diff --git a/test/Extensions/Orleans.AWS.Tests/Orleans.AWS.Tests.csproj b/test/Extensions/Orleans.AWS.Tests/Orleans.AWS.Tests.csproj
index 714ddf05d5a..db15a872371 100644
--- a/test/Extensions/Orleans.AWS.Tests/Orleans.AWS.Tests.csproj
+++ b/test/Extensions/Orleans.AWS.Tests/Orleans.AWS.Tests.csproj
@@ -3,9 +3,38 @@
$(DefineConstants);AWSUTILS_TESTS
$(TestTargetFrameworks)
true
+ false
enable
+ false
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/test/Extensions/Orleans.AWS.Tests/StorageTests/AWSTestConstants.cs b/test/Extensions/Orleans.AWS.Tests/StorageTests/AWSTestConstants.cs
index a6f0d1dc792..e77ab65bc35 100644
--- a/test/Extensions/Orleans.AWS.Tests/StorageTests/AWSTestConstants.cs
+++ b/test/Extensions/Orleans.AWS.Tests/StorageTests/AWSTestConstants.cs
@@ -2,7 +2,11 @@
using Amazon.DynamoDBv2.Model;
using Amazon.Runtime;
using Microsoft.Extensions.Logging.Abstractions;
+#if TRANSACTIONS_DYNAMODB_TESTS
+using Orleans.Transactions.DynamoDB;
+#else
using Orleans.AWSUtils.Tests;
+#endif
using Orleans.Internal;
using TestExtensions;
diff --git a/test/Extensions/Orleans.AWS.Tests/Streaming/SQSAspireIntegrationTests.cs b/test/Extensions/Orleans.AWS.Tests/Streaming/SQSAspireIntegrationTests.cs
new file mode 100644
index 00000000000..5524f5b3a02
--- /dev/null
+++ b/test/Extensions/Orleans.AWS.Tests/Streaming/SQSAspireIntegrationTests.cs
@@ -0,0 +1,294 @@
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Options;
+using Orleans.Configuration;
+using Orleans.Hosting;
+using Orleans.Streaming.SQS.Streams;
+using Orleans.Streams;
+using OrleansAWSUtils.Storage;
+using OrleansAWSUtils.Streams;
+using SqsMessage = Amazon.SQS.Model.Message;
+using Xunit;
+
+namespace AWSUtils.Tests.Streaming;
+
+[Collection(SQSStreamProviderBuilderTestCollection.CollectionName)]
+[TestSuite("BVT")]
+[TestProvider("SQS")]
+[TestArea("Streaming")]
+[TestCategory("AWS"), TestCategory("SQS"), TestCategory("BVT")]
+public sealed class SQSAspireIntegrationTests
+{
+ private const string ServiceKey = "orleans-sqs";
+ private const string ServiceId = "aspire-sqs-service";
+ private const string RegionAdapterKey = "aspire-region-adapter";
+
+ [Fact]
+ public async Task AspireAppModel_Region_ProducesWorkingSiloConfiguration()
+ {
+ await using var app = await CreateRegionAppAsync();
+ var environment = await app.GetSiloEnvironmentAsync();
+ var configuration = SqsAspireTestApp.NormalizeConfiguration(environment);
+ var provider = GetProviderConfiguration(configuration, app.ProviderName);
+
+ AssertExactConfiguration(
+ new Dictionary
+ {
+ ["ProviderType"] = "SQS",
+ ["PartitionCount"] = "4",
+ ["FifoQueue"] = "False",
+ ["ReceiveWaitTimeSeconds"] = "12",
+ ["VisibilityTimeoutSeconds"] = "45",
+ ["ReceiveMessageAttributes:0"] = "TraceId",
+ ["ReceiveMessageAttributes:1"] = "Tenant",
+ ["ReceiveMessageSystemAttributes:0"] = "SentTimestamp",
+ ["ReceiveMessageSystemAttributes:1"] = "SequenceNumber",
+ ["DataAdapterKey"] = RegionAdapterKey,
+ ["CacheSize"] = "2048",
+ },
+ provider);
+ Assert.Equal("integration-profile", environment["AWS_PROFILE"]);
+ Assert.Equal("us-east-1", environment["AWS_REGION"]);
+ Assert.Equal("integration-profile", environment["AWS__Profile"]);
+ Assert.Equal("us-east-1", environment["AWS__Region"]);
+ Assert.DoesNotContain(configuration.Keys, key => key.StartsWith("ConnectionStrings:", StringComparison.Ordinal));
+ AssertSecretFree(environment);
+ var clientEnvironment = await app.GetClientEnvironmentAsync();
+ Assert.Equal(environment["AWS_PROFILE"], clientEnvironment["AWS_PROFILE"]);
+ Assert.Equal(environment["AWS_REGION"], clientEnvironment["AWS_REGION"]);
+ Assert.Equal(environment["AWS__Profile"], clientEnvironment["AWS__Profile"]);
+ Assert.Equal(environment["AWS__Region"], clientEnvironment["AWS__Region"]);
+ Assert.DoesNotContain(
+ app.Model.Resources,
+ resource => resource.GetType().Name.Contains("AWS", StringComparison.OrdinalIgnoreCase));
+ }
+
+ [Fact]
+ public async Task AspireAppModel_CustomEndpoint_ProducesWorkingClientConfiguration()
+ {
+ await using var app = await CreateEndpointAppAsync();
+ var environment = await app.GetClientEnvironmentAsync();
+ var configuration = SqsAspireTestApp.NormalizeConfiguration(environment);
+ var provider = GetProviderConfiguration(configuration, app.ProviderName);
+
+ AssertExactConfiguration(
+ new Dictionary
+ {
+ ["ProviderType"] = "SQS",
+ ["ServiceKey"] = ServiceKey,
+ ["PartitionCount"] = "3",
+ ["FifoQueue"] = "True",
+ ["ReceiveWaitTimeSeconds"] = "7",
+ ["VisibilityTimeoutSeconds"] = "90",
+ ["ReceiveMessageAttributes:0"] = "CorrelationId",
+ ["ReceiveMessageSystemAttributes:0"] = "ApproximateReceiveCount",
+ },
+ provider);
+ Assert.Equal(
+ "Service=http://127.0.0.1:9324",
+ configuration[$"ConnectionStrings:{ServiceKey}"]);
+ AssertSecretFree(environment);
+ Assert.Single(app.Model.Resources, resource => resource.Name == ServiceKey);
+ Assert.DoesNotContain(
+ app.Model.Resources,
+ resource => resource.Name.Contains("queue", StringComparison.OrdinalIgnoreCase)
+ || resource.GetType().Name.Contains("queue", StringComparison.OrdinalIgnoreCase));
+ }
+
+ [Fact]
+ public async Task AspireGeneratedConfiguration_ActivatesSqsProviderOnSilo()
+ {
+ await using var app = await CreateRegionAppAsync();
+ using var host = await app.BuildSiloHostAsync(services =>
+ services.AddKeyedSingleton(
+ RegionAdapterKey,
+ new FakeSqsDataAdapter(RegionAdapterKey)));
+ var sqsOptions = GetOptions(host.Services, app.ProviderName);
+ var partitionOptions = GetOptions(host.Services, app.ProviderName);
+ var cacheOptions = GetOptions(host.Services, app.ProviderName);
+ var adapterFactory = host.Services.GetRequiredKeyedService(app.ProviderName);
+ var streamProvider = host.Services.GetRequiredKeyedService(app.ProviderName);
+ var configuredAdapter = host.Services.GetRequiredKeyedService(app.ProviderName);
+
+ Assert.Equal("Service=us-east-1", sqsOptions.ConnectionString);
+ Assert.False(sqsOptions.FifoQueue);
+ Assert.Equal(4, partitionOptions.TotalQueueCount);
+ Assert.Equal(2048, cacheOptions.CacheSize);
+ Assert.IsType(adapterFactory);
+ Assert.Equal(app.ProviderName, streamProvider.Name);
+ Assert.Equal(RegionAdapterKey, Assert.IsType(configuredAdapter).Id);
+ }
+
+ [Fact]
+ public async Task AspireGeneratedConfiguration_ActivatesSqsProviderOnClient()
+ {
+ await using var app = await CreateEndpointAppAsync();
+ using var host = await app.BuildClientHostAsync();
+ var sqsOptions = GetOptions(host.Services, app.ProviderName);
+ var partitionOptions = GetOptions(host.Services, app.ProviderName);
+ var cacheOptions = GetOptions(host.Services, app.ProviderName);
+ var adapterFactory = host.Services.GetRequiredKeyedService(app.ProviderName);
+ var streamProvider = host.Services.GetRequiredKeyedService(app.ProviderName);
+
+ Assert.Equal("Service=http://127.0.0.1:9324", sqsOptions.ConnectionString);
+ Assert.True(sqsOptions.FifoQueue);
+ Assert.Equal(3, partitionOptions.TotalQueueCount);
+ Assert.Equal(SimpleQueueCacheOptions.DEFAULT_CACHE_SIZE, cacheOptions.CacheSize);
+ Assert.IsType(adapterFactory);
+ Assert.Equal(app.ProviderName, streamProvider.Name);
+ }
+
+ [Fact]
+ public async Task StreamingEnvironmentScope_PreservesTopologyAndAwsConfiguration()
+ {
+ await using var app = await CreateRegionAppAsync();
+ using var environment = await app.CreateEnvironmentScopeAsync(
+ SqsAspireResourceRole.Silo,
+ streamingOnly: true);
+
+ Assert.Equal(ServiceId, Environment.GetEnvironmentVariable("Orleans__ServiceId"));
+ Assert.NotNull(Environment.GetEnvironmentVariable("Orleans__ClusterId"));
+ Assert.Equal("integration-profile", Environment.GetEnvironmentVariable("AWS_PROFILE"));
+ Assert.Equal("us-east-1", Environment.GetEnvironmentVariable("AWS_REGION"));
+ Assert.Equal("us-east-1", Environment.GetEnvironmentVariable("AWS__Region"));
+ }
+
+ [Theory]
+ [InlineData(false)]
+ [InlineData(true)]
+ public async Task AspireGeneratedConfiguration_PreservesExplicitPartitionTopology(bool fifoQueue)
+ {
+ const string providerName = "Topology";
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ providerName,
+ [
+ ("PartitionCount", "4"),
+ ("FifoQueue", fifoQueue.ToString()),
+ ],
+ serviceId: ServiceId,
+ awsRegion: "eu-central-1");
+ using var siloHost = await app.BuildSiloHostAsync();
+ using var clientHost = await app.BuildClientHostAsync();
+ var siloSqsOptions = GetOptions(siloHost.Services, providerName);
+ var clientSqsOptions = GetOptions(clientHost.Services, providerName);
+ var siloMapper = CreateQueueMapper(siloHost.Services, providerName);
+ var clientMapper = CreateQueueMapper(clientHost.Services, providerName);
+ var siloQueueIds = siloMapper.GetAllQueues().Order().ToArray();
+ var clientQueueIds = clientMapper.GetAllQueues().Order().ToArray();
+ var siloPhysicalNames = GetPhysicalQueueNames(siloQueueIds, siloSqsOptions, ServiceId);
+ var clientPhysicalNames = GetPhysicalQueueNames(clientQueueIds, clientSqsOptions, ServiceId);
+ var expectedPhysicalNames = Enumerable.Range(0, 4)
+ .Select(index => $"{ServiceId}-topology-{index}{(fifoQueue ? ".fifo" : string.Empty)}")
+ .Order()
+ .ToArray();
+
+ Assert.Equal(4, siloQueueIds.Length);
+ Assert.Equal(siloQueueIds, clientQueueIds);
+ Assert.Equal(expectedPhysicalNames, siloPhysicalNames);
+ Assert.Equal(siloPhysicalNames, clientPhysicalNames);
+ Assert.All(
+ siloPhysicalNames,
+ name => Assert.Equal(fifoQueue, name.EndsWith(".fifo", StringComparison.Ordinal)));
+ }
+
+ private static Task CreateRegionAppAsync()
+ => SqsAspireTestApp.CreateAsync(
+ "orders-stream",
+ [
+ ("PartitionCount", "4"),
+ ("CacheSize", "2048"),
+ ("FifoQueue", "False"),
+ ("ReceiveWaitTimeSeconds", "12"),
+ ("VisibilityTimeoutSeconds", "45"),
+ ("ReceiveMessageAttributes:0", "TraceId"),
+ ("ReceiveMessageAttributes:1", "Tenant"),
+ ("ReceiveMessageSystemAttributes:0", "SentTimestamp"),
+ ("ReceiveMessageSystemAttributes:1", "SequenceNumber"),
+ ("DataAdapterKey", RegionAdapterKey),
+ ],
+ serviceId: ServiceId,
+ awsProfile: "integration-profile",
+ awsRegion: "us-east-1");
+
+ private static Task CreateEndpointAppAsync()
+ => SqsAspireTestApp.CreateAsync(
+ "critical-orders",
+ [
+ ("ServiceKey", ServiceKey),
+ ("PartitionCount", "3"),
+ ("FifoQueue", "True"),
+ ("ReceiveWaitTimeSeconds", "7"),
+ ("VisibilityTimeoutSeconds", "90"),
+ ("ReceiveMessageAttributes:0", "CorrelationId"),
+ ("ReceiveMessageSystemAttributes:0", "ApproximateReceiveCount"),
+ ],
+ [($"ConnectionStrings:{ServiceKey}", "Service=http://127.0.0.1:9324")],
+ ServiceId);
+
+ private static IReadOnlyDictionary GetProviderConfiguration(
+ IReadOnlyDictionary configuration,
+ string providerName)
+ {
+ var prefix = $"Orleans:Streaming:{providerName}:";
+ return configuration
+ .Where(pair => pair.Key.StartsWith(prefix, StringComparison.Ordinal))
+ .ToDictionary(
+ pair => pair.Key[prefix.Length..],
+ pair => pair.Value,
+ StringComparer.Ordinal);
+ }
+
+ private static void AssertExactConfiguration(
+ IReadOnlyDictionary expected,
+ IReadOnlyDictionary actual)
+ {
+ Assert.Equal(expected.Count, actual.Count);
+ foreach (var (key, expectedValue) in expected)
+ {
+ Assert.True(actual.TryGetValue(key, out var actualValue), $"Missing generated key '{key}'.");
+ Assert.Equal(expectedValue, actualValue);
+ }
+ }
+
+ private static void AssertSecretFree(IReadOnlyDictionary environment)
+ {
+ string[] forbiddenFragments = ["ACCESS_KEY", "SECRET", "SESSION_TOKEN", "QUEUE_URL"];
+ Assert.DoesNotContain(
+ environment,
+ pair => forbiddenFragments.Any(
+ fragment => pair.Key.Contains(fragment, StringComparison.OrdinalIgnoreCase)
+ || pair.Value?.Contains(fragment, StringComparison.OrdinalIgnoreCase) == true));
+ }
+
+ private static TOptions GetOptions(IServiceProvider services, string providerName)
+ where TOptions : class
+ => services.GetRequiredService>().Get(providerName);
+
+ private static HashRingBasedStreamQueueMapper CreateQueueMapper(IServiceProvider services, string providerName)
+ => new(GetOptions(services, providerName), providerName);
+
+ private static string[] GetPhysicalQueueNames(
+ IEnumerable queueIds,
+ SqsOptions options,
+ string serviceId)
+ {
+ return queueIds
+ .Select(queueId => SqsQueueName.Create(queueId.ToString(), options.FifoQueue, serviceId))
+ .Order()
+ .ToArray();
+ }
+
+ private sealed class FakeSqsDataAdapter(string id) : ISQSDataAdapter
+ {
+ public string Id { get; } = id;
+
+ public IBatchContainer FromQueueMessage(SqsMessage queueMessage, long sequenceId)
+ => throw new NotSupportedException();
+
+ public SqsMessage ToQueueMessage(
+ Orleans.Runtime.StreamId streamId,
+ IEnumerable events,
+ StreamSequenceToken? token,
+ Dictionary? requestContext)
+ => throw new NotSupportedException();
+ }
+}
diff --git a/test/Extensions/Orleans.AWS.Tests/Streaming/SQSAspireLiveStreamTests.cs b/test/Extensions/Orleans.AWS.Tests/Streaming/SQSAspireLiveStreamTests.cs
new file mode 100644
index 00000000000..a463044d61a
--- /dev/null
+++ b/test/Extensions/Orleans.AWS.Tests/Streaming/SQSAspireLiveStreamTests.cs
@@ -0,0 +1,125 @@
+using AWSUtils.Tests.StorageTests;
+using Microsoft.Extensions.Configuration;
+using Microsoft.Extensions.Logging.Abstractions;
+using Orleans.Hosting;
+using Orleans.Persistence.FileStorage;
+using Orleans.TestingHost;
+using OrleansAWSUtils.Streams;
+using TestExtensions;
+using UnitTests.Streaming;
+using UnitTests.StreamingTests;
+using Xunit;
+
+namespace AWSUtils.Tests.Streaming;
+
+[Collection(SQSStreamProviderBuilderTestCollection.CollectionName)]
+[TestSuite("Functional")]
+[TestProvider("SQS")]
+[TestArea("Streaming")]
+[TestCategory("AWS"), TestCategory("SQS")]
+public sealed class SQSAspireLiveStreamTests : TestClusterPerTest
+{
+ private const string ProviderName = "AspireSQS";
+ private const string PubSubStoreRootDirectoryKey = "AspireSQS:PubSubStoreRootDirectory";
+ private SqsAspireTestApp _app = null!;
+ private EnvironmentVariableScope _environment = null!;
+ private string? _ownedDirectory;
+ private SingleStreamTestRunner _runner = null!;
+
+ protected override void CheckPreconditionsOrThrow()
+ {
+ if (!AWSTestConstants.IsSqsAvailable)
+ {
+ throw Xunit.Sdk.SkipException.ForSkip("SQS connection string is not configured.");
+ }
+ }
+
+ protected override void ConfigureTestCluster(TestClusterBuilder builder)
+ {
+ var ownedDirectory = _ownedDirectory
+ ?? throw new InvalidOperationException("The PubSubStore directory has not been initialized.");
+ builder.Options.InitialSilosCount = 1;
+ builder.Properties[PubSubStoreRootDirectoryKey] = Path.Combine(
+ ownedDirectory,
+ "PubSubStore");
+ builder.ConfigureHostConfiguration(configuration => configuration.AddEnvironmentVariables());
+ builder.AddSiloBuilderConfigurator();
+ }
+
+ public override async ValueTask InitializeAsync()
+ {
+ EnsurePreconditionsMet();
+ _ownedDirectory = Path.Combine(
+ Path.GetTempPath(),
+ nameof(SQSAspireLiveStreamTests),
+ Guid.NewGuid().ToString("N"));
+ Directory.CreateDirectory(_ownedDirectory);
+ _app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [
+ ("ServiceKey", "orleans-sqs"),
+ ("PartitionCount", "4"),
+ ("FifoQueue", "false"),
+ ],
+ [("ConnectionStrings:orleans-sqs", AWSTestConstants.SqsConnectionString)]);
+ _environment = await _app.CreateEnvironmentScopeAsync(
+ SqsAspireResourceRole.Silo,
+ streamingOnly: true);
+ await base.InitializeAsync();
+ _runner = new SingleStreamTestRunner(HostedCluster, ProviderName);
+ }
+
+ [Fact]
+ public Task AspireGeneratedConfiguration_PublishesConsumesAndAcknowledgesElasticMqStream()
+ {
+#pragma warning disable xUnit1051 // The shared stream runner does not expose a cancellation-token overload.
+ return _runner.StreamTest_04_OneProducerClientOneConsumerClient();
+#pragma warning restore xUnit1051
+ }
+
+ public override async ValueTask DisposeAsync()
+ {
+ if (!PreconditionsMet)
+ {
+ return;
+ }
+
+ try
+ {
+ if (HostedCluster is not null)
+ {
+ var serviceId = HostedCluster.Options.ServiceId;
+ await base.DisposeAsync();
+ await SQSStreamProviderUtils.DeleteAllUsedQueues(
+ ProviderName,
+ serviceId,
+ AWSTestConstants.SqsConnectionString,
+ NullLoggerFactory.Instance);
+ }
+ }
+ finally
+ {
+ _environment?.Dispose();
+ if (_app is not null)
+ {
+ await _app.DisposeAsync();
+ }
+
+ if (_ownedDirectory is not null && Directory.Exists(_ownedDirectory))
+ {
+ Directory.Delete(_ownedDirectory, recursive: true);
+ }
+ }
+ }
+
+ private sealed class SiloConfigurator : ISiloConfigurator
+ {
+ public void Configure(ISiloBuilder siloBuilder)
+ => siloBuilder.AddFileGrainStorage(
+ "PubSubStore",
+ options => options.RootDirectory =
+ siloBuilder.Configuration[PubSubStoreRootDirectoryKey]
+ ?? throw new InvalidOperationException(
+ $"Missing {PubSubStoreRootDirectoryKey} configuration."));
+ }
+}
diff --git a/test/Extensions/Orleans.AWS.Tests/Streaming/SQSStreamProviderBuilderTests.cs b/test/Extensions/Orleans.AWS.Tests/Streaming/SQSStreamProviderBuilderTests.cs
new file mode 100644
index 00000000000..c72b4184d73
--- /dev/null
+++ b/test/Extensions/Orleans.AWS.Tests/Streaming/SQSStreamProviderBuilderTests.cs
@@ -0,0 +1,507 @@
+using System.Reflection;
+using Amazon.Runtime;
+using Amazon.SQS;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Logging.Abstractions;
+using Microsoft.Extensions.Options;
+using Orleans.Configuration;
+using Orleans.Hosting;
+using Orleans.Providers;
+using Orleans.Streaming.SQS.Streams;
+using Orleans.Streams;
+using OrleansAWSUtils.Storage;
+using OrleansAWSUtils.Streams;
+using SqsMessage = Amazon.SQS.Model.Message;
+using Xunit;
+
+namespace AWSUtils.Tests.Streaming;
+
+[CollectionDefinition(CollectionName, DisableParallelization = true)]
+public sealed class SQSStreamProviderBuilderTestCollection
+{
+ public const string CollectionName = "SQS stream provider builder tests";
+}
+
+[Collection(SQSStreamProviderBuilderTestCollection.CollectionName)]
+[TestSuite("BVT")]
+[TestProvider("SQS")]
+[TestArea("Streaming")]
+[TestCategory("AWS"), TestCategory("SQS")]
+public sealed class SQSStreamProviderBuilderTests
+{
+ private const string ProviderName = "orders";
+
+ [Fact]
+ public void Assembly_RegistersSqsProviderForSiloAndClient()
+ {
+ var registrations = typeof(SqsStreamProviderBuilder)
+ .Assembly
+ .GetCustomAttributes()
+ .Where(attribute => attribute.Type == typeof(SqsStreamProviderBuilder))
+ .Select(attribute => (attribute.Name, attribute.Kind, attribute.Target))
+ .ToHashSet();
+
+ Assert.Equal(4, registrations.Count);
+ Assert.Contains(("SQS", "Streaming", "Silo"), registrations);
+ Assert.Contains(("SQS", "Streaming", "Client"), registrations);
+ Assert.Contains(("AmazonSQS", "Streaming", "Silo"), registrations);
+ Assert.Contains(("AmazonSQS", "Streaming", "Client"), registrations);
+ }
+
+ [Fact]
+ public void SqsProvider_UsesAwsSdkV4()
+ => Assert.Equal(4, typeof(AmazonSQSClient).Assembly.GetName().Version?.Major);
+
+ [Fact]
+ public async Task AspireAwsSdkResource_ConfiguresRegionWithoutCredentialMaterial()
+ {
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [],
+ awsProfile: "integration-profile",
+ awsRegion: "ap-southeast-2");
+ using var host = await app.BuildSiloHostAsync();
+ var options = GetOptions(host.Services, ProviderName);
+
+ Assert.Equal("Service=ap-southeast-2", options.ConnectionString);
+ Assert.DoesNotContain("integration-profile", options.ConnectionString, StringComparison.Ordinal);
+ Assert.DoesNotContain(
+ typeof(SqsOptions).GetProperties(),
+ property => property.Name.Contains("Credential", StringComparison.OrdinalIgnoreCase)
+ || property.Name.Contains("AccessKey", StringComparison.OrdinalIgnoreCase)
+ || property.Name.Contains("SecretKey", StringComparison.OrdinalIgnoreCase)
+ || property.Name.Contains("Profile", StringComparison.OrdinalIgnoreCase));
+ }
+
+ [Theory]
+ [InlineData("ServiceKey")]
+ [InlineData("ConnectionName")]
+ public async Task AspireConnectionResource_ResolvesReferencedConnectionString(string referenceKey)
+ {
+ const string serviceKey = "shared-sqs";
+ const string connectionString = "Service=http://localhost:9324";
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [(referenceKey, serviceKey)],
+ [($"ConnectionStrings:{serviceKey}", connectionString)]);
+ using var host = await app.BuildClientHostAsync();
+ var options = GetOptions(host.Services, ProviderName);
+
+ Assert.Equal(connectionString, options.ConnectionString);
+ Assert.DoesNotContain(ProviderName, options.ConnectionString, StringComparison.OrdinalIgnoreCase);
+ Assert.DoesNotContain("/queue/", options.ConnectionString, StringComparison.OrdinalIgnoreCase);
+ }
+
+ [Fact]
+ public async Task AspireConnectionResource_EquivalentReferenceAliasesResolve()
+ {
+ const string serviceKey = "shared-sqs";
+ const string connectionString = "Service=http://localhost:9324";
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [
+ ("ServiceKey", serviceKey.ToUpperInvariant()),
+ ("ConnectionName", serviceKey),
+ ],
+ [($"ConnectionStrings:{serviceKey}", connectionString)]);
+ using var host = await app.BuildClientHostAsync();
+ var options = GetOptions(host.Services, ProviderName);
+
+ Assert.Equal(connectionString, options.ConnectionString);
+ }
+
+ [Fact]
+ public async Task AspireConnectionResource_CustomEndpointPreservesExplicitCredentials()
+ {
+ const string connectionString =
+ "Service=https://sqs.example.com;AccessKey=access;SecretKey=secret;SessionToken=token";
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [("ConnectionString", connectionString)]);
+ using var host = await app.BuildClientHostAsync();
+ var options = GetOptions(host.Services, ProviderName);
+
+ Assert.Equal(connectionString, options.ConnectionString);
+ }
+
+ [Theory]
+ [InlineData(false)]
+ [InlineData(true)]
+ public void SqsStorage_CustomEndpointUsesExplicitCredentials(bool useSessionToken)
+ {
+ var connectionString = "Service=https://sqs.example.com;AccessKey=access;SecretKey=secret";
+ if (useSessionToken)
+ {
+ connectionString += ";SessionToken=token";
+ }
+
+ var storage = new SQSStorage(
+ NullLoggerFactory.Instance,
+ "queue",
+ new SqsOptions { ConnectionString = connectionString },
+ "service");
+ var client = Assert.IsType(
+ typeof(SQSStorage)
+ .GetField("sqsClient", BindingFlags.Instance | BindingFlags.NonPublic)!
+ .GetValue(storage));
+ var credentials = Assert.IsAssignableFrom(GetExplicitCredentials(client));
+ var immutableCredentials = credentials.GetCredentials();
+
+ Assert.Equal("access", immutableCredentials.AccessKey);
+ Assert.Equal("secret", immutableCredentials.SecretKey);
+ Assert.Equal(useSessionToken ? "token" : string.Empty, immutableCredentials.Token);
+ }
+
+ [Fact]
+ public void SqsStorage_HttpEndpointUsesDummyCredentials()
+ {
+ var storage = new SQSStorage(
+ NullLoggerFactory.Instance,
+ "queue",
+ new SqsOptions { ConnectionString = "Service=http://localhost:9324" },
+ "service");
+ var client = Assert.IsType(
+ typeof(SQSStorage)
+ .GetField("sqsClient", BindingFlags.Instance | BindingFlags.NonPublic)!
+ .GetValue(storage));
+ var credentials = Assert.IsType(GetExplicitCredentials(client));
+ var immutableCredentials = credentials.GetCredentials();
+
+ Assert.Equal("dummy", immutableCredentials.AccessKey);
+ Assert.Equal("dummyKey", immutableCredentials.SecretKey);
+ }
+
+ [Fact]
+ public void SqsStorage_HttpsEndpointUsesDefaultCredentialChain()
+ {
+ var storage = new SQSStorage(
+ NullLoggerFactory.Instance,
+ "queue",
+ new SqsOptions { ConnectionString = "Service=https://sqs.example.com" },
+ "service");
+ var client = Assert.IsType(
+ typeof(SQSStorage)
+ .GetField("sqsClient", BindingFlags.Instance | BindingFlags.NonPublic)!
+ .GetValue(storage));
+
+ Assert.Null(GetExplicitCredentials(client));
+ }
+
+ [Theory]
+ [InlineData("DataAdapterServiceKey")]
+ [InlineData("DataAdapterKey")]
+ public async Task AspireGeneratedConfiguration_ResolvesKeyedDataAdapter(string adapterConfigurationKey)
+ {
+ const string adapterServiceKey = "custom-adapter";
+ var defaultAdapter = new FakeSqsDataAdapter("default");
+ var keyedAdapter = new FakeSqsDataAdapter("keyed");
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [(adapterConfigurationKey, adapterServiceKey)],
+ awsRegion: "us-east-1");
+ using var host = await app.BuildSiloHostAsync(services =>
+ {
+ services.AddSingleton(defaultAdapter);
+ services.AddKeyedSingleton(adapterServiceKey, keyedAdapter);
+ });
+ var configuredAdapter = host.Services.GetRequiredKeyedService(ProviderName);
+ var adapterFactory = SQSAdapterFactory.Create(host.Services, ProviderName);
+ var factoryAdapter = typeof(SQSAdapterFactory)
+ .GetField("dataAdapter", BindingFlags.Instance | BindingFlags.NonPublic)!
+ .GetValue(adapterFactory);
+
+ Assert.Same(keyedAdapter, configuredAdapter);
+ Assert.NotSame(defaultAdapter, configuredAdapter);
+ Assert.Same(keyedAdapter, factoryAdapter);
+ }
+
+ [Fact]
+ public async Task AspireGeneratedConfiguration_OrdersIndexedAttributeListsNumerically()
+ {
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ [
+ ("Region", "us-east-1"),
+ ("ReceiveMessageAttributes:0", "zero"),
+ ("ReceiveMessageAttributes:1", "one"),
+ ("ReceiveMessageAttributes:10", "ten"),
+ ("ReceiveMessageAttributes:2", "two"),
+ ("ReceiveMessageSystemAttributes:0", "system-zero"),
+ ("ReceiveMessageSystemAttributes:1", "system-one"),
+ ("ReceiveMessageSystemAttributes:10", "system-ten"),
+ ("ReceiveMessageSystemAttributes:2", "system-two"),
+ ]);
+ using var host = await app.BuildSiloHostAsync();
+ var options = GetOptions(host.Services, ProviderName);
+
+ Assert.Equal(["zero", "one", "two", "ten"], options.ReceiveMessageAttributes);
+ Assert.Equal(
+ ["system-zero", "system-one", "system-two", "system-ten"],
+ options.ReceiveMessageSystemAttributes);
+ }
+
+ [Theory]
+ [MemberData(nameof(InvalidConfigurations))]
+ public async Task AspireGeneratedConfiguration_InvalidSiloConfigurationThrowsActionableError(
+ string caseId,
+ string[] providerValues,
+ string[] rootValues,
+ string expectedMessage)
+ {
+ await using var app = await SqsAspireTestApp.CreateAsync(
+ ProviderName,
+ ToPairs(providerValues),
+ ToPairs(rootValues));
+
+ var exception = await Assert.ThrowsAsync(async () =>
+ {
+ using var host = await app.BuildSiloHostAsync();
+ _ = GetOptions(host.Services, ProviderName);
+ });
+
+ Assert.NotEmpty(caseId);
+ Assert.Contains(expectedMessage, exception.Message, StringComparison.OrdinalIgnoreCase);
+ }
+
+ [Fact]
+ public async Task AspireGeneratedConfiguration_MissingServiceLocationThrowsActionableError()
+ {
+ await using var app = await SqsAspireTestApp.CreateAsync(ProviderName, []);
+ using var host = await app.BuildClientHostAsync();
+
+ var exception = Assert.Throws(
+ () => GetOptions(host.Services, ProviderName));
+
+ Assert.Contains("SQS streaming", exception.Message, StringComparison.OrdinalIgnoreCase);
+ Assert.Contains("service location", exception.Message, StringComparison.OrdinalIgnoreCase);
+ Assert.Contains("ServiceKey", exception.Message, StringComparison.Ordinal);
+ Assert.Contains("AWS_REGION", exception.Message, StringComparison.Ordinal);
+ Assert.Contains("AWS_DEFAULT_REGION", exception.Message, StringComparison.Ordinal);
+ }
+
+ public static TheoryData InvalidConfigurations => new()
+ {
+ {
+ "ConnectionStringMalformed",
+ ["ConnectionString", "not-a-connection-string"],
+ [],
+ "key=value"
+ },
+ {
+ "ConnectionStringMissingService",
+ ["ConnectionString", "AccessKey=access;SecretKey=secret"],
+ [],
+ "non-empty Service"
+ },
+ {
+ "ConnectionStringMissingSecretKey",
+ ["ConnectionString", "Service=us-east-1;AccessKey=access"],
+ [],
+ "AccessKey and SecretKey together"
+ },
+ {
+ "ConnectionStringMissingAccessKey",
+ ["ConnectionString", "Service=us-east-1;SecretKey=secret"],
+ [],
+ "AccessKey and SecretKey together"
+ },
+ {
+ "ConnectionStringSessionTokenWithoutCredentials",
+ ["ConnectionString", "Service=us-east-1;SessionToken=token"],
+ [],
+ "SessionToken"
+ },
+ {
+ "ConflictingServiceKeyAndConnectionName",
+ ["ServiceKey", "primary", "ConnectionName", "secondary", "Region", "us-east-1"],
+ [],
+ "ServiceKey and ConnectionName"
+ },
+ {
+ "ConflictingConnectionStringAndLocation",
+ ["ConnectionString", "Service=us-east-1", "Region", "us-west-2"],
+ [],
+ "connection string and a service location"
+ },
+ {
+ "ConflictingRegionAndServiceEndpoint",
+ ["Region", "us-east-1", "ServiceEndpoint", "http://localhost:9324"],
+ [],
+ "multiple values among Service, Region, and ServiceEndpoint"
+ },
+ {
+ "ConflictingServiceEndpointAliases",
+ [
+ "ServiceEndpoint", "http://localhost:9324",
+ "Endpoint", "http://localhost:9325",
+ ],
+ [],
+ "ServiceEndpoint and Endpoint"
+ },
+ {
+ "ConflictingServiceAndRegion",
+ ["Service", "us-east-1", "Region", "us-west-2"],
+ [],
+ "multiple values among Service, Region, and ServiceEndpoint"
+ },
+ {
+ "ConflictingServiceAndEndpoint",
+ ["Service", "us-east-1", "Endpoint", "http://localhost:9324"],
+ [],
+ "multiple values among Service, Region, and ServiceEndpoint"
+ },
+ {
+ "ConflictingAwsRegionVariables",
+ [],
+ ["AWS_REGION", "us-east-1", "AWS_DEFAULT_REGION", "us-west-2"],
+ "AWS_REGION and AWS_DEFAULT_REGION"
+ },
+ {
+ "FifoQueueNotBoolean",
+ ["Region", "us-east-1", "FifoQueue", "sometimes"],
+ [],
+ "'FifoQueue' must be true or false"
+ },
+ {
+ "PartitionCountNotInteger",
+ ["Region", "us-east-1", "PartitionCount", "many"],
+ [],
+ "'PartitionCount' must be an integer"
+ },
+ {
+ "PartitionCountNotPositive",
+ ["Region", "us-east-1", "PartitionCount", "0"],
+ [],
+ "'PartitionCount' must be greater than zero"
+ },
+ {
+ "CacheSizeNotInteger",
+ ["Region", "us-east-1", "CacheSize", "large"],
+ [],
+ "'CacheSize' must be an integer"
+ },
+ {
+ "CacheSizeNotPositive",
+ ["Region", "us-east-1", "CacheSize", "-1"],
+ [],
+ "'CacheSize' must be greater than zero"
+ },
+ {
+ "ReceiveWaitTimeSecondsNegative",
+ ["Region", "us-east-1", "ReceiveWaitTimeSeconds", "-1"],
+ [],
+ "'ReceiveWaitTimeSeconds' must be between 0 and 20"
+ },
+ {
+ "ReceiveWaitTimeSecondsAboveMaximum",
+ ["Region", "us-east-1", "ReceiveWaitTimeSeconds", "21"],
+ [],
+ "'ReceiveWaitTimeSeconds' must be between 0 and 20"
+ },
+ {
+ "VisibilityTimeoutSecondsNegative",
+ ["Region", "us-east-1", "VisibilityTimeoutSeconds", "-1"],
+ [],
+ "'VisibilityTimeoutSeconds' must be between 0 and 43200"
+ },
+ {
+ "VisibilityTimeoutSecondsAboveMaximum",
+ ["Region", "us-east-1", "VisibilityTimeoutSeconds", "43201"],
+ [],
+ "'VisibilityTimeoutSeconds' must be between 0 and 43200"
+ },
+ {
+ "ConflictingDataAdapterKeys",
+ ["Region", "us-east-1", "DataAdapterKey", "primary", "DataAdapterServiceKey", "secondary"],
+ [],
+ "DataAdapterKey and DataAdapterServiceKey"
+ },
+ {
+ "DataAdapterKeyMatchesProviderName",
+ ["Region", "us-east-1", "DataAdapterKey", ProviderName],
+ [],
+ "DataAdapterKey must differ from the stream provider name"
+ },
+ {
+ "ConnectionStringAndReference",
+ ["ConnectionString", "Service=us-east-1", "ServiceKey", "shared-sqs"],
+ ["ConnectionStrings:shared-sqs", "Service=us-west-2"],
+ "both a connection reference and ConnectionString"
+ },
+ {
+ "ConnectionReferenceNotFound",
+ ["ServiceKey", "missing-sqs"],
+ [],
+ "connection reference 'missing-sqs' did not resolve"
+ },
+ {
+ "ServiceEndpointNotAbsoluteHttp",
+ ["ServiceEndpoint", "localhost:9324"],
+ [],
+ "ServiceEndpoint values must be absolute HTTP or HTTPS URIs"
+ },
+ {
+ "ServiceNotAbsoluteHttp",
+ ["Service", "localhost:9324"],
+ [],
+ "service endpoint values must be absolute HTTP or HTTPS URIs"
+ },
+ {
+ "RegionNotRecognized",
+ ["Region", "us-east-99"],
+ [],
+ "region 'us-east-99' is not recognized"
+ },
+ {
+ "ConnectionStringServiceNotAbsoluteHttp",
+ ["ConnectionString", "Service=localhost:9324"],
+ [],
+ "service endpoint values must be absolute HTTP or HTTPS URIs"
+ },
+ {
+ "ConnectionStringRegionNotRecognized",
+ ["ConnectionString", "Service=us-east-99"],
+ [],
+ "region 'us-east-99' is not recognized"
+ },
+ {
+ "ConnectionStringDuplicateService",
+ ["ConnectionString", "Service=us-east-1;Service=us-west-2"],
+ [],
+ "property 'Service' is configured more than once"
+ },
+ };
+
+ private static (string Key, string? Value)[] ToPairs(string[] values)
+ {
+ Assert.Equal(0, values.Length % 2);
+ return values
+ .Chunk(2)
+ .Select(pair => (pair[0], (string?)pair[1]))
+ .ToArray();
+ }
+
+ private static TOptions GetOptions(IServiceProvider services, string providerName)
+ where TOptions : class
+ => services.GetRequiredService>().Get(providerName);
+
+ private static object? GetExplicitCredentials(AmazonServiceClient client)
+ => typeof(AmazonServiceClient)
+ .GetProperty("ExplicitAWSCredentials", BindingFlags.Instance | BindingFlags.NonPublic)!
+ .GetValue(client);
+
+ private sealed class FakeSqsDataAdapter(string id) : ISQSDataAdapter
+ {
+ public string Id { get; } = id;
+
+ public IBatchContainer FromQueueMessage(SqsMessage queueMessage, long sequenceId)
+ => throw new NotSupportedException();
+
+ public SqsMessage ToQueueMessage(
+ Orleans.Runtime.StreamId streamId,
+ IEnumerable events,
+ StreamSequenceToken? token,
+ Dictionary? requestContext)
+ => throw new NotSupportedException();
+ }
+}
diff --git a/test/Extensions/Orleans.AWS.Tests/Streaming/SqsAspireTestApp.cs b/test/Extensions/Orleans.AWS.Tests/Streaming/SqsAspireTestApp.cs
new file mode 100644
index 00000000000..d4ad56392b7
--- /dev/null
+++ b/test/Extensions/Orleans.AWS.Tests/Streaming/SqsAspireTestApp.cs
@@ -0,0 +1,332 @@
+using Aspire.Hosting;
+using Aspire.Hosting.ApplicationModel;
+using Aspire.Hosting.Orleans;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Hosting;
+using Orleans.Hosting;
+
+namespace AWSUtils.Tests.Streaming;
+
+internal sealed class SqsAspireTestApp : IAsyncDisposable
+{
+ private const string SiloResourceName = "silo";
+ private const string ClientResourceName = "client";
+ private readonly DistributedApplication _application;
+ private readonly IResource _silo;
+ private readonly IResource _client;
+
+ private SqsAspireTestApp(
+ DistributedApplication application,
+ IResource silo,
+ IResource client,
+ string providerName,
+ string serviceId)
+ {
+ _application = application;
+ _silo = silo;
+ _client = client;
+ ProviderName = providerName;
+ ServiceId = serviceId;
+ }
+
+ public string ProviderName { get; }
+
+ public string ServiceId { get; }
+
+ public DistributedApplicationModel Model
+ => _application.Services.GetRequiredService();
+
+ public static Task CreateAsync(
+ string providerName,
+ IEnumerable<(string Key, string? Value)> providerValues,
+ IEnumerable<(string Key, string? Value)>? rootValues = null,
+ string serviceId = "aspire-sqs-service",
+ string? awsProfile = null,
+ string? awsRegion = null)
+ {
+ var builder = DistributedApplication.CreateBuilder(
+ new DistributedApplicationOptions
+ {
+ Args = [],
+ DisableDashboard = true,
+ });
+ var connectionResources = new List>();
+ var environmentValues = new List<(string Key, string? Value)>();
+ foreach (var (key, value) in rootValues ?? [])
+ {
+ const string connectionStringPrefix = "ConnectionStrings:";
+ if (key.StartsWith(connectionStringPrefix, StringComparison.OrdinalIgnoreCase))
+ {
+ connectionResources.Add(
+ builder.AddConnectionString(
+ key[connectionStringPrefix.Length..],
+ ReferenceExpression.Create($"{value}")));
+ }
+ else
+ {
+ environmentValues.Add((key, value));
+ }
+ }
+
+ var provider = new SqsProviderConfiguration(
+ awsProfile,
+ awsRegion,
+ connectionResources,
+ providerValues,
+ environmentValues);
+ var orleans = builder.AddOrleans("cluster")
+ .WithClustering(new TestClusteringConfiguration())
+ .WithServiceId(serviceId)
+ .WithStreaming(providerName, provider);
+ var silo = builder.AddContainer(SiloResourceName, "unused")
+ .WithReference(orleans);
+ var client = builder.AddContainer(ClientResourceName, "unused")
+ .WithReference(orleans.AsClient());
+ var application = builder.Build();
+
+ return Task.FromResult(
+ new SqsAspireTestApp(
+ application,
+ silo.Resource,
+ client.Resource,
+ providerName,
+ serviceId));
+ }
+
+ public Task> GetSiloEnvironmentAsync()
+ => GetEnvironmentVariablesAsync(_silo);
+
+ public Task> GetClientEnvironmentAsync()
+ => GetEnvironmentVariablesAsync(_client);
+
+ public async Task CreateEnvironmentScopeAsync(
+ SqsAspireResourceRole role,
+ bool streamingOnly = false)
+ {
+ var resource = role == SqsAspireResourceRole.Silo ? _silo : _client;
+ var values = await GetEnvironmentVariablesAsync(resource);
+ if (streamingOnly)
+ {
+ var streamingPrefix = $"Orleans__Streaming__{ProviderName}__";
+ values = values
+ .Where(pair => pair.Key.StartsWith(streamingPrefix, StringComparison.Ordinal)
+ || pair.Key is "Orleans__ClusterId" or "Orleans__ServiceId"
+ || pair.Key.StartsWith("ConnectionStrings__", StringComparison.Ordinal)
+ || pair.Key.StartsWith("AWS_", StringComparison.Ordinal)
+ || pair.Key.StartsWith("AWS__", StringComparison.Ordinal))
+ .ToDictionary(StringComparer.Ordinal);
+ }
+
+ return new EnvironmentVariableScope(values);
+ }
+
+ public async Task BuildSiloHostAsync(Action? configureServices = null)
+ {
+ using var environment = await CreateEnvironmentScopeAsync(SqsAspireResourceRole.Silo);
+ var hostBuilder = Host.CreateApplicationBuilder();
+ configureServices?.Invoke(hostBuilder.Services);
+ hostBuilder.UseOrleans();
+ return hostBuilder.Build();
+ }
+
+ public async Task BuildClientHostAsync(Action? configureServices = null)
+ {
+ using var environment = await CreateEnvironmentScopeAsync(SqsAspireResourceRole.Client);
+ var hostBuilder = Host.CreateApplicationBuilder();
+ configureServices?.Invoke(hostBuilder.Services);
+ hostBuilder.UseOrleansClient();
+ return hostBuilder.Build();
+ }
+
+ public async ValueTask DisposeAsync()
+ => await _application.DisposeAsync();
+
+ public static IReadOnlyDictionary NormalizeConfiguration(
+ IReadOnlyDictionary environment)
+ => environment.ToDictionary(
+ pair => pair.Key.StartsWith("Orleans__", StringComparison.Ordinal)
+ || pair.Key.StartsWith("ConnectionStrings__", StringComparison.Ordinal)
+ ? pair.Key.Replace("__", ":", StringComparison.Ordinal)
+ : pair.Key,
+ pair => pair.Value,
+ StringComparer.Ordinal);
+
+ private async Task> GetEnvironmentVariablesAsync(IResource resource)
+ {
+ var executionContext = new DistributedApplicationExecutionContext(
+ new DistributedApplicationExecutionContextOptions(DistributedApplicationOperation.Run)
+ {
+ ServiceProvider = _application.Services,
+ });
+ var values = new Dictionary();
+ var callbackContext = new EnvironmentCallbackContext(executionContext, resource, values);
+ foreach (var annotation in resource.Annotations.OfType())
+ {
+ await annotation.Callback(callbackContext).WaitAsync(TimeSpan.FromSeconds(10));
+ }
+
+ var valueContext = new ValueProviderContext
+ {
+ Caller = resource,
+ ExecutionContext = executionContext,
+ Network = KnownNetworkIdentifiers.LocalhostNetwork,
+ };
+ var result = new Dictionary(StringComparer.Ordinal);
+ foreach (var (key, value) in values)
+ {
+ if (!IsRelevantEnvironmentVariable(key))
+ {
+ continue;
+ }
+
+ try
+ {
+ result[key] = value switch
+ {
+ IValueProvider provider => await provider
+ .GetValueAsync(valueContext)
+ .AsTask()
+ .WaitAsync(TimeSpan.FromSeconds(10)),
+ _ => value.ToString(),
+ };
+ }
+ catch (TimeoutException exception)
+ {
+ throw new TimeoutException(
+ $"Timed out resolving environment variable '{key}' for resource '{resource.Name}'.",
+ exception);
+ }
+ }
+
+ return result;
+ }
+
+ private static bool IsRelevantEnvironmentVariable(string name)
+ => name.StartsWith("Orleans__Streaming__", StringComparison.Ordinal)
+ || name.StartsWith("Orleans__Clustering__", StringComparison.Ordinal)
+ || name is "Orleans__ClusterId" or "Orleans__ServiceId"
+ || name.StartsWith("ConnectionStrings__", StringComparison.Ordinal)
+ || name.StartsWith("AWS_", StringComparison.Ordinal)
+ || name.StartsWith("AWS__", StringComparison.Ordinal);
+
+ private sealed class TestClusteringConfiguration : IProviderConfiguration
+ {
+ public void ConfigureResource(
+ IResourceBuilder resourceBuilder,
+ string configSectionPath)
+ where T : IResourceWithEnvironment
+ {
+ var prefix = $"Orleans__{configSectionPath.Replace(":", "__", StringComparison.Ordinal)}";
+ resourceBuilder.WithEnvironment($"{prefix}__ProviderType", "Development");
+ if (resourceBuilder.Resource.Name == SiloResourceName)
+ {
+ resourceBuilder.WithEnvironment(
+ $"{prefix}__PrimarySiloEndPoint",
+ "127.0.0.1:11111");
+ }
+ else
+ {
+ resourceBuilder.WithEnvironment(
+ $"{prefix}__Gateways__0",
+ "gwy.tcp://127.0.0.1:30000/0");
+ }
+ }
+ }
+
+ private sealed class SqsProviderConfiguration(
+ string? awsProfile,
+ string? awsRegion,
+ IReadOnlyList> connections,
+ IEnumerable<(string Key, string? Value)> providerValues,
+ IEnumerable<(string Key, string? Value)> environmentValues) : IProviderConfiguration
+ {
+ private readonly (string Key, string? Value)[] _providerValues = providerValues.ToArray();
+ private readonly (string Key, string? Value)[] _environmentValues = environmentValues.ToArray();
+
+ public void ConfigureResource(
+ IResourceBuilder resourceBuilder,
+ string configSectionPath)
+ where T : IResourceWithEnvironment
+ {
+ var prefix = $"Orleans__{configSectionPath.Replace(":", "__", StringComparison.Ordinal)}";
+ resourceBuilder.WithEnvironment($"{prefix}__ProviderType", "SQS");
+ if (awsProfile is not null)
+ {
+ resourceBuilder
+ .WithEnvironment("AWS_PROFILE", awsProfile)
+ .WithEnvironment("AWS__Profile", awsProfile);
+ }
+
+ if (awsRegion is not null)
+ {
+ resourceBuilder
+ .WithEnvironment("AWS_REGION", awsRegion)
+ .WithEnvironment("AWS__Region", awsRegion);
+ }
+
+ foreach (var connection in connections)
+ {
+ resourceBuilder.WithReference(connection);
+ }
+
+ foreach (var (key, value) in _providerValues)
+ {
+ resourceBuilder.WithEnvironment(
+ $"{prefix}__{key.Replace(":", "__", StringComparison.Ordinal)}",
+ value);
+ }
+
+ foreach (var (key, value) in _environmentValues)
+ {
+ resourceBuilder.WithEnvironment(
+ key.Replace(":", "__", StringComparison.Ordinal),
+ value);
+ }
+ }
+ }
+}
+
+internal enum SqsAspireResourceRole
+{
+ Silo,
+ Client,
+}
+
+internal sealed class EnvironmentVariableScope : IDisposable
+{
+ private readonly Dictionary _previousValues = new(StringComparer.OrdinalIgnoreCase);
+
+ public EnvironmentVariableScope(IReadOnlyDictionary values)
+ {
+ foreach (var key in new[]
+ {
+ "AWS_REGION",
+ "AWS_DEFAULT_REGION",
+ "AWS_PROFILE",
+ "AWS__Region",
+ "AWS__Profile",
+ })
+ {
+ SaveAndSet(key, null);
+ }
+
+ foreach (var (key, value) in values)
+ {
+ SaveAndSet(key, value);
+ }
+ }
+
+ public void Dispose()
+ {
+ foreach (var (key, value) in _previousValues)
+ {
+ Environment.SetEnvironmentVariable(key, value);
+ }
+ }
+
+ private void SaveAndSet(string key, string? value)
+ {
+ _previousValues.TryAdd(key, Environment.GetEnvironmentVariable(key));
+ Environment.SetEnvironmentVariable(key, value);
+ }
+}
diff --git a/test/Extensions/Orleans.AWS.Tests/Streaming/SqsStreamingResourceTests.cs b/test/Extensions/Orleans.AWS.Tests/Streaming/SqsStreamingResourceTests.cs
new file mode 100644
index 00000000000..93db0e65031
--- /dev/null
+++ b/test/Extensions/Orleans.AWS.Tests/Streaming/SqsStreamingResourceTests.cs
@@ -0,0 +1,840 @@
+using Amazon;
+using Amazon.CDK.AWS.SQS;
+using Aspire.Hosting;
+using Aspire.Hosting.ApplicationModel;
+using Aspire.Hosting.AWS.CDK;
+using Aspire.Hosting.Orleans;
+using Microsoft.Extensions.DependencyInjection;
+using Microsoft.Extensions.Hosting;
+using Microsoft.Extensions.Options;
+using Orleans.Configuration;
+using Orleans.Hosting;
+using Orleans.Streams;
+using OrleansAWSUtils.Storage;
+using TestExtensions;
+using Xunit;
+
+namespace AWSUtils.Tests.Streaming;
+
+[Collection(SQSStreamProviderBuilderTestCollection.CollectionName)]
+[TestSuite("BVT")]
+[TestProvider("SQS")]
+[TestArea("Streaming")]
+[TestCategory("AWS"), TestCategory("SQS"), TestCategory("BVT")]
+public sealed class SqsStreamingResourceTests
+{
+ private const string ProviderName = "Orders";
+ private const string ServiceId = "orders-service";
+
+ [Fact]
+ public async Task AddSqsStreaming_StandardTopology_ConfiguresCdkEnvironmentAndWaits()
+ {
+ await using var app = await CreateAppAsync(
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ PartitionCount = 3,
+ ReceiveWaitTimeSeconds = 12,
+ VisibilityTimeoutSeconds = 45,
+ CacheSize = 2048,
+ DataAdapterKey = "orders-adapter",
+ ReceiveMessageAttributes = ["TraceId", "Tenant"],
+ ReceiveMessageSystemAttributes = ["SentTimestamp"],
+ });
+
+ Assert.Equal(ProviderName, app.Streaming.Name);
+ Assert.Equal(ServiceId, app.Streaming.Options.ServiceId);
+ Assert.Equal("us-east-1", app.Streaming.AwsSdkConfig.Region?.SystemName);
+ Assert.Equal(3, app.Streaming.Queues.Count);
+ Assert.Equal("cluster-orders-sqs", app.Streaming.Stack.Resource.Name);
+
+ var queues = GetQueues(app.Streaming);
+ Assert.Equal(
+ ["orders-service-orders-0", "orders-service-orders-1", "orders-service-orders-2"],
+ queues.Select(queue => queue.QueueName));
+ Assert.All(queues, queue =>
+ {
+ Assert.NotEqual(true, queue.FifoQueue);
+ Assert.Equal(12, queue.ReceiveMessageWaitTimeSeconds);
+ Assert.Equal(45, queue.VisibilityTimeout);
+ });
+
+ var siloEnvironment = await app.GetSiloEnvironmentAsync();
+ var clientEnvironment = await app.GetClientEnvironmentAsync();
+ AssertProviderEnvironment(siloEnvironment, fifoQueue: false);
+ AssertProviderEnvironment(clientEnvironment, fifoQueue: false);
+ AssertWaitsForStack(app.Silo, app.Streaming);
+ AssertWaitsForStack(app.Client, app.Streaming);
+ }
+
+ [Fact]
+ public async Task AddSqsStreaming_FifoTopology_UsesRuntimeQueueNamesAndFifoProperties()
+ {
+ await using var app = await CreateAppAsync(
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ PartitionCount = 4,
+ FifoQueue = true,
+ });
+
+ var expectedNames = new HashRingBasedStreamQueueMapper(
+ new HashRingStreamQueueMapperOptions { TotalQueueCount = 4 },
+ ProviderName)
+ .GetAllQueues()
+ .Select(queue => SqsQueueName.Create(queue.ToString(), fifoQueue: true, ServiceId))
+ .Order()
+ .ToArray();
+ var queues = GetQueues(app.Streaming);
+
+ Assert.Equal(expectedNames, queues.Select(queue => queue.QueueName).Order());
+ Assert.All(queues, queue =>
+ {
+ Assert.Equal(true, queue.FifoQueue);
+ Assert.Equal(true, queue.ContentBasedDeduplication);
+ Assert.Equal("messageGroup", queue.DeduplicationScope);
+ Assert.Equal("perMessageGroupId", queue.FifoThroughputLimit);
+ });
+ }
+
+ [Fact]
+ public async Task WithSqsStreaming_ReturnsOrleansServiceAndActivatesSiloAndClientProviders()
+ {
+ _ = typeof(SqsStreamProviderBuilder);
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig()
+ .WithProfile("integration-profile")
+ .WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster")
+ .WithDevelopmentClustering()
+ .WithSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ PartitionCount = 2,
+ FifoQueue = true,
+ });
+ var silo = builder.AddContainer("silo", "unused").WithReference(orleans);
+ var client = builder.AddContainer("client", "unused").WithReference(orleans.AsClient());
+ await using var application = builder.Build();
+
+ using (new EnvironmentVariableScope(
+ await GetEnvironmentAsync(application.Services, silo.Resource)))
+ {
+ using var siloHost = Host.CreateApplicationBuilder().UseOrleans().Build();
+ AssertActivatedProvider(siloHost.Services, partitionCount: 2, fifoQueue: true);
+ }
+
+ using (new EnvironmentVariableScope(
+ await GetEnvironmentAsync(application.Services, client.Resource)))
+ {
+ using var clientHost = Host.CreateApplicationBuilder().UseOrleansClient().Build();
+ AssertActivatedProvider(clientHost.Services, partitionCount: 2, fifoQueue: true);
+ }
+ }
+
+ [Theory]
+ [InlineData("", 1, null, null, "ServiceId")]
+ [InlineData("service", 0, null, null, "PartitionCount")]
+ [InlineData("service", 1, -1, null, "ReceiveWaitTimeSeconds")]
+ [InlineData("service", 1, 21, null, "ReceiveWaitTimeSeconds")]
+ [InlineData("service", 1, null, 43_201, "VisibilityTimeoutSeconds")]
+ public void AddSqsStreaming_InvalidOptions_ThrowsActionableError(
+ string serviceId,
+ int partitionCount,
+ int? receiveWaitTimeSeconds,
+ int? visibilityTimeoutSeconds,
+ string expectedMessage)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var options = new SqsStreamingOptions
+ {
+ ServiceId = serviceId,
+ PartitionCount = partitionCount,
+ ReceiveWaitTimeSeconds = receiveWaitTimeSeconds,
+ VisibilityTimeoutSeconds = visibilityTimeoutSeconds,
+ };
+
+ var exception = Assert.ThrowsAny(
+ () => orleans.AddSqsStreaming(ProviderName, aws, options));
+
+ Assert.Contains(expectedMessage, exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_MissingAwsRegion_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig();
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId }));
+
+ Assert.Contains("concrete AWS region", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_NullArguments_ThrowWithParameterNames()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var options = new SqsStreamingOptions { ServiceId = ServiceId };
+
+ Assert.Equal(
+ "orleansService",
+ Assert.Throws(
+ () => OrleansSqsStreamingExtensions.AddSqsStreaming(null!, ProviderName, aws, options)).ParamName);
+ Assert.Equal(
+ "awsSdkConfig",
+ Assert.Throws(
+ () => orleans.AddSqsStreaming(ProviderName, null!, options)).ParamName);
+ Assert.Equal(
+ "options",
+ Assert.Throws(
+ () => orleans.AddSqsStreaming(ProviderName, aws, null!)).ParamName);
+ Assert.Equal(
+ "name",
+ Assert.Throws(
+ () => orleans.AddSqsStreaming(" ", aws, options)).ParamName);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_InvalidCacheSize_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ CacheSize = 0,
+ }));
+
+ Assert.Contains("CacheSize", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_DataAdapterKeyMatchingProviderName_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ DataAdapterKey = ProviderName,
+ }));
+
+ Assert.Contains("must differ", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Theory]
+ [InlineData(true)]
+ [InlineData(false)]
+ public void AddSqsStreaming_NullAttributeList_ThrowsWithPropertyName(bool messageAttributes)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var options = messageAttributes
+ ? new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ ReceiveMessageAttributes = null!,
+ }
+ : new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ ReceiveMessageSystemAttributes = null!,
+ };
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(ProviderName, aws, options));
+
+ Assert.Equal(
+ messageAttributes ? "ReceiveMessageAttributes" : "ReceiveMessageSystemAttributes",
+ exception.ParamName);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_InvalidGeneratedQueueName_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ "Orders.With.Invalid.Characters",
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId }));
+
+ Assert.Contains("generated SQS queue name", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_ProviderNameWithDoubleUnderscore_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ "orders__priority",
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId }));
+
+ Assert.Equal("name", exception.ParamName);
+ Assert.Contains("configuration path segment", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_OverlongGeneratedQueueName_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions { ServiceId = new string('s', 75) }));
+
+ Assert.Contains("at most 80 characters", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_NonalphanumericProviderName_UsesStableConstructIds()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var resource = orleans.AddSqsStreaming(
+ "-",
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ PartitionCount = 2,
+ });
+
+ Assert.Equal(
+ [$"{resource.Stack.Resource.Name}-0", $"{resource.Stack.Resource.Name}-1"],
+ resource.Queues.Select(queue => queue.Resource.Name));
+ Assert.Equal(
+ ["orders-service---0", "orders-service---1"],
+ GetQueues(resource).Select(queue => queue.QueueName));
+ }
+
+ [Fact]
+ public void AddSqsStreaming_NormalizedResourceNamesRemainDistinct()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var hyphenated = orleans.AddSqsStreaming(
+ "orders-primary",
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId });
+ var underscored = orleans.AddSqsStreaming(
+ "orders_primary",
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId });
+
+ Assert.Equal("cluster-orders-primary-sqs", hyphenated.Stack.Resource.Name);
+ Assert.StartsWith("cluster-orders-primary-sqs-", underscored.Stack.Resource.Name);
+ Assert.NotEqual(hyphenated.Stack.Resource.Name, underscored.Stack.Resource.Name);
+ Assert.Equal(
+ Enumerable.Range(0, 8).Select(index => $"orders-service-orders-primary-{index}"),
+ GetQueues(hyphenated).Select(queue => queue.QueueName));
+ Assert.Equal(
+ Enumerable.Range(0, 8).Select(index => $"orders-service-orders_primary-{index}"),
+ GetQueues(underscored).Select(queue => queue.QueueName));
+ Assert.Equal(2, orleans.Streaming.Count);
+ }
+
+ [Theory]
+ [InlineData("Orders", "Orders", false)]
+ [InlineData("Orders", "orders", false)]
+ [InlineData("orders_primary", "Orders_primary", false)]
+ [InlineData("orders-primary", "Orders-primary", true)]
+ public void AddSqsStreaming_ConflictingProviderIdentity_PreservesResourcesAndConfiguration(
+ string existingName,
+ string conflictingName,
+ bool fifoQueue)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var existing = orleans.AddSqsStreaming(
+ existingName, aws, new SqsStreamingOptions { ServiceId = ServiceId });
+ var resources = builder.Resources.ToArray();
+ var constructs = existing.Stack.Resource.Stack.Node.FindAll().ToArray();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ conflictingName,
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId, PartitionCount = 2, FifoQueue = fifoQueue }));
+
+ Assert.Contains("stream provider", exception.Message, StringComparison.Ordinal);
+ Assert.Contains(existingName, exception.Message, StringComparison.Ordinal);
+ Assert.Contains(conflictingName, exception.Message, StringComparison.Ordinal);
+ Assert.Equal(resources, builder.Resources);
+ Assert.Equal(constructs, existing.Stack.Resource.Stack.Node.FindAll());
+ var provider = Assert.Single(orleans.Streaming);
+ Assert.Equal(existingName, provider.Key);
+ Assert.Same(existing, provider.Value);
+ }
+
+ [Theory]
+ [InlineData("Orders", "orders")]
+ [InlineData("orders_primary", "Orders_primary")]
+ public async Task AddSqsStreaming_ExistingOtherProvider_PreservesResourcesAndServiceId(
+ string existingName,
+ string conflictingName)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering().WithMemoryStreaming(existingName);
+ var existing = orleans.Streaming[existingName];
+ var silo = builder.AddContainer("silo", "unused").WithReference(orleans);
+ await using var application = builder.Build();
+ var environment = await GetEnvironmentAsync(application.Services, silo.Resource);
+ var resources = builder.Resources.ToArray();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ conflictingName, aws, new SqsStreamingOptions { ServiceId = ServiceId }));
+
+ Assert.Contains("stream provider", exception.Message, StringComparison.Ordinal);
+ Assert.Equal(resources, builder.Resources);
+ Assert.Same(existing, Assert.Single(orleans.Streaming).Value);
+ Assert.Equal(
+ environment.OrderBy(pair => pair.Key),
+ (await GetEnvironmentAsync(application.Services, silo.Resource)).OrderBy(pair => pair.Key));
+ }
+
+ [Theory]
+ [InlineData("service", "Orders", "service", "orders", false)]
+ [InlineData("service", "orders_primary", "service", "Orders_primary", true)]
+ [InlineData("service-orders", "primary", "service", "orders-primary", false)]
+ [InlineData("service-orders", "primary", "service", "orders-primary", true)]
+ public async Task AddSqsStreaming_ConflictingPhysicalTopology_PreservesResourcesAndConfiguration(
+ string existingServiceId,
+ string existingName,
+ string conflictingServiceId,
+ string conflictingName,
+ bool fifoQueue)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var first = builder.AddOrleans("first").WithDevelopmentClustering();
+ var existing = first.AddSqsStreaming(
+ existingName,
+ aws,
+ new SqsStreamingOptions { ServiceId = existingServiceId, PartitionCount = 1, FifoQueue = fifoQueue });
+ var second = builder.AddOrleans("second").WithDevelopmentClustering();
+ var silo = builder.AddContainer("silo", "unused").WithReference(second);
+ await using var application = builder.Build();
+ var environment = await GetEnvironmentAsync(application.Services, silo.Resource);
+ var resources = builder.Resources.ToArray();
+ var constructs = existing.Stack.Resource.Stack.Node.FindAll().ToArray();
+
+ var exception = Assert.Throws(
+ () => second.AddSqsStreaming(
+ conflictingName,
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = conflictingServiceId,
+ PartitionCount = 2,
+ FifoQueue = fifoQueue,
+ }));
+
+ var queueName = Assert.IsType(Assert.Single(GetQueues(existing)).QueueName);
+ Assert.Contains(queueName, exception.Message, StringComparison.Ordinal);
+ Assert.Contains("us-east-1", exception.Message, StringComparison.Ordinal);
+ Assert.Equal(resources, builder.Resources);
+ Assert.Equal(constructs, existing.Stack.Resource.Stack.Node.FindAll());
+ Assert.Same(existing, Assert.Single(first.Streaming).Value);
+ Assert.Empty(second.Streaming);
+ Assert.Equal(
+ environment.OrderBy(pair => pair.Key),
+ (await GetEnvironmentAsync(application.Services, silo.Resource)).OrderBy(pair => pair.Key));
+ }
+
+ [Fact]
+ public void AddSqsStreaming_SameQueueNamesInDifferentRegions_CreatesSeparateTopologies()
+ {
+ var builder = CreateBuilder();
+ var east = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var west = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USWest2);
+ var options = new SqsStreamingOptions { ServiceId = ServiceId, PartitionCount = 1 };
+ var first = builder.AddOrleans("first").AddSqsStreaming(ProviderName, east, options);
+ var second = builder.AddOrleans("second").AddSqsStreaming(ProviderName, west, options);
+
+ Assert.Equal(Assert.Single(GetQueues(first)).QueueName, Assert.Single(GetQueues(second)).QueueName);
+ Assert.NotEqual(first.Stack.Resource.Name, second.Stack.Resource.Name);
+ Assert.Equal(2, builder.Resources.OfType().Count());
+ }
+
+ [Fact]
+ public void AddSqsStreaming_SameProviderInDifferentServices_CreatesSeparateTopologies()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var first = builder.AddOrleans("first").AddSqsStreaming(
+ ProviderName, aws, new SqsStreamingOptions { ServiceId = "first", PartitionCount = 1 });
+ var second = builder.AddOrleans("second").AddSqsStreaming(
+ ProviderName, aws, new SqsStreamingOptions { ServiceId = "second", PartitionCount = 1 });
+
+ Assert.Equal("first-orders-0", Assert.Single(GetQueues(first)).QueueName);
+ Assert.Equal("second-orders-0", Assert.Single(GetQueues(second)).QueueName);
+ Assert.Equal("first-orders-sqs-0", Assert.Single(first.Queues).Resource.Name);
+ Assert.Equal("second-orders-sqs-0", Assert.Single(second.Queues).Resource.Name);
+ }
+
+ [Theory]
+ [InlineData("cluster-orders-sqs")]
+ [InlineData("cluster-orders-sqs-1")]
+ public void AddSqsStreaming_ConflictingResourceName_PreservesResourcesAndConfiguration(string resourceName)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ builder.AddContainer(resourceName, "unused");
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var resources = builder.Resources.ToArray();
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ ProviderName, aws, new SqsStreamingOptions { ServiceId = ServiceId, PartitionCount = 2 }));
+
+ Assert.Contains(resourceName, exception.Message, StringComparison.Ordinal);
+ Assert.Equal(resources, builder.Resources);
+ Assert.Empty(orleans.Streaming);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_ConflictingServiceIdsAcrossProviders_ThrowsActionableError()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ orleans.AddSqsStreaming(
+ "Orders",
+ aws,
+ new SqsStreamingOptions { ServiceId = "orders-service" });
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ "Billing",
+ aws,
+ new SqsStreamingOptions { ServiceId = "billing-service" }));
+
+ Assert.Contains("conflicts with Orleans ServiceId", exception.Message, StringComparison.Ordinal);
+ Assert.Contains("orders-service", exception.Message, StringComparison.Ordinal);
+ Assert.Contains("billing-service", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_ParameterizedServiceId_ThrowsBeforeReplacingIt()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var serviceId = builder.AddParameter("service-id");
+ var orleans = builder.AddOrleans("cluster")
+ .WithDevelopmentClustering()
+ .WithServiceId(serviceId);
+
+ var exception = Assert.Throws(
+ () => orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId }));
+
+ Assert.Contains("concrete string", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public async Task WithServiceId_AfterAddSqsStreaming_ConflictingValueFailsBeforeDeployment()
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions { ServiceId = ServiceId });
+ orleans.WithServiceId("different-service");
+ var silo = builder.AddContainer("silo", "unused").WithReference(orleans);
+ await using var application = builder.Build();
+
+ var exception = await Assert.ThrowsAsync(
+ () => GetEnvironmentAsync(application.Services, silo.Resource));
+
+ Assert.Contains(ServiceId, exception.Message, StringComparison.Ordinal);
+ Assert.Contains("different-service", exception.Message, StringComparison.Ordinal);
+ }
+
+ [Fact]
+ public void AddSqsStreaming_DefensivelyCopiesAttributeLists()
+ {
+ var messageAttributes = new[] { "TraceId" };
+ var systemAttributes = new[] { "SentTimestamp" };
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig().WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var resource = orleans.AddSqsStreaming(
+ ProviderName,
+ aws,
+ new SqsStreamingOptions
+ {
+ ServiceId = ServiceId,
+ ReceiveMessageAttributes = messageAttributes,
+ ReceiveMessageSystemAttributes = systemAttributes,
+ });
+
+ messageAttributes[0] = "Changed";
+ systemAttributes[0] = "Changed";
+
+ Assert.Equal("TraceId", Assert.Single(resource.Options.ReceiveMessageAttributes));
+ Assert.Equal("SentTimestamp", Assert.Single(resource.Options.ReceiveMessageSystemAttributes));
+ }
+
+ [Fact]
+ public void EnvironmentVariableScope_ClearsAndRestoresInheritedAwsConfiguration()
+ {
+ var originalDefaultRegion = Environment.GetEnvironmentVariable("AWS_DEFAULT_REGION");
+ var originalStructuredRegion = Environment.GetEnvironmentVariable("AWS__Region");
+ try
+ {
+ Environment.SetEnvironmentVariable("AWS_DEFAULT_REGION", "us-west-2");
+ Environment.SetEnvironmentVariable("AWS__Region", "us-west-2");
+
+ using (new EnvironmentVariableScope(
+ new Dictionary { ["AWS_REGION"] = "us-east-1" }))
+ {
+ Assert.Equal("us-east-1", Environment.GetEnvironmentVariable("AWS_REGION"));
+ Assert.Null(Environment.GetEnvironmentVariable("AWS_DEFAULT_REGION"));
+ Assert.Null(Environment.GetEnvironmentVariable("AWS__Region"));
+ }
+
+ Assert.Equal("us-west-2", Environment.GetEnvironmentVariable("AWS_DEFAULT_REGION"));
+ Assert.Equal("us-west-2", Environment.GetEnvironmentVariable("AWS__Region"));
+ }
+ finally
+ {
+ Environment.SetEnvironmentVariable("AWS_DEFAULT_REGION", originalDefaultRegion);
+ Environment.SetEnvironmentVariable("AWS__Region", originalStructuredRegion);
+ }
+ }
+
+ private static Task CreateAppAsync(SqsStreamingOptions options)
+ {
+ var builder = CreateBuilder();
+ var aws = builder.AddAWSSDKConfig()
+ .WithProfile("integration-profile")
+ .WithRegion(RegionEndpoint.USEast1);
+ var orleans = builder.AddOrleans("cluster").WithDevelopmentClustering();
+ var streaming = orleans.AddSqsStreaming(ProviderName, aws, options);
+ var silo = builder.AddContainer("silo", "unused").WithReference(orleans);
+ var client = builder.AddContainer("client", "unused").WithReference(orleans.AsClient());
+ var application = builder.Build();
+ return Task.FromResult(new TestApp(application, streaming, silo.Resource, client.Resource));
+ }
+
+ private static IDistributedApplicationBuilder CreateBuilder()
+ => DistributedApplication.CreateBuilder(
+ new DistributedApplicationOptions
+ {
+ Args = [],
+ DisableDashboard = true,
+ });
+
+ private static CfnQueue[] GetQueues(SqsStreamingResource streaming)
+ => streaming.Queues
+ .Select(queue => Assert.IsType(queue.Resource.Construct.Node.DefaultChild))
+ .OrderBy(queue => queue.QueueName)
+ .ToArray();
+
+ private static async Task> GetEnvironmentAsync(
+ IServiceProvider services,
+ IResource resource)
+ {
+ Assert.IsAssignableFrom(resource);
+ var executionContext = new DistributedApplicationExecutionContext(
+ new DistributedApplicationExecutionContextOptions(DistributedApplicationOperation.Run)
+ {
+ ServiceProvider = services,
+ });
+ var values = new Dictionary();
+ var callbackContext = new EnvironmentCallbackContext(executionContext, resource, values);
+ foreach (var annotation in resource.Annotations.OfType())
+ {
+ await annotation.Callback(callbackContext).WaitAsync(TimeSpan.FromSeconds(10));
+ }
+
+ var valueContext = new ValueProviderContext
+ {
+ Caller = resource,
+ ExecutionContext = executionContext,
+ Network = KnownNetworkIdentifiers.LocalhostNetwork,
+ };
+ var result = new Dictionary(StringComparer.Ordinal);
+ foreach (var (key, value) in values)
+ {
+ if (!IsRelevantEnvironmentVariable(key))
+ {
+ continue;
+ }
+
+ try
+ {
+ result[key] = value switch
+ {
+ IValueProvider provider => await provider
+ .GetValueAsync(valueContext)
+ .AsTask()
+ .WaitAsync(TimeSpan.FromSeconds(10)),
+ _ => value.ToString(),
+ };
+ }
+ catch (TimeoutException exception)
+ {
+ throw new TimeoutException(
+ $"Timed out resolving environment variable '{key}' for resource '{resource.Name}'.",
+ exception);
+ }
+ }
+ return result;
+ }
+
+ private static bool IsRelevantEnvironmentVariable(string name)
+ => name.StartsWith("Orleans__Streaming__", StringComparison.Ordinal)
+ || name.StartsWith("Orleans__Clustering__", StringComparison.Ordinal)
+ || name is "Orleans__ClusterId" or "Orleans__ServiceId"
+ || name.StartsWith("AWS_", StringComparison.Ordinal);
+
+ private static void AssertProviderEnvironment(
+ IReadOnlyDictionary environment,
+ bool fifoQueue)
+ {
+ const string prefix = "Orleans__Streaming__Orders__";
+ Assert.Equal(ServiceId, environment["Orleans__ServiceId"]);
+ Assert.Equal("integration-profile", environment["AWS_PROFILE"]);
+ Assert.Equal("us-east-1", environment["AWS_REGION"]);
+ Assert.Equal("SQS", environment[$"{prefix}ProviderType"]);
+ Assert.Equal("us-east-1", environment[$"{prefix}Region"]);
+ Assert.Equal("3", environment[$"{prefix}PartitionCount"]);
+ Assert.Equal(fifoQueue.ToString(), environment[$"{prefix}FifoQueue"]);
+ Assert.Equal("12", environment[$"{prefix}ReceiveWaitTimeSeconds"]);
+ Assert.Equal("45", environment[$"{prefix}VisibilityTimeoutSeconds"]);
+ Assert.Equal("2048", environment[$"{prefix}CacheSize"]);
+ Assert.Equal("orders-adapter", environment[$"{prefix}DataAdapterKey"]);
+ Assert.Equal("TraceId", environment[$"{prefix}ReceiveMessageAttributes__0"]);
+ Assert.Equal("Tenant", environment[$"{prefix}ReceiveMessageAttributes__1"]);
+ Assert.Equal("SentTimestamp", environment[$"{prefix}ReceiveMessageSystemAttributes__0"]);
+ }
+
+ private static void AssertWaitsForStack(IResource resource, SqsStreamingResource streaming)
+ => Assert.Contains(
+ resource.Annotations.OfType(),
+ annotation => ReferenceEquals(annotation.Resource, streaming.Stack.Resource));
+
+ private static void AssertActivatedProvider(
+ IServiceProvider services,
+ int partitionCount,
+ bool fifoQueue)
+ {
+ var sqsOptions = services.GetRequiredService>().Get(ProviderName);
+ var partitionOptions = services
+ .GetRequiredService>()
+ .Get(ProviderName);
+
+ Assert.Equal("Service=us-east-1", sqsOptions.ConnectionString);
+ Assert.Equal(fifoQueue, sqsOptions.FifoQueue);
+ Assert.Equal(partitionCount, partitionOptions.TotalQueueCount);
+ }
+
+ private sealed class TestApp(
+ DistributedApplication application,
+ SqsStreamingResource streaming,
+ IResource silo,
+ IResource client) : IAsyncDisposable
+ {
+ public SqsStreamingResource Streaming { get; } = streaming;
+
+ public IResource Silo { get; } = silo;
+
+ public IResource Client { get; } = client;
+
+ public Task> GetSiloEnvironmentAsync()
+ => GetEnvironmentAsync(application.Services, Silo);
+
+ public Task> GetClientEnvironmentAsync()
+ => GetEnvironmentAsync(application.Services, Client);
+
+ public ValueTask DisposeAsync() => application.DisposeAsync();
+ }
+
+ private sealed class EnvironmentVariableScope : IDisposable
+ {
+ private readonly Dictionary _previousValues = new(StringComparer.OrdinalIgnoreCase);
+
+ public EnvironmentVariableScope(IReadOnlyDictionary values)
+ {
+ foreach (var key in new[]
+ {
+ "AWS_REGION",
+ "AWS_DEFAULT_REGION",
+ "AWS_PROFILE",
+ "AWS__Region",
+ "AWS__Profile",
+ })
+ {
+ SaveAndSet(key, null);
+ }
+
+ foreach (var (key, value) in values)
+ {
+ SaveAndSet(key, value);
+ }
+ }
+
+ public void Dispose()
+ {
+ foreach (var (key, value) in _previousValues)
+ {
+ Environment.SetEnvironmentVariable(key, value);
+ }
+ }
+
+ private void SaveAndSet(string key, string? value)
+ {
+ _previousValues.TryAdd(key, Environment.GetEnvironmentVariable(key));
+ Environment.SetEnvironmentVariable(key, value);
+ }
+ }
+}
diff --git a/test/Transactions/Orleans.Transactions.DynamoDB.Test/Orleans.Transactions.DynamoDB.Test.csproj b/test/Transactions/Orleans.Transactions.DynamoDB.Test/Orleans.Transactions.DynamoDB.Test.csproj
index fb691258a44..0e604f2f1fc 100644
--- a/test/Transactions/Orleans.Transactions.DynamoDB.Test/Orleans.Transactions.DynamoDB.Test.csproj
+++ b/test/Transactions/Orleans.Transactions.DynamoDB.Test/Orleans.Transactions.DynamoDB.Test.csproj
@@ -4,12 +4,13 @@
Orleans.Transactions.DynamoDB.Tests
$(TestTargetFrameworks)
true
+ $(DefineConstants);TRANSACTIONS_DYNAMODB_TESTS
-
+