Kanject.Core.Queue.Provider.AwsSqs 3.7.7

Prefix Reserved
There is a newer version of this package available.
See the version list below for details.
dotnet add package Kanject.Core.Queue.Provider.AwsSqs --version 3.7.7
                    
NuGet\Install-Package Kanject.Core.Queue.Provider.AwsSqs -Version 3.7.7
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs" Version="3.7.7" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Kanject.Core.Queue.Provider.AwsSqs" Version="3.7.7" />
                    
Directory.Packages.props
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Kanject.Core.Queue.Provider.AwsSqs --version 3.7.7
                    
#r "nuget: Kanject.Core.Queue.Provider.AwsSqs, 3.7.7"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Kanject.Core.Queue.Provider.AwsSqs@3.7.7
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Kanject.Core.Queue.Provider.AwsSqs&version=3.7.7
                    
Install as a Cake Addin
#tool nuget:?package=Kanject.Core.Queue.Provider.AwsSqs&version=3.7.7
                    
Install as a Cake Tool

Kanject.Core.Queue.Provider.AwsSqs

The Amazon SQS implementation of Kanject.Core.Queue.Abstractions. It covers publishing (single, batch and routed messages), consuming (AWS Lambda SQS events or long-polling hosted consumers in any .NET host), partial-batch failure reporting, and optional queue and dead-letter-queue creation.

Use it in services that exchange work through SQS. Declare messages and consumers with attributes; the Kanject.Core.Queue.Provider.AwsSqs.Annotations generator writes the registration, send and Lambda-dispatch code for them.

Installation

<ItemGroup>
  <PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs" />
  
  <PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs.Annotations" PrivateAssets="all" />
</ItemGroup>

or

dotnet add package Kanject.Core.Queue.Provider.AwsSqs
dotnet add package Kanject.Core.Queue.Provider.AwsSqs.Annotations

Targets .NET 8, .NET 9 and .NET 10. This package depends on the Annotations package, so the generators normally reach your project through it. Add the direct Annotations reference whenever you use the attributes anyway, so the generator version is explicit.

Quick start

1. Declare the message and the consumer

using Amazon.Lambda.SQSEvents;
using Kanject.Core.Queue.Abstractions.Models;
using Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes;

namespace Shop.Orders;

[QueueMessage("orders")]                       // SQS queue "orders_queue"
public sealed class OrderPlaced
{
    public string OrderId { get; set; } = string.Empty;
    public decimal Total { get; set; }
}

[QueueConsumer(Message = typeof(OrderPlaced))] // queue name comes from [QueueMessage]
[QueueConsumerDependency(typeof(IFulfilmentService))]
public partial class OrderPlacedConsumer
{
    protected override async Task ConsumeAsync(
        List<MessageContext<OrderPlaced>> messages, CancellationToken cancellationToken)
    {
        foreach (var context in messages)
        {
            try
            {
                if (context.Message is null)            // body could not be deserialized
                    throw new InvalidOperationException("Unreadable message.");

                await FulfilmentService.ShipAsync(context.Message, cancellationToken);
                await AcknowledgeAsync(context);
            }
            catch (Exception)
            {
                Response.BatchItemFailures.Add(new SQSBatchResponse.BatchItemFailure
                {
                    ItemIdentifier = context.MessageId
                });
            }
        }
    }
}

The generator makes OrderPlacedConsumer derive from AbstractQueueConsumer<OrderPlaced>. It also adds a constructor that takes IQueueRegistry plus each declared dependency, and exposes each dependency as a property (FulfilmentService here). Finally, it emits AddOrderPlacedConsumer(), ProcessOrderPlacedEventAsync(...), AddOrdersQueue() and EnqueueOrderPlacedAsync(...).

2. Consume in AWS Lambda

using Amazon.Lambda.Core;
using Amazon.Lambda.SQSEvents;
using Kanject.Core.Queue.Provider.AwsSqs.Extensions;
using Microsoft.Extensions.DependencyInjection;
using Shop.Orders;

public sealed class Function
{
    private static readonly IServiceProvider Services = BuildServices();

    private static IServiceProvider BuildServices()
    {
        var services = new ServiceCollection();
        services.AddSingleton<IFulfilmentService, FulfilmentService>();

        services
            .AddAwsSqsGlobalQueueConfiguration(options =>
            {
                // No keys: the SQS client uses the default AWS credential chain, which in Lambda
                // is the function's execution role. Set AWSAccessKey/AWSSecretKey only for static keys.
                options.AWSRegion = "eu-west-1";
            })
            .AddOrderPlacedConsumer();

        return services.BuildServiceProvider();
    }

    // SQS event source mapping with ReportBatchItemFailures enabled.
    public Task<SQSBatchResponse> HandleAsync(SQSEvent sqsEvent, ILambdaContext context)
        => Services.ProcessOrderPlacedEventAsync(sqsEvent);
}

ProcessOrderPlacedEventAsync deserializes each record into a MessageContext<OrderPlaced>, passes the batch to ConsumeAsync, and returns the consumer's Response. Only the ids you added to BatchItemFailures are retried. Every invocation starts with an empty Response, so a failure recorded in one invocation is never returned again by the next one on a warm container. It also takes an optional CancellationToken, which it forwards to ConsumeAsync.

3. Send

using Kanject.Core.Queue.Abstractions.Interfaces;
using Shop.Orders;

// Registration: services.AddAwsSqsGlobalQueueConfiguration(...).AddOrdersQueue();

public sealed class CheckoutService(IQueueManagerService queues)
{
    public async Task CompleteAsync(OrderPlaced order)
    {
        if (!await queues.EnqueueOrderPlacedAsync(order))
            throw new InvalidOperationException($"Order {order.OrderId} was not queued.");
    }
}

The message is serialized with System.Text.Json. Other generated overloads take metadata (sent as SQS message attributes), an IList<OrderPlaced> batch, or (message, metadata) tuples.

Consuming in a long-running host

Every registered consumer is also an IHostedService. Polling is opt-in: set WatchQueue and each consumer long-polls its queue when the host starts. With WatchQueue off (the default) the hosted service does nothing, so a host that starts its hosted services inside Lambda never polls alongside the event source mapping.

using Kanject.Core.Queue.Provider.AwsSqs.Abstractions.Extensions;
using Kanject.Core.Queue.Provider.AwsSqs.Extensions;
using Shop.Orders;

var builder = WebApplication.CreateBuilder(args);

builder.Services
    .AddAwsSqsGlobalQueueConfiguration(options =>
    {
        // credentials and region as above
        options.WatchQueue = true;       // hosted consumers poll their queue
        options.WatcherSleepTime = 5;    // seconds between polls
        options.UseDeadLetterQueue = true;
    })
    .AddOrdersQueue()
    .AddOrderPlacedConsumer();

var app = builder.Build();

// Resolves every registered queue in the background, creating missing ones
// when CreateQueueIfNotExist is true (the default).
app.UseAwsSqsQueueProvider();

app.Run();

For a generic host, call host.Services.UseAwsSqsQueueProvider(). Each poll receives up to MaximumReceiveMessageCount messages and passes them to ConsumeAsync with a token that's cancelled when the host stops the consumer. Polling failures back off exponentially, up to 30 seconds. In this mode, AcknowledgeAsync is what deletes a processed message from the queue.

Routed queues

One queue can carry several message kinds, dispatched by a route tag:

[RouteQueueConsumer(QueueName = "order-events", Message = typeof(OrderEvent))]
[QueueConsumerRoute("order.created")]
[QueueConsumerRoute("order.updated")]
public partial class OrderChangedHandler
{
    protected override Task ConsumeAsync(
        List<MessageContext<OrderEvent>> messages, CancellationToken cancellationToken)
    {
        /* handle, acknowledge, or add to Response.BatchItemFailures */
        return Task.CompletedTask;
    }
}
// Registration: maps every [QueueConsumerRoute] declared for "order-events"
services.AddOrderEventsQueueQueueRouter();

// Lambda handler
Task<SQSBatchResponse> HandleAsync(SQSEvent sqsEvent, ILambdaContext context)
    => Services.RouteSqsEventToOrderEventsQueueAsync(sqsEvent);

// Publisher: the route travels as the "QueueRoute" message attribute
await queues["order-events"].EnqueueAsync("order.created", orderEvent);

Route consumers are registered as scoped services and resolved in a fresh scope per event. With WatchQueue = true, a hosted SqsQueueConsumerRouter polls and dispatches the queue instead. Without the generator, the same wiring is services.AddQueueConsumerRouter(q => q.UseQueue("order-events").MapConsumerRoute<OrderChangedHandler>("order.created")). MapDefaultConsumerRoute<T>() catches routes that have no mapped consumer, and records that carry no QueueRoute attribute at all. A record that reaches no consumer — no mapped route and no default route, or a mapped consumer that isn't registered — is reported in BatchItemFailures, so SQS retries it and eventually moves it to the dead-letter queue instead of deleting it.

Configuration

AddAwsSqsGlobalQueueConfiguration sets the defaults for every queue and consumer registered after it. Each registration copies the global settings at the moment it's made, and you can override them per queue through the generated Add…(options => …) overloads.

AwsSqsQueueConfiguration property Default Effect
AWSAccessKey, AWSSecretKey empty Static credentials for the SQS client. Leave either empty to use the default AWS credential chain (see below)
AWSRegion empty Region for the SQS client. Empty uses the SDK's region resolution (AWS_REGION in Lambda)
EnforceQueueOrdering false Use a FIFO queue (.fifo name). Every send carries a MessageGroupId (see below)
CreateQueueIfNotExist true Create the queue during UseAwsSqsQueueProvider() if it doesn't exist
UseDeadLetterQueue false Also create <queue>_dlq and attach a redrive policy
MaximumMessageRetry 10 Redrive maxReceiveCount before a message moves to the DLQ
WatchQueue false Hosted consumers and the hosted router poll the queue only when this is true
WatcherSleepTime 10 Seconds between polls
MaximumReceiveMessageCount 10 Messages per receive call
ReceiveMessageWaitTimeInSeconds 20 Long-poll wait time
VisibilityWaitTimeInSeconds 20 Visibility timeout
DelaySeconds 0 Delivery delay for every message
MessageRetentionPeriod / DlqMessageRetentionPeriod 4 / 5 Retention, in days

Queue attributes such as delay, retention, wait time and visibility timeout are applied only when the provider creates a queue. Existing queues are left unchanged.

Things to know

  • Credentials: when both AWSAccessKey and AWSSecretKey are set, the queue's SQS client signs with them. Otherwise it uses the default AWS credential chain: the Lambda or ECS task role, an EC2 instance profile, the AWS_* environment variables, or the shared profile. The provider calls sqs:GetQueueUrl, sqs:GetQueueAttributes, sqs:SendMessage, sqs:ReceiveMessage and sqs:DeleteMessage. It also calls sqs:CreateQueue when it's allowed to create queues.
  • Queue names: orders becomes the SQS queue orders_queue (see QueueNameHelper in Kanject.Core.Queue.Abstractions). The queue must exist, either provisioned by your infrastructure or created at startup with UseAwsSqsQueueProvider().
  • Send results: enqueue calls return false instead of throwing when SQS rejects the request or the queue can't be resolved; the exception is written to the console. Always check the result. Using a queue key that isn't registered throws KeyNotFoundException from the IQueueManagerService indexer.
  • Batches: SQS accepts at most 10 entries per SendMessageBatch request, so the batch overloads send larger batches in order, 10 messages per request, one request at a time. A batch returns true only if SQS accepted every message. It returns false as soon as a request is rejected or lists any entry under Failed (SQS reports those with HTTP 200), and the rest of the batch isn't sent. The failed entries are written to the console. Earlier requests may already have been delivered, so resending the whole batch can deliver those messages twice. Metadata and routes passed to a batch overload are sent with every message in it.
  • FIFO queues: with EnforceQueueOrdering, every send carries a MessageGroupId, which SQS requires on a FIFO queue. By default all messages share one group named after the queue, which gives one strict order. Put a MessageGroupId entry in the metadata (AwsQueueConstants.MessageGroupIdMetadataKey) to order per group instead. Queues the provider creates have content-based deduplication enabled; a FIFO queue created elsewhere needs it too, because the provider doesn't send a deduplication id.
  • Unreadable messages: in the Lambda path, a record whose body can't be deserialized still reaches ConsumeAsync, with Message left at its default value. Guard for it and report the failure. The polling path logs such messages and skips them.
  • Lifetimes: typed consumers ([QueueConsumer], AddQueueConsumer<T>) are singletons. Any service they depend on is captured for the lifetime of the app, so use IServiceScopeFactory for scoped work. Response is per invocation: inside ConsumeAsync it is always the current batch's response, even if invocations overlap; read after the call returns, it is the most recent invocation's.
  • Missing consumer: if a consumer isn't registered, ProcessSqsEventWithQueueConsumerAsync<T>, ProcessSqsEventWithRouteQueueConsumerAsync<T> and the generated Process…EventAsync methods throw InvalidOperationException. The invocation fails, so SQS retries the whole batch and eventually moves it to the dead-letter queue, whether or not the event source mapping has ReportBatchItemFailures enabled.
  • Unrouted records: a record without a QueueRoute attribute goes to the default route. A record that no consumer receives (no mapped route and no default route, or a mapped consumer that isn't registered) is reported as a failure, not dropped.

Public surface at a glance

Type or member Purpose
AddAwsSqsGlobalQueueConfiguration(Action<AwsSqsQueueConfiguration>) Global defaults
CreateQueue(...), AddSqsQueue(name, ...) Register a queue for publishing
AddQueueConsumer<T>(...) Register a consumer as a singleton and hosted service
AddQueueConsumerRouter(...) with UseQueue, MapConsumerRoute<T>, MapDefaultConsumerRoute<T> Register a routed queue
RouteIncomingSqsQueueEventAsync(sqsEvent, queue) Dispatch a Lambda SQSEvent to route consumers without the generator. Return its SQSBatchResponse from the handler. It replaces the obsolete SqsQueueConsumerRouter.ProcessSqsEventAsync, which can't return a response and throws when any record fails
SQSEvent.CategorizeRecordsByFilter(attributeName, allowedValues) Split an event into buckets by a JSON body property
AbstractQueueConsumer<T>, AbstractRawQueueConsumer Base classes for typed and raw (Amazon.SQS.Model.Message) consumers
AbstractRouteQueueConsumer<T>, AbstractRouteRawQueueConsumer Base classes for route consumers
SqsQueueManager, SqsQueueManagerService, QueueRegistry, SqsQueueConsumerRouter The implementations behind the contracts

The Lambda helpers (ProcessSqsEventWithQueueConsumerAsync<T>, ProcessSqsEventWithRouteQueueConsumerAsync<T>), UseAwsSqsQueueProvider(), and AcknowledgeOrFailAsync for [Parallel] outcomes ship in Kanject.Core.Queue.Provider.AwsSqs.Abstractions.

Package Role Availability
Kanject.Core.Queue.Abstractions Provider-neutral queue contracts nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Abstractions SQS contracts, AwsSqsQueueConfiguration, Lambda handler and partial-batch helpers nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Annotations Source generator and analyzers nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes The queue attributes nuget.org
Kanject.Core Base library, including the [Parallel] feature nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Extensions.EventBridge Scheduled and recurring delivery through EventBridge Scheduler Commercial license (not on nuget.org)
Kanject.Core.Queue.Provider.AwsSqs.Extensions.EventBridge.Annotations Generator for typed schedule helpers Commercial license (not on nuget.org)

License

Licensed under the Kanject Code Libraries License Agreement (KCLLA); the full text ships in this package as LICENSE.md. Organizations whose trailing-twelve-month gross revenue and total funding raised are each below US$250,000 may use it at no cost under the Free Tier. At or above either threshold a commercial license is required — contact commercial@kanjectbusiness.com.

Product Compatible and additional computed target framework versions.
.NET net8.0 is compatible.  net8.0-android was computed.  net8.0-browser was computed.  net8.0-ios was computed.  net8.0-maccatalyst was computed.  net8.0-macos was computed.  net8.0-tvos was computed.  net8.0-windows was computed.  net9.0 is compatible.  net9.0-android was computed.  net9.0-browser was computed.  net9.0-ios was computed.  net9.0-maccatalyst was computed.  net9.0-macos was computed.  net9.0-tvos was computed.  net9.0-windows was computed.  net10.0 is compatible.  net10.0-android was computed.  net10.0-browser was computed.  net10.0-ios was computed.  net10.0-maccatalyst was computed.  net10.0-macos was computed.  net10.0-tvos was computed.  net10.0-windows was computed. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages

This package is not used by any NuGet packages.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
3.9.0 39 10/2/2026
3.8.2 39 10/1/2026
3.8.1 107 9/27/2026
3.8.0 83 9/27/2026
3.7.7 90 9/26/2026
3.7.6 119 9/7/2026
3.7.5 102 8/27/2026
3.7.4 107 8/22/2026
3.7.3 119 8/10/2026
3.7.2 114 8/9/2026
3.7.1 115 8/5/2026
3.7.0 112 8/5/2026
3.6.0 115 8/3/2026
3.5.9 129 7/30/2026
3.5.8 122 7/18/2026
3.5.7 141 7/13/2026
3.5.6 121 7/11/2026
3.5.5 129 7/11/2026
3.5.4 129 7/9/2026
3.5.3 125 7/9/2026