cubefs/blobstore/common/taskswitch/task_switch.go
slasher b34dfb8de0 chore(util): add default of switch, update err judgment
@formatter:off

Signed-off-by: slasher <shenjie1@oppo.com>
Signed-off-by: JasonHu520 <huzongchao@oppo.com>
2024-07-29 15:13:21 +08:00

179 lines
3.5 KiB
Go

// Copyright 2022 The CubeFS Authors.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// 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 taskswitch
import (
"context"
"errors"
"strings"
"sync"
"time"
errcode "github.com/cubefs/cubefs/blobstore/common/errors"
"github.com/cubefs/cubefs/blobstore/common/trace"
)
type ISwitcher interface {
Enabled() bool
WaitEnable()
}
const (
syncTaskStatusIntervalS = 15
SwitchOpen = "true"
SwitchClose = "false"
)
var (
ErrConflictSwitch = errors.New("switch has existed")
ErrNoSuchSwitch = errors.New("no such switch")
)
type TaskSwitch struct {
mu sync.Mutex
enabled bool
wg sync.WaitGroup
}
func newTaskSwitch() *TaskSwitch {
c := &TaskSwitch{
enabled: true,
}
c.Disable()
return c
}
func NewEnabledTaskSwitch() *TaskSwitch {
taskSwitch := newTaskSwitch()
taskSwitch.Enable()
return taskSwitch
}
func (s *TaskSwitch) Enable() {
s.mu.Lock()
defer s.mu.Unlock()
if s.enabled {
return
}
s.enabled = true
s.wg.Done()
}
func (s *TaskSwitch) Disable() {
s.mu.Lock()
defer s.mu.Unlock()
if !s.enabled {
return
}
s.enabled = false
s.wg.Add(1)
}
func (s *TaskSwitch) Enabled() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.enabled
}
func (s *TaskSwitch) WaitEnable() {
s.wg.Wait()
}
type Accessor interface {
GetConfig(ctx context.Context, key string) (value string, err error)
SetConfig(ctx context.Context, key, value string) error
}
type SwitchMgr struct {
switchs map[string]*TaskSwitch
mu sync.Mutex
accessor Accessor
}
func NewSwitchMgr(accessor Accessor) *SwitchMgr {
sm := SwitchMgr{
switchs: make(map[string]*TaskSwitch),
accessor: accessor,
}
go sm.loopUpdate()
return &sm
}
func (sm *SwitchMgr) loopUpdate() {
for {
sm.update()
time.Sleep(syncTaskStatusIntervalS * time.Second)
}
}
func (sm *SwitchMgr) update() {
sm.mu.Lock()
defer sm.mu.Unlock()
span, ctx := trace.StartSpanFromContext(context.Background(), "")
for switchName, taskSwitch := range sm.switchs {
statusStr, err := sm.accessor.GetConfig(ctx, switchName)
if err != nil {
span.Errorf("Get Fail switchName %s err %v", switchName, err)
if strings.Contains(err.Error(), errcode.ErrNotFound.Error()) {
if err = sm.accessor.SetConfig(ctx, switchName, SwitchClose); err != nil {
span.Errorf("Set Fail switchName %s err %v", switchName, err)
}
}
continue
}
if switchStatus(statusStr) {
taskSwitch.Enable()
continue
}
taskSwitch.Disable()
}
}
func (sm *SwitchMgr) AddSwitch(switchName string) (*TaskSwitch, error) {
sm.mu.Lock()
defer sm.mu.Unlock()
if _, ok := sm.switchs[switchName]; ok {
return nil, ErrConflictSwitch
}
sm.switchs[switchName] = newTaskSwitch()
return sm.switchs[switchName], nil
}
func (sm *SwitchMgr) DelSwitch(switchName string) error {
sm.mu.Lock()
defer sm.mu.Unlock()
if _, ok := sm.switchs[switchName]; ok {
delete(sm.switchs, switchName)
return nil
}
return ErrNoSuchSwitch
}
func switchStatus(statusStr string) (open bool) {
switch statusStr {
case SwitchOpen:
return true
case SwitchClose:
return false
default:
return false
}
}