Kanject.Core.Recurring.Abstractions
1.7.0
Prefix Reserved
dotnet add package Kanject.Core.Recurring.Abstractions --version 1.7.0
NuGet\Install-Package Kanject.Core.Recurring.Abstractions -Version 1.7.0
<PackageReference Include="Kanject.Core.Recurring.Abstractions" Version="1.7.0" />
<PackageVersion Include="Kanject.Core.Recurring.Abstractions" Version="1.7.0" />
<PackageReference Include="Kanject.Core.Recurring.Abstractions" />
paket add Kanject.Core.Recurring.Abstractions --version 1.7.0
#r "nuget: Kanject.Core.Recurring.Abstractions, 1.7.0"
#:package Kanject.Core.Recurring.Abstractions@1.7.0
#addin nuget:?package=Kanject.Core.Recurring.Abstractions&version=1.7.0
#tool nuget:?package=Kanject.Core.Recurring.Abstractions&version=1.7.0
Kanject.Core.Recurring.Abstractions
The storage contract for coordinating recurring background jobs. IRecurringDataProvider supplies two things:
- Leases: named locks that expire on their own, so only one instance runs a job's iteration at a time.
- Checkpoints: durable records of a job's last successful run.
Kanject.Core's [RecurringHosted] source generator calls this contract when a job sets LeaseKey or CheckpointKey.
Most applications get this package through an implementation such as Kanject.Core.Recurring.Provider.CacheDb. Reference it directly when you write your own implementation or read checkpoints from application code.
Installation
dotnet add package Kanject.Core.Recurring.Abstractions
Targets .NET 8, .NET 9 and .NET 10. The package is marked trimming / Native AOT compatible. It depends on Kanject.Core.CacheDb.Abstractions, which brings in Kanject.Core.
Quick start: how it plugs into [RecurringHosted]
Kanject.Core ships two attributes:
[Recurring]generates a{Method}RecurringAsyncextension that runs a method on a drift-corrected schedule.[RecurringHosted]also generates aBackgroundServiceand anAdd{Type}{Method}Recurringregistration method.
By default, every running instance runs its own loop. Set LeaseKey and/or CheckpointKey to opt in to coordination:
using Kanject.Core.Annotations.Attributes.Recurring;
using Kanject.Core.Annotations.Attributes.Recurring.Enums;
namespace Shop.Jobs;
public sealed class OutboxRelay(IOutboxPublisher publisher)
{
[Recurring(15, RecurringRateUnit.Seconds)]
[RecurringHosted(LeaseKey = "outbox-relay", LeaseDurationSeconds = 60, CheckpointKey = "outbox-relay")]
public Task PublishAsync(CancellationToken cancellationToken) =>
publisher.PublishPendingAsync(cancellationToken);
}
public interface IOutboxPublisher
{
Task PublishPendingAsync(CancellationToken cancellationToken);
}
using Kanject.Core.Recurring.Abstractions;
using Shop.Jobs;
builder.Services.AddSingleton<IOutboxPublisher, SqlOutboxPublisher>(); // your implementation
builder.Services.AddSingleton<IRecurringDataProvider, MyRecurringDataProvider>(); // or an implementation package's registration
builder.Services.AddOutboxRelayPublishAsyncRecurring(); // generated
The generated service gets IRecurringDataProvider from DI and does this on every tick:
- Acquire the lease (
LeaseKeyonly):TryAcquireLeaseAsync(LeaseKey, TimeSpan.FromSeconds(LeaseDurationSeconds)). IfAcquiredisfalse, it skips this tick. - Run your method.
- Save a checkpoint (
CheckpointKeyonly, and only if the method didn't throw): read the stored one withTryGetCheckpointAsync(CheckpointKey), thenSaveCheckpointAsync(CheckpointKey, new RecurringCheckpoint(stored.Iteration + 1, DateTimeOffset.UtcNow, stored.OpaqueStateJson)). With no stored checkpoint, the iteration is 1 and the state isnull. - Release the lease (
LeaseKeyonly):ReleaseLeaseAsync(LeaseKey), in afinallyblock, but only ifLeaseDurationSecondshasn't elapsed since the acquire call.
What the generated integration guarantees
- No overlapping iterations across instances. The lease is taken and released on every iteration. There's no permanent leader, and it doesn't guarantee exactly one run per interval: an instance whose tick comes after the previous holder released will also run.
- Lease length.
LeaseDurationSecondsdefaults to 60, and values of 0 or less also become 60. Set it longer than your slowest iteration, so the lease doesn't expire while the method is still running. - Overrunning the lease.
ReleaseLeaseAsynctakes only the key, so implementations delete whatever lease is stored there. If an iteration runs pastLeaseDurationSeconds, another instance may already hold the key. The generated service therefore skips the release once the lease time is up. Its own lease has expired by then anyway. Once the lease expires, a second instance can still start an overlapping iteration, so size the lease for your slowest run. If you callReleaseLeaseAsyncyourself, apply the same check. - Errors. If
TryAcquireLeaseAsyncthrows (a store failure, for example), the exception is treated as that iteration's failure and handled by the job's exception policy. "Already held" is not an exception; it comes back asAcquired = false. - Setup requirements.
- The generated code references
IRecurringDataProviderby name, so a project that usesLeaseKeyorCheckpointKeymust reference this package, directly or through an implementation package. - An implementation must be registered, or the hosted service can't be constructed when the host starts.
- The generated code references
- Checkpoint contents.
Iterationcontinues from the stored checkpoint, so it keeps counting across restarts. WithLeaseKeyset, it also counts across instances, because the read and the save happen under the lease. Without a lease, replicas race and can write the same number.LastCompletedAtis the time the save ran.OpaqueStateJsonis copied from the stored checkpoint unchanged.- The checkpoint isn't passed to your method, because hosted methods take at most a
CancellationToken. To resume from it, injectIRecurringDataProviderand callTryGetCheckpointAsync(CheckpointKey). - To persist a cursor, save it under the same key from inside your method (see below). The generated save runs after your method returns and reads the checkpoint again, so it keeps your
OpaqueStateJson.
Reading and extending a checkpoint
From inside the job, read the checkpoint to resume, and write your cursor into OpaqueStateJson. Keep the stored Iteration, because the generated save adds one to it.
using Kanject.Core.Annotations.Attributes.Recurring;
using Kanject.Core.Annotations.Attributes.Recurring.Enums;
using Kanject.Core.Recurring.Abstractions;
using Kanject.Core.Recurring.Abstractions.Models;
public sealed class OrderExporter(IRecurringDataProvider recurring, IOrderFeed feed)
{
private const string Key = "order-export";
[Recurring(1, RecurringRateUnit.Minutes)]
[RecurringHosted(LeaseKey = Key, CheckpointKey = Key)]
public async Task ExportAsync(CancellationToken cancellationToken)
{
var stored = await recurring.TryGetCheckpointAsync(Key, cancellationToken);
var cursor = await feed.ExportAfterAsync(stored?.OpaqueStateJson, cancellationToken);
var checkpoint = stored ?? new RecurringCheckpoint(0, DateTimeOffset.UtcNow, null);
await recurring.SaveCheckpointAsync(Key, checkpoint with { OpaqueStateJson = cursor }, cancellationToken: cancellationToken);
}
}
public interface IOrderFeed
{
/// <summary>Exports orders after <paramref name="cursorJson"/> and returns the new cursor as JSON.</summary>
Task<string?> ExportAfterAsync(string? cursorJson, CancellationToken cancellationToken);
}
From anywhere else, for example a status endpoint:
using Kanject.Core.Recurring.Abstractions;
app.MapGet("/jobs/outbox-relay", async (IRecurringDataProvider recurring, CancellationToken ct) =>
{
var checkpoint = await recurring.TryGetCheckpointAsync("outbox-relay", ct);
return checkpoint is { } last
? Results.Ok(new { last.Iteration, last.LastCompletedAt })
: Results.NotFound();
});
Implementing IRecurringDataProvider
| Member | Contract |
|---|---|
TryAcquireLeaseAsync(leaseKey, duration, holderData, cancellationToken) |
Never waits for the lease. If the caller gets it: Acquired = true, LeaseId, RemainingDuration = TimeSpan.Zero. If someone else holds it: Acquired = false, the blocking lease's LeaseId, the holder's CurrentHolderData when available, and roughly how long until it expires. |
ReleaseLeaseAsync(leaseKey, cancellationToken) |
Idempotent: releasing twice, or after expiry, does nothing. |
TryGetCheckpointAsync(checkpointKey, cancellationToken) |
The most recently saved checkpoint, or null if none was ever stored. |
SaveCheckpointAsync(checkpointKey, checkpoint, retention, cancellationToken) |
May overwrite without conditions. When the job also sets LeaseKey, the generated service reads and saves under the lease, so under correct usage no other writer competes. retention: null means the implementation's default retention. |
Register the implementation as a singleton. It shouldn't hold any mutable state besides its store. This single-process version works well as a test double:
using System.Collections.Concurrent;
using Kanject.Core.Recurring.Abstractions;
using Kanject.Core.Recurring.Abstractions.Models;
public sealed class InProcessRecurringDataProvider : IRecurringDataProvider
{
private readonly object _gate = new();
private readonly Dictionary<string, (DateTimeOffset ExpiresAt, string? Holder)> _leases = new();
private readonly ConcurrentDictionary<string, RecurringCheckpoint> _checkpoints = new();
public Task<RecurringLeaseResult> TryAcquireLeaseAsync(
string leaseKey, TimeSpan duration, string? holderData = null,
CancellationToken cancellationToken = default)
{
var now = DateTimeOffset.UtcNow;
lock (_gate)
{
if (_leases.TryGetValue(leaseKey, out var held) && held.ExpiresAt > now)
return Task.FromResult(new RecurringLeaseResult(false, leaseKey, held.Holder, held.ExpiresAt - now));
_leases[leaseKey] = (now + duration, holderData);
return Task.FromResult(new RecurringLeaseResult(true, leaseKey, null, TimeSpan.Zero));
}
}
public Task ReleaseLeaseAsync(string leaseKey, CancellationToken cancellationToken = default)
{
lock (_gate) _leases.Remove(leaseKey);
return Task.CompletedTask;
}
public Task<RecurringCheckpoint?> TryGetCheckpointAsync(
string checkpointKey, CancellationToken cancellationToken = default) =>
Task.FromResult(_checkpoints.TryGetValue(checkpointKey, out var checkpoint)
? checkpoint
: (RecurringCheckpoint?)null);
public Task SaveCheckpointAsync(
string checkpointKey, RecurringCheckpoint checkpoint, TimeSpan? retention = null,
CancellationToken cancellationToken = default)
{
_checkpoints[checkpointKey] = checkpoint;
return Task.CompletedTask;
}
}
Public surface at a glance
| Type | Purpose |
|---|---|
IRecurringDataProvider |
Lease and checkpoint storage contract |
RecurringLeaseResult(bool Acquired, string LeaseId, string? CurrentHolderData, TimeSpan RemainingDuration) |
readonly record struct returned by TryAcquireLeaseAsync |
RecurringCheckpoint(long Iteration, DateTimeOffset LastCompletedAt, string? OpaqueStateJson) |
readonly record struct for a job's last successful run; OpaqueStateJson is caller-owned JSON state |
Related packages
| Package | Role | Availability |
|---|---|---|
Kanject.Core |
The [Recurring] / [RecurringHosted] attributes, generators and scheduling engine |
nuget.org |
Kanject.Core.Recurring.Provider.CacheDb |
IRecurringDataProvider implementation over any ICacheDb |
nuget.org |
Kanject.Core.CacheDb.Provider.InMemory |
In-process ICacheDb, for local development with the CacheDb implementation |
nuget.org |
Kanject.Core.CacheDb.Provider.DynamoDb |
Durable ICacheDb for multi-instance leases and checkpoints |
Commercial license (not on nuget.org) |
Kanject.Core.CacheDb.Provider.S3Express |
Durable ICacheDb for multi-instance leases and checkpoints |
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
- Kanject.Core.CacheDb.Abstractions (>= 3.9.0)
-
net8.0
- Kanject.Core.CacheDb.Abstractions (>= 3.9.0)
-
net9.0
- Kanject.Core.CacheDb.Abstractions (>= 3.9.0)
NuGet packages (1)
Showing the top 1 NuGet packages that depend on Kanject.Core.Recurring.Abstractions:
| Package | Downloads |
|---|---|
|
Kanject.Core.Recurring.Provider.CacheDb
IRecurringDataProvider implementation adapting Kanject.Core.CacheDb.Abstractions.ICacheDb for leases and checkpoints |
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 1.7.0 | 36 | 10/2/2026 |
| 1.6.1 | 131 | 9/27/2026 |
| 1.6.0 | 100 | 9/27/2026 |
| 1.5.7 | 104 | 9/26/2026 |
| 1.5.6 | 152 | 9/7/2026 |
| 1.5.5 | 119 | 8/27/2026 |
| 1.5.4 | 130 | 8/22/2026 |
| 1.5.3 | 144 | 8/10/2026 |
| 1.5.2 | 122 | 8/9/2026 |
| 1.5.1 | 124 | 8/5/2026 |
| 1.5.0 | 131 | 8/5/2026 |
| 1.4.0 | 138 | 8/3/2026 |
| 1.3.8 | 140 | 7/30/2026 |
| 1.3.7 | 148 | 7/18/2026 |
| 1.3.6 | 142 | 7/13/2026 |
| 1.3.5 | 136 | 7/11/2026 |
| 1.3.4 | 142 | 7/11/2026 |
| 1.3.3 | 178 | 7/9/2026 |
| 1.3.2 | 133 | 7/9/2026 |