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
4 changes: 3 additions & 1 deletion client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions docs/PLAN.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 7 additions & 1 deletion docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
49 changes: 49 additions & 0 deletions reliability_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading