mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-19 02:37:24 +08:00
Merge pull request #6280 from yousong/bugfix/yousong-elect
cloudcommon: elect: notifyOne on subscribe
This commit is contained in:
@@ -72,6 +72,7 @@ type Elect struct {
|
||||
config *EtcdConfig
|
||||
|
||||
mutex *sync.Mutex
|
||||
latestEv electEvent
|
||||
subscribers []chan electEvent
|
||||
}
|
||||
|
||||
@@ -104,10 +105,12 @@ func NewElect(config *EtcdConfig, key string) (*Elect, error) {
|
||||
return nil, errors.Wrap(err, "new etcd client")
|
||||
}
|
||||
elect := &Elect{
|
||||
cli: cli,
|
||||
path: config.LockPrefix + "/" + key,
|
||||
ttl: config.LockTTL,
|
||||
mutex: &sync.Mutex{},
|
||||
cli: cli,
|
||||
path: config.LockPrefix + "/" + key,
|
||||
ttl: config.LockTTL,
|
||||
|
||||
mutex: &sync.Mutex{},
|
||||
latestEv: electEventInit,
|
||||
}
|
||||
return elect, nil
|
||||
}
|
||||
@@ -173,23 +176,35 @@ func (elect *Elect) do(ctx context.Context) (*ticket, error) {
|
||||
return r, err
|
||||
}
|
||||
|
||||
func (elect *Elect) subscribe(ch chan electEvent) {
|
||||
func (elect *Elect) subscribe(ctx context.Context, ch chan electEvent) {
|
||||
elect.mutex.Lock()
|
||||
defer elect.mutex.Unlock()
|
||||
elect.subscribers = append(elect.subscribers, ch)
|
||||
if ev := elect.latestEv; ev != electEventInit {
|
||||
elect.notifyOne(ctx, ev, ch)
|
||||
}
|
||||
}
|
||||
|
||||
func (elect *Elect) notifyOne(ctx context.Context, ev electEvent, ch chan electEvent) {
|
||||
sent := false
|
||||
select {
|
||||
case ch <- ev:
|
||||
sent = true
|
||||
case <-ctx.Done():
|
||||
default:
|
||||
}
|
||||
if !sent {
|
||||
log.Errorf("elect event '%s' missed by %#v", ev, ch)
|
||||
}
|
||||
}
|
||||
|
||||
func (elect *Elect) notify(ctx context.Context, ev electEvent) {
|
||||
elect.mutex.Lock()
|
||||
defer elect.mutex.Unlock()
|
||||
|
||||
elect.latestEv = ev
|
||||
for _, ch := range elect.subscribers {
|
||||
select {
|
||||
case ch <- ev:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
elect.notifyOne(ctx, ev, ch)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -197,7 +212,7 @@ func (elect *Elect) SubscribeWithAction(ctx context.Context, onWin, onLost func(
|
||||
go func() {
|
||||
ch := make(chan electEvent, 3)
|
||||
var ev electEvent
|
||||
elect.subscribe(ch)
|
||||
elect.subscribe(ctx, ch)
|
||||
for {
|
||||
select {
|
||||
case ev = <-ch:
|
||||
@@ -229,6 +244,7 @@ type electEvent int
|
||||
const (
|
||||
electEventWin electEvent = iota
|
||||
electEventLost
|
||||
electEventInit
|
||||
)
|
||||
|
||||
func (ev electEvent) String() string {
|
||||
@@ -237,6 +253,8 @@ func (ev electEvent) String() string {
|
||||
return "win"
|
||||
case electEventLost:
|
||||
return "lost"
|
||||
case electEventInit:
|
||||
return "init"
|
||||
default:
|
||||
return "unexpected"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user