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
18 changes: 15 additions & 3 deletions client.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,7 @@ func NewClient[TTx any](d driver.Driver[TTx], cfg *Config) (*Client[TTx], error)
exec: d.Executor(),
caps: d.Capabilities(),
cfg: resolved,
tuning: defaultTuning,
tuning: defaultTuning.withLease(resolved.LeaseTTL),
logger: resolved.Logger,
workers: resolved.Workers,
leaderPoke: make(chan struct{}, 1),
Expand Down Expand Up @@ -205,7 +205,7 @@ func (c *Client[TTx]) startProducer(name string, qcfg QueueConfig, paused, limit
claim: func(ctx context.Context, limit int, limited bool) (driver.JobClaimResult, error) {
return c.exec.JobClaim(ctx, driver.JobClaimParams{Queue: name, ClientID: c.clientID.Load(), Limit: limit, Limited: limited})
},
work: func(row *driver.JobRow) driver.JobFinalize { return c.execute(c.generation(), row, qcfg) },
work: func(row *driver.JobRow) driver.JobFinalize { return c.execute(c.jobContext(row), row, qcfg) },
submit: c.finalizer.submit,
wake: make(chan struct{}, 1),
freed: make(chan struct{}, 1),
Expand Down Expand Up @@ -295,7 +295,19 @@ func (c *Client[TTx]) newGeneration() {
c.genCtx, c.genCancel = context.WithCancel(c.workCtx) //nolint:gosec // cancelled by fence or Stop
}

// generation returns the context jobs run under.
// jobContext returns the context a claimed job runs under. A job with
// attempts left runs under the lease generation, so that a client which
// loses its lease stops it before another client runs it again. A job on its
// last attempt will not run again, so cancelling it would only discard its
// work: it runs under the work context, which only a hard stop cancels.
func (c *Client[TTx]) jobContext(row *driver.JobRow) context.Context {
if row.Attempt >= row.MaxAttempts {
return c.workCtx
}
return c.generation()
}

// generation returns the context of the current lease generation.
func (c *Client[TTx]) generation() context.Context {
c.genMu.Lock()
defer c.genMu.Unlock()
Expand Down
21 changes: 21 additions & 0 deletions client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,7 @@ func TestNewClientValidation(t *testing.T) {
"empty queue name": {Queues: map[string]hopper.QueueConfig{"": {MaxWorkers: 1}}, Workers: workers},
"negative attempts": {MaxAttempts: -1},
"negative timeout": {JobTimeout: -1},
"lease under a second": {LeaseTTL: 500 * time.Millisecond},
"strict without workers": {StrictKinds: true},
}
for name, cfg := range cases {
Expand All @@ -209,6 +210,26 @@ func TestNewClientValidation(t *testing.T) {
}
}

func TestLeaseTTLSetsTheLeaseAndItsRenewal(t *testing.T) {
t.Parallel()
h := newHarness(t)
for name, tc := range map[string]struct {
cfg *hopper.Config
ttl, renew time.Duration
}{
"default": {nil, 15 * time.Second, 5 * time.Second},
"configured": {&hopper.Config{LeaseTTL: 3 * time.Minute}, 3 * time.Minute, time.Minute},
} {
c, err := hopper.NewClient(h.d, tc.cfg)
if err != nil {
t.Fatalf("%s: NewClient: %v", name, err)
}
if ttl, renew := c.LeaseTuning(); ttl != tc.ttl || renew != tc.renew {
t.Errorf("%s: lease = %s renewed every %s, want %s every %s", name, ttl, renew, tc.ttl, tc.renew)
}
}
}

func TestWorkersRejectDuplicateKind(t *testing.T) {
t.Parallel()
workers := hopper.NewWorkers()
Expand Down
28 changes: 26 additions & 2 deletions config.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,12 @@ 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
// 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
// out a longer database outage before interrupting jobs, and delays the
// rescue of a crashed client's jobs by as much. Defaults to 15 seconds.
LeaseTTL time.Duration
// RescueStuckAfter, if set, rescues jobs that have been running longer
// than this even though their client's lease is live, for workers that
// ignore their context. Off by default; timeouts cover most cases.
Expand Down Expand Up @@ -111,6 +117,7 @@ const (
DefaultMaxAttempts = 25
DefaultStopTimeout = 30 * time.Second
DefaultPollInterval = time.Second
DefaultLeaseTTL = 15 * time.Second
DefaultCompletedRetention = 24 * time.Hour
DefaultFailedRetention = 7 * 24 * time.Hour
DefaultStreamRetention = 7 * 24 * time.Hour
Expand Down Expand Up @@ -140,6 +147,9 @@ func (cfg *Config) withDefaults() (Config, error) {
if out.PollInterval == 0 {
out.PollInterval = DefaultPollInterval
}
if out.LeaseTTL == 0 {
out.LeaseTTL = DefaultLeaseTTL
}
if out.CompletedRetention == 0 {
out.CompletedRetention = DefaultCompletedRetention
}
Expand All @@ -163,6 +173,9 @@ func (cfg *Config) withDefaults() (Config, error) {
if out.JobTimeout < 0 || out.StopTimeout < 0 || out.PollInterval < 0 || out.RescueStuckAfter < 0 {
return out, errors.New("hopper: Config durations must be positive")
}
if out.LeaseTTL < time.Second {
return out, errors.New("hopper: Config.LeaseTTL must be at least one second")
}
if len(out.Queues) > 0 && out.Workers == nil {
return out, errors.New("hopper: Config.Workers is required to work queues")
}
Expand Down Expand Up @@ -232,9 +245,20 @@ type tuning struct {
streamBatch int
}

// withLease returns t with the client lease lasting ttl.
func (t tuning) withLease(ttl time.Duration) tuning {
t.leaseTTL = ttl
t.leaseRenew = ttl / leaseRenewalsPerTTL
return t
}

// leaseRenewalsPerTTL is how many renewals fit in one lease, so that two can
// fail before the lease lapses.
const leaseRenewalsPerTTL = 3

var defaultTuning = tuning{
leaseTTL: 15 * time.Second,
leaseRenew: 5 * time.Second,
leaseTTL: DefaultLeaseTTL,
leaseRenew: DefaultLeaseTTL / leaseRenewalsPerTTL,
claimCooldown: 20 * time.Millisecond,
finalizeInterval: 25 * time.Millisecond,
finalizeBatch: 500,
Expand Down
9 changes: 6 additions & 3 deletions docs/PLAN.md
Original file line number Diff line number Diff line change
Expand Up @@ -701,7 +701,8 @@ SELECT id FROM retry UNION ALL SELECT id FROM done;

### 7.6 Liveness and rescue
Liveness is tracked **per client, not per job**. Each client holds a lease row in
`hopper_clients` and renews it every 5s (the TTL is 15s):
`hopper_clients` and renews it three times per TTL (`Config.LeaseTTL`, 15s by default,
so every 5s):

```sql
UPDATE hopper_clients SET expires_at = now() + $ttl WHERE id = $me AND expires_at > now();
Expand All @@ -716,8 +717,10 @@ the largest source of write amplification in heartbeat-based designs.
`discarded` if they are out of attempts). A crashed pod's jobs run again within about
20s.
- **Fencing.** If a renewal matches zero rows (after a long GC pause or a partition),
the client has lost its lease. It cancels every in-flight job context and re-registers
under a new ID. Results still in flight are submitted anyway: the finalize fence (§7.5)
the client has lost its lease. It cancels the context of every in-flight job that has
attempts left and re-registers under a new ID. A job on its last attempt keeps
running, because the rescuer discards it instead of running it again, so no second
attempt can overlap it and cancelling would only lose its work. Results still in flight are submitted anyway: the finalize fence (§7.5)
rejects any whose job was rescued meanwhile and applies the rest, which is equivalent
to a rescue. A leader that finds its own jobs among the rescue candidates fences
itself first, so it does not re-claim them under the lapsed ID before its lease loop
Expand Down
25 changes: 17 additions & 8 deletions docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,15 +85,24 @@ for a queue's completed jobs entirely, for maximum throughput.

## Failure handling

- **A process dies.** Its lease (renewed every 5 s, 15 s TTL) expires, and the
leader moves its running jobs back to `retryable` with the error `hopper:
client lost`. Expect them to run again within about 20 seconds. Jobs out of
attempts are dead-lettered instead.
- **A process pauses** (long GC, VM stall, partition) longer than the TTL and
then resumes. Its next renewal fails; it cancels every running job's context
and re-registers under a new client ID. Results still in flight are applied
- **A process dies.** Its lease (`Config.LeaseTTL`, 15 s by default, renewed
three times per TTL) expires, and the leader moves its running jobs back to
`retryable` with the error `hopper: client lost`. Expect them to run again
within about 20 seconds with the default lease. Jobs out of attempts are
dead-lettered instead.
- **A process pauses, or loses the database** (long GC, VM stall, partition)
for longer than the lease and then resumes. Its next renewal fails; it
cancels the context of every running job that has attempts left and
re-registers under a new client ID. Results still in flight are applied
only if the job was not rescued meanwhile, so a job never runs two attempts
at once for long.
at once for long. A job on its last attempt is left running: nothing will
run it again, so cancelling it would only throw its work away. Its result
is recorded if it lands before the leader's rescue, and otherwise the job
stays dead-lettered as `client lost` although its work finished.
- **Long jobs.** When jobs are long calls to another system, raise
`Config.LeaseTTL` so that a brief database outage does not interrupt them.
The cost is that a crashed process's jobs wait that much longer to be
rescued. The lease is per client, so each program picks its own.
- **A worker hangs** ignoring its context. Timeouts (`Config.JobTimeout`, one
minute by default, or the worker's `Timeout`) cover the common case.
`Config.RescueStuckAfter` additionally rescues any job running longer than
Expand Down
5 changes: 5 additions & 0 deletions export_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,11 @@ func (c *Client[TTx]) SetTuning(t Tuning) {
c.notifier.interval = c.tuning.notifyInterval
}

// LeaseTuning returns the client lease's length and renewal interval.
func (c *Client[TTx]) LeaseTuning() (ttl, renew time.Duration) {
return c.tuning.leaseTTL, c.tuning.leaseRenew
}

// IsLeader reports whether the client currently holds the leader lease.
func (c *Client[TTx]) IsLeader() bool { return c.isLeader.Load() }

Expand Down
59 changes: 59 additions & 0 deletions reliability_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,65 @@ func TestFencedClientCancelsJobsAndReregisters(t *testing.T) {
}
}

// lastAttemptWorker reports when it starts, waits to be released, and records
// whether its context was cancelled while it waited.
type lastAttemptWorker struct {
hopper.WorkerDefaults[noop]
started chan struct{}
release chan struct{}
runs atomic.Int64
cancelled atomic.Bool
}

func (w *lastAttemptWorker) Work(ctx context.Context, _ *hopper.Job[noop]) error {
w.runs.Add(1)
w.started <- struct{}{}
select {
case <-w.release:
return nil
case <-ctx.Done():
w.cancelled.Store(true)
return ctx.Err()
}
}

func (w *lastAttemptWorker) Timeout(*hopper.Job[noop]) time.Duration { return -1 }

func TestFencedClientLetsALastAttemptFinish(t *testing.T) {
t.Parallel()
h := newHarness(t)
ctx := context.Background()
w := &lastAttemptWorker{started: make(chan struct{}, 10), release: make(chan struct{})}
workers := hopper.NewWorkers()
hopper.AddWorker(workers, w)
c := h.started(workers, 1)
oldID := c.ClientID()

res, err := c.Insert(ctx, noop{}, &hopper.InsertOpts{MaxAttempts: 1})
if err != nil {
t.Fatal(err)
}
<-w.started

if _, err := h.pool.Exec(ctx, "UPDATE hopper_clients SET expires_at = now() - interval '1s' WHERE id = $1", oldID); err != nil {
t.Fatal(err)
}
waitFor(t, func() bool { return c.ClientID() != oldID })

// The job has no attempt left, so nothing will run it again and the
// fenced client leaves it running. Whether its result or the leader's
// rescue lands first decides the recorded state, so only the run itself
// is asserted.
close(w.release)
job := waitForJob(t, c, res.Job.ID, hopper.JobStateCompleted, hopper.JobStateDiscarded)
if w.cancelled.Load() {
t.Errorf("the last attempt was cancelled when the lease was lost; job = %+v", job)
}
if runs := w.runs.Load(); runs != 1 {
t.Errorf("job ran %d times, want 1", runs)
}
}

func TestLeaderFailover(t *testing.T) {
t.Parallel()
h := newHarness(t)
Expand Down
Loading