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
3 changes: 3 additions & 0 deletions internal/leadership/elector.go
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,7 @@ func (e *Elector) attemptGainLeadershipLoop(ctx context.Context) error {

elected, err := attemptElectOrReelect(ctx, e.exec, false, &riverdriver.LeaderElectParams{
LeaderID: e.config.ClientID,
Now: e.Time.NowUTCOrNil(),
Schema: e.config.Schema,
TTL: e.leaderTTL(),
})
Expand Down Expand Up @@ -333,6 +334,7 @@ func (e *Elector) keepLeadershipLoop(ctx context.Context) error {

reelected, err := attemptElectOrReelect(ctx, e.exec, true, &riverdriver.LeaderElectParams{
LeaderID: e.config.ClientID,
Now: e.Time.NowUTCOrNil(),
Schema: e.config.Schema,
TTL: e.leaderTTL(),
})
Expand Down Expand Up @@ -515,6 +517,7 @@ func attemptElectOrReelect(ctx context.Context, exec riverdriver.Executor, alrea

return dbutil.WithTxV(ctx, exec, func(ctx context.Context, exec riverdriver.ExecutorTx) (bool, error) {
if _, err := exec.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Now: params.Now,
Schema: params.Schema,
}); err != nil {
return false, err
Expand Down
117 changes: 78 additions & 39 deletions internal/riverinternaltest/riverdrivertest/riverdrivertest.go
Original file line number Diff line number Diff line change
Expand Up @@ -2204,37 +2204,6 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,

const leaderTTL = 10 * time.Second

t.Run("LeaderDeleteExpired", func(t *testing.T) {
Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I just moved this set of tests down so it's in alphabetical order and easier to find.

t.Parallel()

exec, _ := setup(ctx, t)

now := time.Now().UTC()

{
numDeleted, err := exec.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Schema: "",
})
require.NoError(t, err)
require.Zero(t, numDeleted)
}

_ = testfactory.Leader(ctx, t, exec, &testfactory.LeaderOpts{
ElectedAt: ptrutil.Ptr(now.Add(-2 * time.Hour)),
ExpiresAt: ptrutil.Ptr(now.Add(-1 * time.Hour)),
LeaderID: ptrutil.Ptr(clientID),
Schema: "",
})

{
numDeleted, err := exec.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Schema: "",
})
require.NoError(t, err)
require.Equal(t, 1, numDeleted)
}
})

t.Run("LeaderAttemptElect", func(t *testing.T) {
t.Parallel()

Expand All @@ -2243,8 +2212,11 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,

exec, _ := setup(ctx, t)

now := time.Now()

elected, err := exec.LeaderAttemptElect(ctx, &riverdriver.LeaderElectParams{
LeaderID: clientID,
Now: &now,
Schema: "",
TTL: leaderTTL,
})
Expand All @@ -2255,8 +2227,8 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,
Schema: "",
})
require.NoError(t, err)
require.WithinDuration(t, time.Now(), leader.ElectedAt, 100*time.Millisecond)
require.WithinDuration(t, time.Now().Add(leaderTTL), leader.ExpiresAt, 100*time.Millisecond)
require.WithinDuration(t, now, leader.ElectedAt, time.Microsecond)
require.WithinDuration(t, now.Add(leaderTTL), leader.ExpiresAt, time.Microsecond)
})

t.Run("CannotElectTwiceInARow", func(t *testing.T) {
Expand Down Expand Up @@ -2296,8 +2268,11 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,

exec, _ := setup(ctx, t)

now := time.Now()

elected, err := exec.LeaderAttemptReelect(ctx, &riverdriver.LeaderElectParams{
LeaderID: clientID,
Now: &now,
Schema: "",
TTL: leaderTTL,
})
Expand All @@ -2308,8 +2283,8 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,
Schema: "",
})
require.NoError(t, err)
require.WithinDuration(t, time.Now(), leader.ElectedAt, 100*time.Millisecond)
require.WithinDuration(t, time.Now().Add(leaderTTL), leader.ExpiresAt, 100*time.Millisecond)
require.WithinDuration(t, now, leader.ElectedAt, time.Microsecond)
require.WithinDuration(t, now.Add(leaderTTL), leader.ExpiresAt, time.Microsecond)
})

t.Run("ReelectsSameLeader", func(t *testing.T) {
Expand Down Expand Up @@ -2343,19 +2318,80 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,
})
})

t.Run("LeaderDeleteExpired", func(t *testing.T) {
t.Parallel()

t.Run("DeletesExpired", func(t *testing.T) {
t.Parallel()

exec, _ := setup(ctx, t)

now := time.Now().UTC()

{
numDeleted, err := exec.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Schema: "",
})
require.NoError(t, err)
require.Zero(t, numDeleted)
}

_ = testfactory.Leader(ctx, t, exec, &testfactory.LeaderOpts{
ElectedAt: ptrutil.Ptr(now.Add(-2 * time.Hour)),
ExpiresAt: ptrutil.Ptr(now.Add(-1 * time.Hour)),
LeaderID: ptrutil.Ptr(clientID),
Schema: "",
})

{
numDeleted, err := exec.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Schema: "",
})
require.NoError(t, err)
require.Equal(t, 1, numDeleted)
}
})

t.Run("WithInjectedNow", func(t *testing.T) {
t.Parallel()

exec, _ := setup(ctx, t)

now := time.Now().UTC()

// Elected in the future.
_ = testfactory.Leader(ctx, t, exec, &testfactory.LeaderOpts{
ElectedAt: ptrutil.Ptr(now.Add(1 * time.Hour)),
ExpiresAt: ptrutil.Ptr(now.Add(2 * time.Hour)),
LeaderID: ptrutil.Ptr(clientID),
Schema: "",
})

numDeleted, err := exec.LeaderDeleteExpired(ctx, &riverdriver.LeaderDeleteExpiredParams{
Now: ptrutil.Ptr(now.Add(2*time.Hour + 1*time.Second)),
Schema: "",
})
require.NoError(t, err)
require.Equal(t, 1, numDeleted)
})
})

t.Run("LeaderInsert", func(t *testing.T) {
t.Parallel()

exec, _ := setup(ctx, t)

now := time.Now()

leader, err := exec.LeaderInsert(ctx, &riverdriver.LeaderInsertParams{
LeaderID: clientID,
Now: &now,
Schema: "",
TTL: leaderTTL,
})
require.NoError(t, err)
require.WithinDuration(t, time.Now(), leader.ElectedAt, 500*time.Millisecond)
require.WithinDuration(t, time.Now().Add(leaderTTL), leader.ExpiresAt, 500*time.Millisecond)
require.WithinDuration(t, now, leader.ElectedAt, time.Microsecond)
require.WithinDuration(t, now.Add(leaderTTL), leader.ExpiresAt, time.Microsecond)
require.Equal(t, clientID, leader.LeaderID)
})

Expand All @@ -2364,17 +2400,20 @@ func Exercise[TTx any](ctx context.Context, t *testing.T,

exec, _ := setup(ctx, t)

now := time.Now()

_ = testfactory.Leader(ctx, t, exec, &testfactory.LeaderOpts{
LeaderID: ptrutil.Ptr(clientID),
Now: &now,
Schema: "",
})

leader, err := exec.LeaderGetElectedLeader(ctx, &riverdriver.LeaderGetElectedLeaderParams{
Schema: "",
})
require.NoError(t, err)
require.WithinDuration(t, time.Now(), leader.ElectedAt, 500*time.Millisecond)
require.WithinDuration(t, time.Now().Add(leaderTTL), leader.ExpiresAt, 500*time.Millisecond)
require.WithinDuration(t, now, leader.ElectedAt, time.Microsecond)
require.WithinDuration(t, now.Add(leaderTTL), leader.ExpiresAt, time.Microsecond)
require.Equal(t, clientID, leader.LeaderID)
})

Expand Down
3 changes: 3 additions & 0 deletions riverdriver/river_driver_interface.go
Original file line number Diff line number Diff line change
Expand Up @@ -509,6 +509,7 @@ type Leader struct {
}

type LeaderDeleteExpiredParams struct {
Now *time.Time
Schema string
}

Expand All @@ -520,12 +521,14 @@ type LeaderInsertParams struct {
ElectedAt *time.Time
ExpiresAt *time.Time
LeaderID string
Now *time.Time
Schema string
TTL time.Duration
}

type LeaderElectParams struct {
LeaderID string
Now *time.Time
Schema string
TTL time.Duration
}
Expand Down
26 changes: 15 additions & 11 deletions riverdriver/riverdatabasesql/internal/dbsqlc/river_leader.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion riverdriver/riverdatabasesql/river_database_sql_driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -601,6 +601,7 @@ func (e *Executor) JobUpdate(ctx context.Context, params *riverdriver.JobUpdateP
func (e *Executor) LeaderAttemptElect(ctx context.Context, params *riverdriver.LeaderElectParams) (bool, error) {
numElectionsWon, err := dbsqlc.New().LeaderAttemptElect(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.LeaderAttemptElectParams{
LeaderID: params.LeaderID,
Now: params.Now,
TTL: params.TTL,
})
if err != nil {
Expand All @@ -612,6 +613,7 @@ func (e *Executor) LeaderAttemptElect(ctx context.Context, params *riverdriver.L
func (e *Executor) LeaderAttemptReelect(ctx context.Context, params *riverdriver.LeaderElectParams) (bool, error) {
numElectionsWon, err := dbsqlc.New().LeaderAttemptReelect(schemaTemplateParam(ctx, params.Schema), e.dbtx, &dbsqlc.LeaderAttemptReelectParams{
LeaderID: params.LeaderID,
Now: params.Now,
TTL: params.TTL,
})
if err != nil {
Expand All @@ -621,7 +623,7 @@ func (e *Executor) LeaderAttemptReelect(ctx context.Context, params *riverdriver
}

func (e *Executor) LeaderDeleteExpired(ctx context.Context, params *riverdriver.LeaderDeleteExpiredParams) (int, error) {
numDeleted, err := dbsqlc.New().LeaderDeleteExpired(schemaTemplateParam(ctx, params.Schema), e.dbtx)
numDeleted, err := dbsqlc.New().LeaderDeleteExpired(schemaTemplateParam(ctx, params.Schema), e.dbtx, params.Now)
if err != nil {
return 0, interpretError(err)
}
Expand All @@ -641,6 +643,7 @@ func (e *Executor) LeaderInsert(ctx context.Context, params *riverdriver.LeaderI
ElectedAt: params.ElectedAt,
ExpiresAt: params.ExpiresAt,
LeaderID: params.LeaderID,
Now: params.Now,
TTL: params.TTL,
})
if err != nil {
Expand Down
12 changes: 6 additions & 6 deletions riverdriver/riverpgxv5/internal/dbsqlc/river_leader.sql
Original file line number Diff line number Diff line change
Expand Up @@ -9,22 +9,22 @@ CREATE UNLOGGED TABLE river_leader(

-- name: LeaderAttemptElect :execrows
INSERT INTO /* TEMPLATE: schema */river_leader (leader_id, elected_at, expires_at)
VALUES (@leader_id, now(), now() + @ttl::interval)
VALUES (@leader_id, coalesce(sqlc.narg('now')::timestamptz, now()), coalesce(sqlc.narg('now')::timestamptz, now()) + @ttl::interval)
ON CONFLICT (name)
DO NOTHING;

-- name: LeaderAttemptReelect :execrows
INSERT INTO /* TEMPLATE: schema */river_leader (leader_id, elected_at, expires_at)
VALUES (@leader_id, now(), now() + @ttl::interval)
VALUES (@leader_id, coalesce(sqlc.narg('now')::timestamptz, now()), coalesce(sqlc.narg('now')::timestamptz, now()) + @ttl::interval)
ON CONFLICT (name)
DO UPDATE SET
expires_at = now() + @ttl
expires_at = coalesce(sqlc.narg('now')::timestamptz, now()) + @ttl
WHERE
river_leader.leader_id = @leader_id;

-- name: LeaderDeleteExpired :execrows
DELETE FROM /* TEMPLATE: schema */river_leader
WHERE expires_at < now();
WHERE expires_at < coalesce(sqlc.narg('now')::timestamptz, now());

-- name: LeaderGetElectedLeader :one
SELECT *
Expand All @@ -36,8 +36,8 @@ INSERT INTO /* TEMPLATE: schema */river_leader(
expires_at,
leader_id
) VALUES (
coalesce(sqlc.narg('elected_at')::timestamptz, now()),
coalesce(sqlc.narg('expires_at')::timestamptz, now() + @ttl::interval),
coalesce(sqlc.narg('elected_at')::timestamptz, coalesce(sqlc.narg('now')::timestamptz, now())),
coalesce(sqlc.narg('expires_at')::timestamptz, coalesce(sqlc.narg('now')::timestamptz, now()) + @ttl::interval),
@leader_id
) RETURNING *;

Expand Down
Loading