代码拉取完成,页面将自动刷新
/**
* Tencent is pleased to support the open source community by making polaris-go available.
*
* Copyright (C) 2019 THL A29 Limited, a Tencent company. All rights reserved.
*
* Licensed under the BSD 3-Clause License (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* https://opensource.org/licenses/BSD-3-Clause
*
* Unless required by applicable law or agreed to in writing, software distributed
* under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR
* CONDITIONS OF ANY KIND, either express or implied. See the License for the
* specific language governing permissions and limitations under the License.
*/
package registerstate
import (
"context"
"fmt"
"sync"
"time"
"gitee.com/meng_mengs_boys/polaris-go/pkg/log"
"gitee.com/meng_mengs_boys/polaris-go/pkg/model"
)
type (
registerFunc func(instance *model.InstanceRegisterRequest, header map[string]string) (*model.InstanceRegisterResponse, error)
heartbeatFunc func(instance *model.InstanceHeartbeatRequest) error
)
const (
_maxHeartbeatErrorCount = 2
_headerKeyAsyncRegis = "async-regis"
_headerValueAsyncRegis = "true"
)
func NewRegisterStateManager(minRegisterInterval time.Duration) *RegisterStateManager {
return &RegisterStateManager{
minRegisterInterval: minRegisterInterval,
states: map[string]*registerState{},
}
}
type RegisterStateManager struct {
mu sync.RWMutex
minRegisterInterval time.Duration
states map[string]*registerState
}
type registerState struct {
instance *model.InstanceRegisterRequest
lastRegisterTime time.Time
cancel context.CancelFunc
}
func (c *RegisterStateManager) Destroy() {
c.mu.Lock()
pre := c.states
c.states = make(map[string]*registerState)
c.mu.Unlock()
for _, state := range pre {
state.cancel()
}
}
func (c *RegisterStateManager) PutRegister(instance *model.InstanceRegisterRequest, regis registerFunc, beat heartbeatFunc) (*registerState, bool) {
key := buildRegisterStateKey(instance.Namespace, instance.Service, instance.Host, instance.Port)
c.mu.Lock()
defer c.mu.Unlock()
_, ok := c.states[key]
if ok {
return nil, false
}
ctx, cancel := context.WithCancel(context.Background())
state := ®isterState{
instance: instance,
lastRegisterTime: time.Now(),
cancel: cancel,
}
c.states[key] = state
go c.runHeartbeat(ctx, state, regis, beat)
return state, true
}
func (c *RegisterStateManager) RemoveRegister(instance *model.InstanceDeRegisterRequest) {
key := buildRegisterStateKey(instance.Namespace, instance.Service, instance.Host, instance.Port)
c.mu.Lock()
defer c.mu.Unlock()
state, ok := c.states[key]
if ok {
state.cancel()
delete(c.states, key)
}
}
func buildRegisterStateKey(namespace string, service string, host string, port int) string {
return fmt.Sprintf("%s##%s##%s##%d", namespace, service, host, port)
}
func (c *RegisterStateManager) runHeartbeat(ctx context.Context, state *registerState, regis registerFunc, beat heartbeatFunc) {
instance := state.instance
log.GetBaseLogger().Infof("[Provider][Heartbeat] instance heartbeat task started {%s, %s, %s:%d}",
instance.Namespace, instance.Service, instance.Host, instance.Port)
ticker := time.NewTicker(time.Duration(*instance.TTL) * time.Second)
defer ticker.Stop()
errCnt := 0
minInterval := c.minRegisterInterval
for {
select {
case <-ctx.Done():
log.GetBaseLogger().Infof("[Provider][Heartbeat] instance heartbeat task stopped {%s, %s, %s:%d}",
instance.Namespace, instance.Service, instance.Host, instance.Port)
return
case <-ticker.C:
hbReq := &model.InstanceHeartbeatRequest{
Namespace: instance.Namespace,
Service: instance.Service,
Host: instance.Host,
Port: instance.Port,
ServiceToken: instance.ServiceToken,
InstanceID: instance.InstanceId,
}
start := time.Now()
if err := beat(hbReq); err != nil {
log.GetBaseLogger().Errorf("[Provider][Heartbeat] heartbeat failed {%s, %s, %s:%d}",
instance.Namespace, instance.Service, instance.Host, instance.Port, err)
errCnt++
needRegis := errCnt > _maxHeartbeatErrorCount && time.Since(state.lastRegisterTime) > minInterval
if needRegis {
// 重新记录注册的时间
state.lastRegisterTime = time.Now()
_, err = regis(instance, CreateRegisterV2Header())
if err == nil {
log.GetBaseLogger().Infof("[Provider][Heartbeat] re-register instatnce success {%s, %s, %s:%d}",
instance.Namespace, instance.Service, instance.Host, instance.Port)
} else {
log.GetBaseLogger().Warnf("[Provider][Heartbeat] re-register instatnce failed {%s, %s, %s:%d}",
instance.Namespace, instance.Service, instance.Host, instance.Port, err)
}
}
break
}
log.GetBaseLogger().Debugf("[Provider][Heartbeat] success {%s, %s, %s:%d} cost:%d ms",
instance.Namespace, instance.Service, instance.Host, instance.Port, time.Since(start).Milliseconds())
errCnt = 0
break
}
}
}
func CreateRegisterV2Header() map[string]string {
header := map[string]string{
_headerKeyAsyncRegis: _headerValueAsyncRegis,
}
return header
}
此处可能存在不合适展示的内容,页面不予展示。您可通过相关编辑功能自查并修改。
如您确认内容无涉及 不当用语 / 纯广告导流 / 暴力 / 低俗色情 / 侵权 / 盗版 / 虚假 / 无价值内容或违法国家有关法律法规的内容,可点击提交进行申诉,我们将尽快为您处理。