Kanject.Core.Queue.Provider.AwsSqs.Abstractions
3.8.2
Prefix Reserved
See the version list below for details.
dotnet add package Kanject.Core.Queue.Provider.AwsSqs.Abstractions --version 3.8.2
NuGet\Install-Package Kanject.Core.Queue.Provider.AwsSqs.Abstractions -Version 3.8.2
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs.Abstractions" Version="3.8.2" />
<PackageVersion Include="Kanject.Core.Queue.Provider.AwsSqs.Abstractions" Version="3.8.2" />
<PackageReference Include="Kanject.Core.Queue.Provider.AwsSqs.Abstractions" />
paket add Kanject.Core.Queue.Provider.AwsSqs.Abstractions --version 3.8.2
#r "nuget: Kanject.Core.Queue.Provider.AwsSqs.Abstractions, 3.8.2"
#:package Kanject.Core.Queue.Provider.AwsSqs.Abstractions@3.8.2
#addin nuget:?package=Kanject.Core.Queue.Provider.AwsSqs.Abstractions&version=3.8.2
#tool nuget:?package=Kanject.Core.Queue.Provider.AwsSqs.Abstractions&version=3.8.2
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:
AwsSqsQueueConfigurationand the queue registry contracts.- Consumer interfaces that accept Lambda
SQSEvents and expose anSQSBatchResponse. IServiceProvider/IApplicationBuilderextensions (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
BatchItemFailurekeyed by the message'sMessageId. That covers faulted, timed-out and cancelled outcomes, plus any outcome whoseValue 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.
Related packages
| 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 | Versions 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. |
-
net10.0
- Amazon.Lambda.SQSEvents (>= 3.0.1)
- AWSSDK.SQS (>= 4.0.100.15)
- Kanject.Core.Queue.Abstractions (>= 3.7.1)
- Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes (>= 3.8.1)
-
net8.0
- Amazon.Lambda.SQSEvents (>= 3.0.1)
- AWSSDK.SQS (>= 4.0.100.15)
- Kanject.Core.Queue.Abstractions (>= 3.7.1)
- Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes (>= 3.8.1)
-
net9.0
- Amazon.Lambda.SQSEvents (>= 3.0.1)
- AWSSDK.SQS (>= 4.0.100.15)
- Kanject.Core.Queue.Abstractions (>= 3.7.1)
- Kanject.Core.Queue.Provider.AwsSqs.Annotations.Attributes (>= 3.8.1)
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 |