From 31edca9a2eb38c2ead75cd604c0c58873c6a38e8 Mon Sep 17 00:00:00 2001 From: Michael McQuade Date: Fri, 2 Oct 2026 20:52:58 -0500 Subject: [PATCH] feat(core): a client can work jobs without standing for leader --- client.go | 4 +++- config.go | 7 +++++++ docs/PLAN.md | 4 ++++ docs/operations.md | 8 +++++++- reliability_test.go | 49 +++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 70 insertions(+), 2 deletions(-) diff --git a/client.go b/client.go index 7810d00..e07a7e2 100644 --- a/client.go +++ b/client.go @@ -174,7 +174,9 @@ func (c *Client[TTx]) Start(ctx context.Context) error { if c.caps.Listen { c.listenWG.Go(func() { c.listenLoop(c.bgCtx) }) } - c.leaderWG.Go(func() { c.leaderLoop(c.claimCtx) }) + if !c.cfg.NeverLead { + c.leaderWG.Go(func() { c.leaderLoop(c.claimCtx) }) + } c.leaseWG.Go(func() { c.leaseLoop(c.bgCtx) }) c.state = clientStarted diff --git a/config.go b/config.go index 90c5bc5..60ff36d 100644 --- a/config.go +++ b/config.go @@ -60,6 +60,13 @@ type Config struct { // consumers that start from the earliest event or are moved back. // Defaults to 7 days; negative keeps them forever. StreamRetention time.Duration + // NeverLead keeps this client out of leader elections. It still claims + // and works jobs, but never runs the leader's duties: rescuing jobs from + // lost clients, retention, maintenance, periodic jobs and stream + // delivery. Use it for clients that share an install but must not run it, + // such as a developer's process against a shared environment. At least + // one client of every install must be able to lead. + NeverLead bool // LeaseTTL is how long this client's lease lasts without a renewal. It // is renewed three times per TTL; when it lapses, the client's running // jobs are rescued and the client cancels them. A longer lease rides diff --git a/docs/PLAN.md b/docs/PLAN.md index 734e953..9e74203 100644 --- a/docs/PLAN.md +++ b/docs/PLAN.md @@ -767,6 +767,10 @@ Leader duties: periodic jobs, the rescuer, expiring TTL'd jobs, creating history partitions ahead of time and dropping expired ones, removing stale `hopper_clients` rows, and resolving batch and workflow completion. +A client with `Config.NeverLead` never attempts the lease. It claims and works jobs +like any other, so an install can include processes that should do work but not run +the install, such as a developer's process against a shared environment. + ### 7.9 Periodic jobs `hopper.Every(d, …)` and `hopper.Cron(spec, …)` are both available from v0.1. Cron supports standard five-field syntax, an optional seconds field, `@hourly`-style diff --git a/docs/operations.md b/docs/operations.md index ed04a94..bca7e28 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -10,9 +10,15 @@ traffic (claims, finalizes, leases, leader duties) plus one dedicated connection for `LISTEN`, outside the pool. Give the pool what the application needs on top of that. -- **Replicas.** Adding replicas adds throughput. Every replica competes for +- **Replicas.** Adding replicas adds throughput. Every replica works jobs; + every replica except one with `Config.NeverLead` also competes for leadership; only the leader runs periodic jobs, rescue and retention, and a replacement takes over within 15 seconds of the leader stopping. +- **Clients that must not lead.** Set `Config.NeverLead` on a client that + shares an install it should not run, such as a developer's process pointed + at a shared environment. It claims and works jobs like any other, but never + rescues, prunes, runs periodic jobs or delivers streams. An install whose + clients all set it has no leader, so nothing is ever rescued or pruned. ## Postgres diff --git a/reliability_test.go b/reliability_test.go index 2889a41..9e7d053 100644 --- a/reliability_test.go +++ b/reliability_test.go @@ -305,6 +305,55 @@ func TestFencedClientLetsALastAttemptFinish(t *testing.T) { } } +func TestNeverLeadClientWorksJobsWithoutLeading(t *testing.T) { + t.Parallel() + h := newHarness(t) + ctx := context.Background() + + worked := make(chan hopper.JobID, 1) + workers := hopper.NewWorkers() + hopper.AddWorkFunc(workers, func(_ context.Context, job *hopper.Job[noop]) error { + worked <- job.ID + return nil + }) + follower := h.client(&hopper.Config{ + Queues: map[string]hopper.QueueConfig{hopper.QueueDefault: {MaxWorkers: 1}}, + Workers: workers, + NeverLead: true, + }) + if err := follower.Start(ctx); err != nil { + t.Fatal(err) + } + + res, err := follower.Insert(ctx, noop{}, nil) + if err != nil { + t.Fatal(err) + } + select { + case id := <-worked: + if id != res.Job.ID { + t.Fatalf("worked %s, want %s", id, res.Job.ID) + } + case <-time.After(10 * time.Second): + t.Fatal("the job was never worked") + } + + // Several leader intervals pass without it standing. + time.Sleep(500 * time.Millisecond) + if follower.IsLeader() { + t.Fatal("a NeverLead client became leader") + } + if n := h.count("SELECT count(*) FROM hopper_leader"); n != 0 { + t.Errorf("leader rows = %d, want 0", n) + } + + leader := h.started(hopper.NewWorkers(), 1) + waitFor(t, leader.IsLeader) + if follower.IsLeader() { + t.Error("the NeverLead client became leader alongside another") + } +} + func TestLeaderFailover(t *testing.T) { t.Parallel() h := newHarness(t)