This repository has been archived by the owner on Mar 11, 2021. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 26
/
worker_lock_repository.go
105 lines (97 loc) · 3.01 KB
/
worker_lock_repository.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
package repository
import (
"context"
"database/sql"
"time"
"github.com/fabric8-services/fabric8-auth/log"
"cirello.io/pglock"
errs "github.com/pkg/errors"
)
// LockRepository the interface for the repository
type LockRepository interface {
AcquireLock(ctx context.Context, owner string, name string, opts ...pglock.ClientOption) (*pglock.Lock, error)
GetLock(ctx context.Context, name string) (*pglock.Lock, error)
}
type lockRepositoryImpl struct {
db *sql.DB
}
// NewLockRepository creates a new storage type.
func NewLockRepository(db *sql.DB) LockRepository {
return &lockRepositoryImpl{
db: db,
}
}
// AcquireLock acquires a lock with the given name for the given owner
// Returns an error if the lock could not be obtained
func (r *lockRepositoryImpl) AcquireLock(ctx context.Context, owner, name string, opts ...pglock.ClientOption) (*pglock.Lock, error) {
log.Info(ctx, map[string]interface{}{
"lock": name,
"owner": owner,
}, "acquiring lock...")
// obtain a lock to prevent other pods to perform this task
var clnOpts []pglock.ClientOption
if opts != nil {
// Useful for testing
clnOpts = opts
} else {
// Use default
clnOpts = []pglock.ClientOption{
pglock.WithCustomTable("worker_lock"),
pglock.WithLeaseDuration(30 * time.Second),
pglock.WithHeartbeatFrequency(10 * time.Second),
pglock.WithLogger(log.Logger()),
}
}
if owner != "" {
// use a specific owner name, otherwise it will be a random value (default behaviour)
clnOpts = append(clnOpts, pglock.WithOwner(owner))
}
c, err := pglock.New(r.db, clnOpts...)
if err != nil {
return nil, errs.Wrap(err, "cannot create worker lock client")
}
l, err := c.Acquire(name) // will wait until it succeeds
if err != nil {
// Try to get the name of the owner who holds the lock if any
var holdByOwner string
existingLock, nerr := c.Get(name)
if nerr != nil {
log.Error(ctx, map[string]interface{}{
"err": nerr,
"lock": name,
"owner": owner,
}, "cannot obtain the existing lock when trying to acquire a new one")
} else {
holdByOwner = existingLock.Owner()
}
return nil, errs.Wrapf(err, "cannot acquire worker lock '%s' which is currently hold by '%s'", name, holdByOwner)
}
log.Info(ctx, map[string]interface{}{
"lock": name,
"owner": owner,
}, "acquired lock")
return l, nil
}
// GetLock returns the lock object from the given name in the table without holding
// it first.
func (r *lockRepositoryImpl) GetLock(ctx context.Context, name string) (*pglock.Lock, error) {
log.Debug(ctx, map[string]interface{}{
"lock": name,
}, "obtaining existing lock...")
opts := []pglock.ClientOption{
pglock.WithCustomTable("worker_lock"),
pglock.WithLogger(log.Logger()),
}
c, err := pglock.New(r.db, opts...)
if err != nil {
return nil, errs.Wrap(err, "cannot create worker lock client")
}
l, err := c.Get(name)
if err != nil {
return nil, errs.Wrapf(err, "cannot obtain the lock '%s'", name)
}
log.Debug(ctx, map[string]interface{}{
"lock": name,
}, "obtained existing lock")
return l, nil
}