-
Notifications
You must be signed in to change notification settings - Fork 1
/
prophet_leader.go
65 lines (52 loc) · 1.33 KB
/
prophet_leader.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
package prophet
import (
"github.com/deepfabric/prophet/util"
)
func (p *defaultProphet) startLeaderLoop() {
go p.member.ElectionLoop(p.ctx)
<-p.completeC
p.client = NewClient(p.cfg.Adapter,
WithRPCTimeout(p.cfg.RPCTimeout.Duration),
WithLeaderGetter(p.GetLeader))
}
func (p *defaultProphet) enableLeader() error {
util.GetLogger().Infof("********%s become to leader now********", p.cfg.Name)
if err := p.createRaftCluster(); err != nil {
util.GetLogger().Errorf("create raft cluster failed with %+v", err)
return err
}
p.createEventNotifer()
p.notifyElectionComplete()
p.cfg.Handler.ProphetBecomeLeader()
return nil
}
func (p *defaultProphet) disableLeader() error {
util.GetLogger().Infof("********%s become to follower now********", p.cfg.Name)
p.stopRaftCluster()
p.notifyElectionComplete()
p.cfg.Handler.ProphetBecomeFollower()
return nil
}
func (p *defaultProphet) notifyElectionComplete() {
p.notifyOnce.Do(func() {
close(p.completeC)
})
}
func (p *defaultProphet) createRaftCluster() error {
if p.cluster.IsRunning() {
return nil
}
return p.cluster.Start(p)
}
func (p *defaultProphet) stopRaftCluster() {
p.cluster.Stop()
}
func (p *defaultProphet) createEventNotifer() {
p.wn = newWatcherNotifier(p.cluster)
p.wn.start()
}
func (p *defaultProphet) stopEventNotifer() {
if p.wn != nil {
p.wn.stop()
}
}