Kanject.Core.Recurring.Abstractions 1.6.0

Prefix Reserved
There is a newer version of this package available.
See the version list below for details.
dotnet add package Kanject.Core.Recurring.Abstractions --version 1.6.0
                    
NuGet\Install-Package Kanject.Core.Recurring.Abstractions -Version 1.6.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.Recurring.Abstractions" Version="1.6.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Kanject.Core.Recurring.Abstractions" Version="1.6.0" />
                    
Directory.Packages.props
<PackageReference Include="Kanject.Core.Recurring.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.Recurring.Abstractions --version 1.6.0
                    
#r "nuget: Kanject.Core.Recurring.Abstractions, 1.6.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.Recurring.Abstractions@1.6.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.Recurring.Abstractions&version=1.6.0
                    
Install as a Cake Addin
#tool nuget:?package=Kanject.Core.Recurring.Abstractions&version=1.6.0
                    
Install as a Cake Tool

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}RecurringAsync extension that runs a method on a drift-corrected schedule.
  • [RecurringHosted] also generates a BackgroundService and an Add{Type}{Method}Recurring registration 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:

  1. Acquire the lease (LeaseKey only): TryAcquireLeaseAsync(LeaseKey, TimeSpan.FromSeconds(LeaseDurationSeconds)). If Acquired is false, it skips this tick.
  2. Run your method.
  3. Save a checkpoint (CheckpointKey only, and only if the method didn't throw): read the stored one with TryGetCheckpointAsync(CheckpointKey), then SaveCheckpointAsync(CheckpointKey, new RecurringCheckpoint(stored.Iteration + 1, DateTimeOffset.UtcNow, stored.OpaqueStateJson)). With no stored checkpoint, the iteration is 1 and the state is null.
  4. Release the lease (LeaseKey only): ReleaseLeaseAsync(LeaseKey), in a finally block, but only if LeaseDurationSeconds hasn'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. LeaseDurationSeconds defaults 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. ReleaseLeaseAsync takes only the key, so implementations delete whatever lease is stored there. If an iteration runs past LeaseDurationSeconds, 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 call ReleaseLeaseAsync yourself, apply the same check.
  • Errors. If TryAcquireLeaseAsync throws (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 as Acquired = false.
  • Setup requirements.
    • The generated code references IRecurringDataProvider by name, so a project that uses LeaseKey or CheckpointKey must 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.
  • Checkpoint contents.
    • Iteration continues from the stored checkpoint, so it keeps counting across restarts. With LeaseKey set, 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.
    • LastCompletedAt is the time the save ran. OpaqueStateJson is 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, inject IRecurringDataProvider and call TryGetCheckpointAsync(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
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 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.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