коммит произвёл
GitHub
родитель
a3db701ccc
Коммит
f2aafac0c9
@@ -34,54 +34,54 @@ func (a *App) NewClusterDiscoveryService() *ClusterDiscoveryService {
|
|||||||
return a.Srv().NewClusterDiscoveryService()
|
return a.Srv().NewClusterDiscoveryService()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (me *ClusterDiscoveryService) Start() {
|
func (cds *ClusterDiscoveryService) Start() {
|
||||||
err := me.srv.Store.ClusterDiscovery().Cleanup()
|
err := cds.srv.Store.ClusterDiscovery().Cleanup()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
mlog.Error("ClusterDiscoveryService failed to cleanup the outdated cluster discovery information", mlog.Err(err))
|
mlog.Error("ClusterDiscoveryService failed to cleanup the outdated cluster discovery information", mlog.Err(err))
|
||||||
}
|
}
|
||||||
|
|
||||||
exists, err := me.srv.Store.ClusterDiscovery().Exists(&me.ClusterDiscovery)
|
exists, err := cds.srv.Store.ClusterDiscovery().Exists(&cds.ClusterDiscovery)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
mlog.Error("ClusterDiscoveryService failed to check if row exists", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()), mlog.Err(err))
|
mlog.Error("ClusterDiscoveryService failed to check if row exists", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()), mlog.Err(err))
|
||||||
} else {
|
} else {
|
||||||
if exists {
|
if exists {
|
||||||
if _, err := me.srv.Store.ClusterDiscovery().Delete(&me.ClusterDiscovery); err != nil {
|
if _, err := cds.srv.Store.ClusterDiscovery().Delete(&cds.ClusterDiscovery); err != nil {
|
||||||
mlog.Error("ClusterDiscoveryService failed to start clean", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()), mlog.Err(err))
|
mlog.Error("ClusterDiscoveryService failed to start clean", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()), mlog.Err(err))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := me.srv.Store.ClusterDiscovery().Save(&me.ClusterDiscovery); err != nil {
|
if err := cds.srv.Store.ClusterDiscovery().Save(&cds.ClusterDiscovery); err != nil {
|
||||||
mlog.Error("ClusterDiscoveryService failed to save", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()), mlog.Err(err))
|
mlog.Error("ClusterDiscoveryService failed to save", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()), mlog.Err(err))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
mlog.Debug("ClusterDiscoveryService ping writer started", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()))
|
mlog.Debug("ClusterDiscoveryService ping writer started", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()))
|
||||||
ticker := time.NewTicker(DISCOVERY_SERVICE_WRITE_PING)
|
ticker := time.NewTicker(DISCOVERY_SERVICE_WRITE_PING)
|
||||||
defer func() {
|
defer func() {
|
||||||
ticker.Stop()
|
ticker.Stop()
|
||||||
if _, err := me.srv.Store.ClusterDiscovery().Delete(&me.ClusterDiscovery); err != nil {
|
if _, err := cds.srv.Store.ClusterDiscovery().Delete(&cds.ClusterDiscovery); err != nil {
|
||||||
mlog.Error("ClusterDiscoveryService failed to cleanup", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()), mlog.Err(err))
|
mlog.Error("ClusterDiscoveryService failed to cleanup", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()), mlog.Err(err))
|
||||||
}
|
}
|
||||||
mlog.Debug("ClusterDiscoveryService ping writer stopped", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()))
|
mlog.Debug("ClusterDiscoveryService ping writer stopped", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()))
|
||||||
}()
|
}()
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
if err := me.srv.Store.ClusterDiscovery().SetLastPingAt(&me.ClusterDiscovery); err != nil {
|
if err := cds.srv.Store.ClusterDiscovery().SetLastPingAt(&cds.ClusterDiscovery); err != nil {
|
||||||
mlog.Error("ClusterDiscoveryService failed to write ping", mlog.String("ClusterDiscovery", me.ClusterDiscovery.ToJson()), mlog.Err(err))
|
mlog.Error("ClusterDiscoveryService failed to write ping", mlog.String("ClusterDiscovery", cds.ClusterDiscovery.ToJson()), mlog.Err(err))
|
||||||
}
|
}
|
||||||
case <-me.stop:
|
case <-cds.stop:
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (me *ClusterDiscoveryService) Stop() {
|
func (cds *ClusterDiscoveryService) Stop() {
|
||||||
me.stop <- true
|
cds.stop <- true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) IsLeader() bool {
|
func (s *Server) IsLeader() bool {
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user