Kanject.Core.Queue.Provider.AwsSqs.Abstractions 3.9.0

Prefix Reserved
dotnet add package Kanject.Core.Queue.Provider.AwsSqs.Abstractions --version 3.9.0
                    
NuGet\Install-Package Kanject.Core.Queue.Provider.AwsSqs.Abstractions -Version 3.9.0
                    
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.Abstractions" Version="3.9.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Kanject.Core.Queue.Provider.AwsSqs.Abstractions" Version="3.9.0" />
                    
Directory.Packages.props
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs.Abstractions" />
                    
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.Abstractions --version 3.9.0
                    
#r "nuget: Kanject.Core.Queue.Provider.AwsSqs.Abstractions, 3.9.0"
                    
#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.Abstractions@3.9.0
                    
#: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.Abstractions&version=3.9.0
                    
Install as a Cake Addin
#tool nuget:?package=Kanject.Core.Queue.Provider.AwsSqs.Abstractions&version=3.9.0
                    
Install as a Cake Tool

Kanject.Core.Queue.Provider.AwsSqs.Abstractions

The SQS-specific contracts and helpers shared by Kanject.Core.Queue.Provider.AwsSqs, its generated code, and your consumers. The package builds on Kanject.Core.Queue.Abstractions and adds:

  • AwsSqsQueueConfiguration and the queue registry contracts.
  • Consumer interfaces that accept Lambda SQSEvents and expose an SQSBatchResponse.
  • IServiceProvider / IApplicationBuilder extensions (SqsHandlerExtensions) that call those consumers from a Lambda SQS handler.
  • A bridge from the [Parallel] feature's per-item outcomes into the Lambda partial-batch response (ParallelItemResultSqsExtensions.AcknowledgeOrFailAsync).

Most applications get this package transitively from Kanject.Core.Queue.Provider.AwsSqs. Reference it directly from a library that needs the SQS contracts without the implementation.

Installation

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

Targets .NET 8, .NET 9 and .NET 10. Depends on Kanject.Core.Queue.Abstractions, AWSSDK.SQS and Amazon.Lambda.SQSEvents.

What's in the package

Type Purpose
AwsSqsQueueConfiguration Per-queue SQS settings: credentials, region, dead-letter queue, polling, retention. Extends QueueConfiguration
IAwsSqsQueueManager One SQS queue: send, receive, acknowledge (typed or raw Message), Initialize(), CreateAsync(), route grouping of an SQSEvent
IAwsQueueManagerService IQueueManagerService plus GetAwsQueueManager(queueName)
IQueueRegistry / IQueueRegistryBuilder The registered queue managers. They're collected while services are registered and frozen into a singleton registry on first resolve
IAwsQueueConsumer Typed or raw consumer: InjectSqsEvent(SQSEvent[, CancellationToken]) and the Response batch
IAwsRouteQueueConsumer Route consumer: InjectSqsEventAsync(...) / InjectSqsMessageAsync(...) and Response
IAwsQueueConsumerRouter Router contract. ProcessSqsEventAsync(SQSEvent) is obsolete: it can't return a partial-batch response, so it throws when any record fails. Return IServiceProvider.RouteIncomingSqsQueueEventAsync(sqsEvent, queue) from the handler instead
SqsHandlerExtensions Dispatch a Lambda SQSEvent to a registered consumer (see below)
ParallelItemResultSqsExtensions AcknowledgeOrFailAsync for partial-batch responses (see below)
UseAwsSqsQueueProvider() On IServiceProvider and IApplicationBuilder: starts IQueueManagerService.Initialize() in the background, which resolves every registered queue and creates missing ones
AwsSqsQueueMessageHelper TryParseSqsMessage, ComposeMessageAttributes, and ProcessSqsMessageAttributes, which unwraps a nested $$body / $$messageAttributes envelope
AwsQueueConstants, SqsQueueManagerMessage, BootManagerHelper Shared constants and helpers

The four abstract consumer base classes that implement these interfaces ship in Kanject.Core.Queue.Provider.AwsSqs:

Base class Implements Consumes
AbstractQueueConsumer<T> IAwsQueueConsumer (also a BackgroundService polling loop) List<MessageContext<T>>
AbstractRawQueueConsumer IAwsQueueConsumer (also a BackgroundService) List<Message>
AbstractRouteQueueConsumer<T> IAwsRouteQueueConsumer List<MessageContext<T>>
AbstractRouteRawQueueConsumer IAwsRouteQueueConsumer IEnumerable<Message>

CancellationToken propagation

Every Inject* and ConsumeAsync entry point on the four consumer base classes has an overload that takes a CancellationToken. The parameterless overloads remain as facades, so consumers written against them keep working unchanged.

Overload pairs

Consumer base Without token With token
AbstractQueueConsumer<T> InjectSqsEvent(SQSEvent) → ConsumeAsync(List<MessageContext<T>>) InjectSqsEvent(SQSEvent, CancellationToken) → ConsumeAsync(List<MessageContext<T>>, CancellationToken)
AbstractRouteQueueConsumer<T> InjectSqsEventAsync(SQSEvent) / InjectSqsMessageAsync(IEnumerable<Message>) same, plus a trailing CancellationToken
AbstractRawQueueConsumer InjectSqsEvent(SQSEvent) → ConsumeAsync(List<Message>) InjectSqsEvent(SQSEvent, CancellationToken) → ConsumeAsync(List<Message>, CancellationToken)
AbstractRouteRawQueueConsumer InjectSqsEventAsync(SQSEvent) / InjectSqsMessageAsync(IEnumerable<Message>) same, plus a trailing CancellationToken

Both ConsumeAsync(...) hooks are virtual. The token-aware overload's default implementation calls the one without a token, so an existing override keeps receiving batches. New consumers can override only the token-aware overload. If a consumer overrides neither, the base throws NotImplementedException naming both overloads.

On the interfaces, the token-aware IAwsQueueConsumer.InjectSqsEvent(SQSEvent, CancellationToken) and the matching IAwsRouteQueueConsumer members are default interface methods. They delegate to the token-less member, so custom implementations don't have to add them.

Lambda → consumer flow

All four SqsHandlerExtensions entry points take an optional CancellationToken cancellationToken = default:

  • IServiceProvider.ProcessSqsEventWithQueueConsumerAsync<TQueueConsumer>(SQSEvent, CancellationToken)
  • IServiceProvider.ProcessSqsEventWithRouteQueueConsumerAsync<TRouteQueueConsumer>(SQSEvent, CancellationToken)
  • IApplicationBuilder.ProcessSqsEventWithQueueConsumerAsync<TQueueConsumer>(SQSEvent, CancellationToken)
  • IApplicationBuilder.ProcessSqsEventWithRouteQueueConsumerAsync<TRouteQueueConsumer>(SQSEvent, CancellationToken)

Each resolves the consumer, injects the event, and returns the consumer's Response. Every inject call starts a new, empty Response, so a failure added during one invocation isn't returned again by the next invocation on a warm container. If the consumer isn't registered, they throw InvalidOperationException rather than return an empty SQSBatchResponse. An empty response would tell the event source mapping that every record succeeded, and it would delete a batch that nothing processed. The failed invocation makes SQS retry the batch and eventually move it to the dead-letter queue, whether or not the mapping has ReportBatchItemFailures enabled, and it shows in the function's error metrics.

using Amazon.Lambda.Core;
using Amazon.Lambda.SQSEvents;
using Kanject.Core.Queue.Provider.AwsSqs.Abstractions.Extensions;

public sealed class Function
{
    private readonly IServiceProvider _serviceProvider = /* your built container */;

    public async Task<SQSBatchResponse> Handler(SQSEvent sqsEvent, ILambdaContext context)
    {
        // ILambdaContext exposes RemainingTime rather than a token: derive one that
        // fires shortly before the invocation times out.
        var budget = context.RemainingTime - TimeSpan.FromSeconds(1);
        using var cts = new CancellationTokenSource(budget > TimeSpan.Zero ? budget : TimeSpan.Zero);

        return await _serviceProvider
            .ProcessSqsEventWithQueueConsumerAsync<OrderQueueConsumer>(sqsEvent, cts.Token);
    }
}

The generated Process{Consumer}EventAsync(sqsEvent, cancellationToken) methods from Kanject.Core.Queue.Provider.AwsSqs.Annotations forward the token the same way.

Opting in inside a consumer

Override the token-aware overload. You don't need the other one when you do:

[QueueConsumer(Message = typeof(Order))]
[QueueConsumerDependency(typeof(IOrderService))]
public sealed partial class OrderQueueConsumer
{
    protected override async Task ConsumeAsync(
        List<MessageContext<Order>> messageContexts,
        CancellationToken cancellationToken)
    {
        foreach (var context in messageContexts)
        {
            cancellationToken.ThrowIfCancellationRequested();
            await OrderService.ProcessAsync(context.Message, cancellationToken);
            await AcknowledgeAsync(context);
        }
    }
}

Here the generator supplies the AbstractQueueConsumer<Order> base class, the constructor and the OrderService property. A class that derives from the base class directly overrides the same method.

When the consumer runs as a hosted service outside Lambda, its polling loop passes its own token into the token-aware overload. StopAsync cancels that token, so a host shutdown surfaces inside ConsumeAsync too.

Partial-batch ack/fail (AcknowledgeOrFailAsync)

ParallelItemResultSqsExtensions.AcknowledgeOrFailAsync bridges a ParallelItemResult<TInput, TResult> into the Lambda SQSBatchResponse. The result comes from the {Method}ParallelOutcomesAsync method that the [Parallel] feature in the Kanject.Core package generates.

  • Successful per-item outcomes invoke the acknowledge callback you supply, typically AcknowledgeAsync.
  • Every other outcome adds a BatchItemFailure keyed by the message's MessageId. That covers faulted, timed-out and cancelled outcomes, plus any outcome whose Value is null.
using Kanject.Core.Annotations.Attributes.Parallel;
using Kanject.Core.Queue.Abstractions.Models;
using Kanject.Core.Queue.Provider.AwsSqs.Abstractions.Extensions;

public sealed partial class OrderQueueConsumer
{
    // The generated extension calls this method from outside the class, so it can't be private.
    [Parallel(MaxDegreeOfParallelism = 8, RetryCount = 1, RetryBackoffMs = 250)]
    internal async ValueTask<bool> ProcessAsync(Order order, CancellationToken ct)
    {
        await OrderService.HandleAsync(order, ct);
        return true;
    }

    protected override async Task ConsumeAsync(
        List<MessageContext<Order>> messageContexts,
        CancellationToken cancellationToken)
    {
        var inputs = messageContexts.Select(messageContext => messageContext.Message).ToArray();

        // A trailing "Async" is dropped: ProcessAsync → ProcessParallelOutcomesAsync.
        // It is an extension method on this type, so call it through `this`.
        var result = await this.ProcessParallelOutcomesAsync(inputs, cancellationToken: cancellationToken);

        await result.AcknowledgeOrFailAsync(messageContexts, Response, AcknowledgeAsync);
    }
}

A second overload takes IReadOnlyList<Message> and a Func<Message, Task> callback for raw consumers.

Index alignment

AcknowledgeOrFailAsync uses outcome.Index to look up messageContexts[outcome.Index]. The list you pass to the helper must be index-aligned with the inputs you handed to ProcessParallelOutcomesAsync. Build the input array with .Select(messageContext => messageContext.Message) over that same list, and don't reorder or filter either side afterwards.

If an outcome index falls outside the list, the helper throws an ArgumentException for messageContexts whose message names the bad outcome index. Accidental reordering or filtering bugs then surface as a clear failure instead of a low-context IndexOutOfRangeException.

Custom success predicate

The default predicate is outcome.IsSuccess && outcome.Value is not null. That works for value-type results, which are never null. For reference-type results, it requires the body to produce a value before the message is acknowledged. Override it when your body's return value isn't the success signal:

await result.AcknowledgeOrFailAsync(
    messageContexts,
    Response,
    AcknowledgeAsync,
    successPredicate: static outcome => outcome.IsSuccess); // value can legitimately be null

Why this exists

Without the helper, each consumer ends up writing the same boilerplate:

foreach (var outcome in result.Outcomes)
{
    if (outcome.IsSuccess && outcome.Value is not null)
        await AcknowledgeAsync(messageContexts[outcome.Index]);
    else
        Response.BatchItemFailures.Add(new SQSBatchResponse.BatchItemFailure
        {
            ItemIdentifier = messageContexts[outcome.Index].MessageId
        });
}

The extension reduces this to a single await result.AcknowledgeOrFailAsync(...) call and keeps the partial-batch contract in one place.

Package Role Availability
Kanject.Core.Queue.Provider.AwsSqs SQS implementation: queue manager, consumer base classes, registration extensions nuget.org
Kanject.Core.Queue.Abstractions Provider-neutral queue contracts this package builds on nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Annotations Source generator for [QueueMessage] / [QueueConsumer] / [RouteQueueConsumer] nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes The queue attributes nuget.org
Kanject.Core The [Parallel] engine and its ParallelItemResult<TInput, TResult> outcomes nuget.org
Kanject.Core.Queue.Provider.AwsSqs.Extensions.EventBridge Scheduled and recurring delivery through EventBridge Scheduler 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 (1)

Showing the top 1 NuGet packages that depend on Kanject.Core.Queue.Provider.AwsSqs.Abstractions:

Package Downloads
Kanject.Core.Queue.Provider.AwsSqs

Kanject Core Queue Provider AWS Sqs

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
3.9.0 37 10/2/2026
3.8.2 58 10/1/2026
3.8.1 119 9/27/2026
3.8.0 94 9/27/2026
3.7.7 102 9/26/2026
3.7.6 133 9/7/2026
3.7.5 106 8/27/2026
3.7.4 128 8/22/2026
3.7.3 132 8/10/2026
3.7.2 122 8/9/2026
3.7.1 140 8/5/2026
3.7.0 131 8/5/2026
3.6.0 135 8/3/2026
3.5.8 147 7/30/2026
3.5.7 140 7/18/2026
3.5.6 137 7/13/2026
3.5.5 152 7/11/2026
3.5.4 153 7/11/2026
3.5.3 174 7/9/2026
3.5.2 144 7/9/2026