Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ Reviewers: treat a hardcoded provider name in code, config, or a committed doc l

### Deploy invariant: exactly ONE serving API instance (#271)

**Run one serving instance.** Every instance runs `DurableJobWorker`, and its poll claims nothing — no `FOR UPDATE SKIP LOCKED`, no lease, no advisory lock. Three independent blockers remain before scaling, and they are not all in #271: background work has no single-runner guarantee (**#271**); the IP-keyed auth limiters (#143) are in-process; and the per-account report concurrency cap (#311) is in-process. The step-up grant registry was the fourth — the blocker with teeth — and **#338 closed it**: replay now lives in the shared `IClaimOnceStore` (#543) and logout revocation in the durable per-user `ApplicationUser.StepUpLogoutEpoch`, an integer compared for equality (never a timestamp), so both survive across replicas without a shared clock. **#307 is closed and does not license scaling.**
**Run one serving instance.** Every instance runs `DurableJobWorker`, but its poll and the three recurring sweeps now run **only under a single-leader gate** — a session-scoped Postgres advisory lock (`pg_try_advisory_lock`) on a dedicated, non-pooled connection (**#271**, closed): at most one instance leads, crash recovery is automatic (a dead leader's session releases the lock), and the contract is at-most-one-leader with at-least-once, idempotent handlers — never "exactly once". That guarantee holds on a **session-pinned** Postgres endpoint (a direct connection or a session-pooled proxy); under a **transaction-pooling** proxy (e.g. PgBouncer in transaction mode) the lock can migrate across backends and single-leader is *not* guaranteed — a backend-PID affinity check narrows but does not close the window, so that topology relies on the at-least-once + idempotent contract until a dedicated session-pinned lease endpoint (**#556**) lands. Two independent blockers still remain before scaling, and neither is #271: the IP-keyed auth limiters (#143) are in-process; and the per-account report concurrency cap (#311) is in-process. The step-up grant registry was another — the blocker with teeth — and **#338 closed it**: replay now lives in the shared `IClaimOnceStore` (#543) and logout revocation in the durable per-user `ApplicationUser.StepUpLogoutEpoch`, an integer compared for equality (never a timestamp), so both survive across replicas without a shared clock. **#307 is closed and does not license scaling.**

**Do not extend that list from memory** — it was twice derived wrongly, both times a process-local limiter. Re-derive it by walking every `AddSingleton`/`AddHostedService` under `src/` plus every in-memory state primitive, then excluding deliberately. The run-then-exit verbs are unaffected: they never start hosted services. The four #543 shared-state registrations (`IConnectionMultiplexer` + the `IClaimOnceStore`/`IFixedWindowCounter`/`ILease` ports) are **not** blockers — they are the shared store: #338 wired the `IClaimOnceStore` grant-replay caller (closed above), and #544/#545 will wire theirs. Redis-backed they are multi-replica-safe, their in-process fallbacks a deliberate alarmed degradation. → [`271-single-serving-instance.md`](docs/decisions/271-single-serving-instance.md)

Expand Down
27 changes: 22 additions & 5 deletions docs/decisions/271-single-serving-instance.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,23 @@
> this file is the relocated rationale, including how the blocker list was
> derived and why that method matters.

**Status:** accepted — the invariant stands until all four blockers close
**Status:** accepted — the invariant stands until the remaining blockers close
· **Date:** 2026-07

> **Update (#271 resolved).** The worker double-run this record's subject describes
> is closed: `DurableJobWorker`'s poll and the three sweeps now run only under a
> session-scoped Postgres advisory-lock **leader gate** (`PostgresLeaderLease`,
> `pg_try_advisory_lock` on a dedicated non-pooled connection), so at most one
> instance leads. Contract is at-most-one-leader with at-least-once, idempotent
> handlers — not "exactly once". Single-leader holds on a **session-pinned**
> endpoint; under a **transaction-pooling** proxy it is not guaranteed (a
> backend-PID affinity check narrows but does not close the window) and relies on
> the at-least-once + idempotent contract — a dedicated session-pinned lease
> endpoint is the fix (**#556**). The remaining single-instance blockers are #143
> and #311 (both in-process); #307 and #338 are closed. The four-blockers analysis
> below is the original (pre-#338) rationale — the holistic AGENTS.md/decision sync
> is #537.

## The rule

**Run one serving instance.** More than one breaks four separate things, and
Expand All @@ -22,10 +36,13 @@ The run-then-exit verbs (`migrate`, `seed`, `recover-admin`, `bootstrap-admin`,

`AddHostedService<DurableJobWorker>()`
(`Hosting/CluckworkJobServiceCollectionExtensions.cs`) means **every** instance
runs the worker loop, and the poll claims nothing — no `FOR UPDATE SKIP LOCKED`,
no lease, no advisory lock. What that exposes today is **the three recurring
sweeps**, which ride the same poll and run unconditionally per instance:
`DailyEntryLockSweep`, `RefreshTokenPurgeSweep`, `IdempotencyRecordPurgeSweep`.
runs the worker loop. Before the fix, the poll claimed nothing — no
`FOR UPDATE SKIP LOCKED`, no lease, no advisory lock — which exposed **the three
recurring sweeps** (`DailyEntryLockSweep`, `RefreshTokenPurgeSweep`,
`IdempotencyRecordPurgeSweep`), which ride the same poll and would run
unconditionally per instance. The leader gate (see the Update note above) closed
this: the poll and the sweeps now run only while this instance holds the advisory
lock.

The durable-job half is still a scaffold that selects pending rows and logs
them — no handlers are registered — so job double-execution is **latent, not
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
namespace Cluckwork.Api.Hosting;

using Cluckwork.Infrastructure.Jobs;
using Microsoft.Extensions.Logging;

internal static class CluckworkJobServiceCollectionExtensions
{
Expand All @@ -11,6 +12,10 @@ public static IServiceCollection AddCluckworkJobs(
services.AddSingleton<DailyEntryLockSweep>();
services.AddSingleton<RefreshTokenPurgeSweep>();
services.AddSingleton<IdempotencyRecordPurgeSweep>();
// #271 — the single-runner gate the worker acquires before polling.
services.AddSingleton<ILeaderLease>(sp => new PostgresLeaderLease(
sp.GetRequiredService<LeaderLeaseConnectionString>().Value,
sp.GetRequiredService<ILogger<PostgresLeaderLease>>()));
Comment on lines +15 to +18

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Synchronize the canonical deployment invariant

This registration introduces the advisory-lock lease, but the canonical deployment guidance and its linked decision still state that the worker poll has no lease or advisory lock and list that absence as the #271 blocker. Leaving those architecture records unchanged makes subsequent scaling audits reason from behavior that this PR explicitly replaces; update them in the same change to describe the new lease and its transaction-pooling limitation while retaining the remaining single-instance blockers.

AGENTS.md reference: AGENTS.md:L126-L130

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed the AGENTS.md deploy-invariant + decision records go stale once this merges. That edit is deliberately owned by #537 (slice T9 — 'Docs: multi-farm tenancy ADR + AGENTS.md / GLOSSARY sync'), which lands alongside whichever slice settles each claim: the file explicitly warns that blocker list 'was derived wrong twice,' so partial edits from a feature PR are exactly what that split avoids. And AGENTS.md still correctly describes main until #555 merges. I've flagged the specific #271 'no lease' wording + the transaction-pooling limitation on #537 so it's captured there.

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Update: on reflection the #271-specific piece of this is a settled doc change, so I folded it into this PR (81c0c57) rather than leaving all of it for #537. The canonical AGENTS.md deploy-invariant paragraph and the 271 decision doc now describe the advisory-lock leader gate, its session-pinned requirement + transaction-pooling limitation (#556), and the re-derived remaining blockers (#143, #311) — the blocker list re-derived by walking every AddSingleton/AddHostedService under src/, not edited from memory (that list 'was derived wrong twice'). The broader multi-slice AGENTS.md/GLOSSARY/ADR reconciliation (tenancy, login, suspension — most blocked on unlanded slices) stays #537.

services.AddHostedService<DurableJobWorker>();
return services;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ namespace Cluckwork.Api.Hosting;

using Cluckwork.Application.Common;
using Cluckwork.Infrastructure.Identity;
using Cluckwork.Infrastructure.Jobs;
using Cluckwork.Infrastructure.Persistence;
using Cluckwork.Infrastructure.Persistence.Interceptors;
using Cluckwork.Infrastructure.Providers;
Expand Down Expand Up @@ -41,6 +42,10 @@ public static CluckworkPersistenceRegistration AddCluckworkPersistence(
configuration.GetValue<bool>("Database:AllowInsecureConnection"),
onWarning: connectionStringWarnings.Add);

// #271 — the leader lease opens its own dedicated, non-pooled connection
// from the same normalised, TLS-floor-validated string the DbContext uses.
services.AddSingleton(new LeaderLeaseConnectionString(connectionString));

services.AddScoped<TenantStampInterceptor>();
services.AddDbContext<AppDbContext>((sp, options) =>
{
Expand Down
74 changes: 61 additions & 13 deletions src/Cluckwork.Infrastructure/Jobs/DurableJobWorker.cs
Original file line number Diff line number Diff line change
Expand Up @@ -21,12 +21,15 @@ public sealed class DurableJob

public enum DurableJobStatus { Pending, Running, Completed, Failed }

// Intervals are injectable for the loop-level tests only; DI fills the
// defaults (ActivatorUtilities resolves optional parameters and prefers the
// registered heartbeat singleton).
// The lease is a REQUIRED dependency (fail-closed: a host that forgets to register
// an ILeaderLease fails at startup rather than silently running every replica as
// leader). The heartbeat, sweeps and intervals are injectable for the tests only;
// DI fills the defaults (ActivatorUtilities resolves optional parameters and prefers
// the registered singletons).
public sealed class DurableJobWorker(
IServiceScopeFactory scopeFactory,
ILogger<DurableJobWorker> logger,
ILeaderLease leaderLease,
DurableJobWorkerHeartbeat? heartbeat = null,
DailyEntryLockSweep? lockSweep = null,
RefreshTokenPurgeSweep? refreshTokenPurgeSweep = null,
Expand All @@ -43,28 +46,52 @@ public sealed class DurableJobWorker(
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
// A transient DB outage must not escape ExecuteAsync: the host default
// (BackgroundServiceExceptionBehavior.StopHost) would take the whole
// API down with it (#65). Failures log and retry with capped backoff;
// StopHost stays as the backstop for anything thrown OUTSIDE the
// guarded iteration (e.g. a fatal bug in the loop itself).
// (BackgroundServiceExceptionBehavior.StopHost) would take the whole API
// down with it (#65). Failures log and retry with capped backoff; StopHost
// stays as the backstop for anything thrown OUTSIDE the guarded iteration
// (e.g. a fatal bug in the loop itself).
//
// One scheduling mechanism only: success waits the poll interval,
// failure waits the backoff — a PeriodicTimer on top would stretch
// every retry back to the poll interval and make the logged backoff
// a lie (codex review of PR #79). Task.Delay throws on cancellation,
// which ends the loop as a normal shutdown.
// One scheduling mechanism only: success waits the poll interval, failure
// waits the backoff — a PeriodicTimer on top would stretch every retry back
// to the poll interval and make the logged backoff a lie (codex review of
// PR #79). Task.Delay throws on cancellation, which ends the loop as a normal
// shutdown.
//
// #271 — before doing any work, this instance must be the single active
// leader. A follower keeps the loop alive (so its health check stays green)
// but runs neither the poll nor the sweeps; at most one instance is the leader
// at a time, so at most one runs the recurring work. A FAULTED acquisition
// (could not reach the lock at all) is NOT a follower: it backs off with the
// heartbeat left unstamped so a sustained fault degrades /health.
heartbeat.MarkStarted();
var backoff = TimeSpan.Zero;
while (!stoppingToken.IsCancellationRequested)
{
if (await TryProcessPendingJobsAsync(stoppingToken))
var leadership = await TryAcquireLeadershipAsync(stoppingToken);

// A healthy follower: the loop ran, checked leadership, and correctly
// stood down. Stamp the heartbeat so the health check reads this as a live
// worker, not a stall (#271) — the sweeps are simply not this instance's
// to run right now.
if (leadership == LeaseStatus.Follower)
{
heartbeat.MarkSuccessfulPoll();
backoff = TimeSpan.Zero;
await Task.Delay(pollInterval, stoppingToken);
continue;
}

if (leadership == LeaseStatus.Leader
&& await TryProcessPendingJobsAsync(stoppingToken))
{
heartbeat.MarkSuccessfulPoll();
backoff = TimeSpan.Zero;
await Task.Delay(pollInterval, stoppingToken);
continue;
}

// A faulted acquisition, or a leader whose poll failed: back off and leave
// the heartbeat unstamped so a sustained fault degrades /health.
backoff = backoff == TimeSpan.Zero
? initialBackoff
: TimeSpan.FromTicks(Math.Min(backoff.Ticks * 2, MaxBackoff.Ticks));
Expand All @@ -73,6 +100,27 @@ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
}
}

// One guarded leadership acquisition. The lease reports faults as
// LeaseStatus.Faulted rather than throwing; this catch is defence in depth for a
// lease impl that breaks that contract — an unexpected throw becomes Faulted so
// the loop backs off instead of the host stopping. Real cancellation propagates.
private async Task<LeaseStatus> TryAcquireLeadershipAsync(CancellationToken ct)
{
try
{
return await leaderLease.TryAcquireAsync(ct);
}
catch (OperationCanceledException) when (ct.IsCancellationRequested)
{
throw;
}
catch (Exception ex)
{
logger.LogError(ex, "Leader-lease acquisition failed.");
return LeaseStatus.Faulted;
}
}

// One guarded poll iteration. Returns false on failure instead of throwing;
// cancellation propagates so shutdown stays prompt.
internal async Task<bool> TryProcessPendingJobsAsync(CancellationToken ct)
Expand Down
50 changes: 50 additions & 0 deletions src/Cluckwork.Infrastructure/Jobs/ILeaderLease.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
namespace Cluckwork.Infrastructure.Jobs;

// #271 — the outcome of one leadership-acquisition attempt. A first-class result
// (not a bool) so the worker can tell "another instance holds the lock" (Follower,
// a healthy steady state) apart from "I could not even ask" (Faulted, a fault that
// must degrade /health). Conflating the two let a DB outage read as a healthy
// follower and silently stop all background work behind a green health check — the
// exact #69 stall-detection guarantee this must not break.
public enum LeaseStatus
{
// This instance holds the lease: run the poll and the sweeps.
Leader,
// Another instance holds the lease: stand down. The loop is alive and healthy —
// a follower does no work by design.
Follower,
// The acquisition attempt itself faulted (e.g. the DB is unreachable): back off
// and do NOT stamp the heartbeat, so a sustained fault degrades /health.
Faulted,
}

// #271 — the background worker's single-runner gate: at most one API instance (the
// "leader") runs the durable-job poll and the recurring sweeps. See
// PostgresLeaderLease for the mechanism and the exact contract (AT MOST ONE ACTIVE
// LEADER, never "exactly once").
public interface ILeaderLease
{
// The leadership status of THIS instance after the call. Safe to call every
// poll: a current leader re-affirms (re-verifying its session) without stacking
// the lock; a follower retries acquisition; a fault is reported as Faulted, not
// thrown, so the worker never lets it reach BackgroundServiceExceptionBehavior.
Task<LeaseStatus> TryAcquireAsync(CancellationToken ct);
}

// A lease that is always the leader — for a deploy that is single-instance by
// construction, and for the worker's resilience unit tests. Deliberately NOT a
// constructor default: the worker REQUIRES an explicit ILeaderLease, so a host that
// forgets to register one fails at startup (fail-closed) rather than silently
// running every replica as leader and re-introducing the #271 double-run.
public sealed class AlwaysLeaderLease : ILeaderLease
{
public Task<LeaseStatus> TryAcquireAsync(CancellationToken ct) =>
Task.FromResult(LeaseStatus.Leader);
}

// Hands the already-normalised, TLS-floor-validated connection string (registered
// by AddCluckworkPersistence) to PostgresLeaderLease without a second configuration
// lookup — the lease opens its own dedicated, non-pooled connection from this exact
// string. A tiny typed wrapper so DI injects the right string rather than an ambient
// one.
public sealed record LeaderLeaseConnectionString(string Value);
Loading
Loading