forked from openshift/assisted-service
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathleaderelector_posgress.go
More file actions
73 lines (62 loc) · 1.57 KB
/
Copy pathleaderelector_posgress.go
File metadata and controls
73 lines (62 loc) · 1.57 KB
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
package leader
import (
"cirello.io/pglock"
"context"
"errors"
"github.com/jinzhu/gorm"
"github.com/lib/pq"
"github.com/sirupsen/logrus"
)
type DbElector struct {
log logrus.FieldLogger
db *gorm.DB
config Config
isLeader bool
leaderName string
}
func NewDbElector(db *gorm.DB, config Config, leaderName string, logger logrus.FieldLogger) *DbElector {
return &DbElector{db: db, log: logger, config: config, isLeader: false, leaderName: leaderName}
}
func (l *DbElector) IsLeader() bool {
return l.isLeader
}
func (l *DbElector) StartLeaderElection(ctx context.Context) error {
c, err := pglock.New(l.db.DB(),
pglock.WithLeaseDuration(l.config.LeaseDuration),
pglock.WithHeartbeatFrequency(l.config.RetryInterval),
pglock.WithCustomTable(l.leaderName))
if err != nil {
l.log.WithError(err).Error("Failed to create db lock")
}
err = c.CreateTable()
if err != nil {
if p, ok := errors.Unwrap(err).(*pq.Error); !ok || p.Code.Name() != "duplicate_table" {
l.log.WithError(err).Infof("CCCCCCCCCCCCCCCCCCCCCCCCCCC")
return err
}
}
go func() {
var lock *pglock.Lock
defer func() {
if lock != nil {
lock.Close()
}
}()
for {
if ctx.Err() != nil {
return
}
l.log.Infof("BBBBBBBBBBBBBBBBBBBBBBBBB")
err = c.Do(ctx, l.leaderName, l.locked)
l.log.WithError(err).Infof("AAAAAAAAAAAAAAAAAAAAAAAA")
}
}()
return nil
}
func (l *DbElector) locked(ctx context.Context, lock *pglock.Lock) error{
l.log.Infof("GGGGGGGGGGGGGGGGGGGGGGGGGGGG")
l.isLeader = true
<-ctx.Done()
l.isLeader = false
return nil
}