mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(mastermeta): add LegacyDataMediaType in master cluster values, and metaNode fetch it from master.
close:#22716686 Signed-off-by: true1064 <tangjingyu@oppo.com>
This commit is contained in:
parent
716aa7970a
commit
168fa0418c
@ -82,6 +82,7 @@ func formatClusterView(cv *proto.ClusterView, cn *proto.ClusterNodeInfo, cp *pro
|
||||
sb.WriteString(fmt.Sprintf(" DecommissionDiskLimit : %v\n", cv.DecommissionDiskLimit))
|
||||
sb.WriteString(fmt.Sprintf(" DpBackupTimeout : %v\n", cv.DpBackupTimeout))
|
||||
sb.WriteString(fmt.Sprintf(" ForbidWriteOpOfProtoVersion0 : %v\n", cv.ForbidWriteOpOfProtoVer0))
|
||||
sb.WriteString(fmt.Sprintf(" LegacyDataMediaType : %v\n", cv.LegacyDataMediaType))
|
||||
return sb.String()
|
||||
}
|
||||
|
||||
|
||||
@ -667,9 +667,6 @@ func (s *DataNode) register(cfg *config.Config) (err error) {
|
||||
if LocalIP == "" {
|
||||
LocalIP = string(ci.Ip)
|
||||
}
|
||||
s.nodeForbidWriteOpOfProtoVer0 = ci.ForbidWriteOpOfProtoVer0
|
||||
nodeForbidWriteOpVerMsg = fmt.Sprintf("action[registerToMaster] from master, node forbid write Operate Of proto version-0: %v",
|
||||
s.nodeForbidWriteOpOfProtoVer0)
|
||||
|
||||
s.localServerAddr = fmt.Sprintf("%s:%v", LocalIP, s.port)
|
||||
if !util.IsIPV4(LocalIP) {
|
||||
@ -680,19 +677,23 @@ func (s *DataNode) register(cfg *config.Config) (err error) {
|
||||
}
|
||||
|
||||
volListForbidWriteOpOfProtoVer0 := make([]string, 0)
|
||||
var volListForbidFromMaster *proto.VolListForbidWriteOpOfProtoVer0
|
||||
if volListForbidFromMaster, err = MasterClient.AdminAPI().GetVolListForbiddenWriteOpOfProtoVer0(); err != nil {
|
||||
var settingsForbidFromMaster *proto.UpgradeCompatibleSettings
|
||||
if settingsForbidFromMaster, err = MasterClient.AdminAPI().GetUpgradeCompatibleSettings(); err != nil {
|
||||
if strings.Contains(err.Error(), proto.KeyWordInHttpApiNotSupportErr) {
|
||||
// master may be lower version and has no this API
|
||||
volsForbidWriteOpVerMsg = fmt.Sprintf("[registerToMaster] master version has no api GetVolListForbiddenWriteOpOfProtoVer0, ues default value(false)")
|
||||
volsForbidWriteOpVerMsg = fmt.Sprintf("[registerToMaster] master version has no api GetUpgradeCompatibleSettings, ues default value(false)")
|
||||
} else {
|
||||
log.LogErrorf("[registerToMaster] failed to get volume list forbidden write op of proto version-0 from master(%v), err: %v",
|
||||
log.LogErrorf("[registerToMaster] GetUpgradeCompatibleSettings from master(%v) err: %v",
|
||||
MasterClient.Leader(), err)
|
||||
timer.Reset(2 * time.Second)
|
||||
continue
|
||||
}
|
||||
} else {
|
||||
volListForbidWriteOpOfProtoVer0 = volListForbidFromMaster.VolsForbidWriteOpOfProtoVer0
|
||||
s.nodeForbidWriteOpOfProtoVer0 = settingsForbidFromMaster.ClusterForbidWriteOpOfProtoVer0
|
||||
nodeForbidWriteOpVerMsg = fmt.Sprintf("action[registerToMaster] from master, cluster node forbid write Operate Of proto version-0: %v",
|
||||
s.nodeForbidWriteOpOfProtoVer0)
|
||||
|
||||
volListForbidWriteOpOfProtoVer0 = settingsForbidFromMaster.VolsForbidWriteOpOfProtoVer0
|
||||
volsForbidWriteOpVerMsg = fmt.Sprintf("[registerToMaster] from master, volumes forbid write operate of proto version-0: %v",
|
||||
volListForbidWriteOpOfProtoVer0)
|
||||
}
|
||||
|
||||
@ -929,6 +929,10 @@ func parseRequestToCreateVol(r *http.Request, req *createVolReq) (err error) {
|
||||
req.volStorageClass = proto.StorageClass_Unspecified
|
||||
} else if proto.IsCold(req.volType) {
|
||||
req.volStorageClass = proto.StorageClass_BlobStore
|
||||
} else {
|
||||
err = fmt.Errorf("invalid volType: %v", req.volType)
|
||||
log.LogErrorf("[parseRequestToCreateVol] err: %v", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@ -832,6 +832,7 @@ func (m *Server) getCluster(w http.ResponseWriter, r *http.Request) {
|
||||
BadPartitionIDs: make([]proto.BadPartitionView, 0),
|
||||
BadMetaPartitionIDs: make([]proto.BadPartitionView, 0),
|
||||
ForbidWriteOpOfProtoVer0: m.cluster.cfg.forbidWriteOpOfProtoVer0,
|
||||
LegacyDataMediaType: m.cluster.legacyDataMediaType,
|
||||
}
|
||||
|
||||
vols := m.cluster.allVolNames()
|
||||
@ -1020,7 +1021,7 @@ func (m *Server) getIPAddr(w http.ResponseWriter, r *http.Request) {
|
||||
doStatAndMetric(proto.AdminGetIP, metric, nil, nil)
|
||||
}()
|
||||
|
||||
m.cluster.loadClusterValue()
|
||||
m.cluster.loadClusterValue(false)
|
||||
batchCount := atomic.LoadUint64(&m.cluster.cfg.MetaNodeDeleteBatchCount)
|
||||
limitRate := atomic.LoadUint64(&m.cluster.cfg.DataNodeDeleteLimitRate)
|
||||
deleteSleepMs := atomic.LoadUint64(&m.cluster.cfg.MetaNodeDeleteWorkerSleepMs)
|
||||
@ -1037,13 +1038,12 @@ func (m *Server) getIPAddr(w http.ResponseWriter, r *http.Request) {
|
||||
DpMaxRepairErrCnt: dpMaxRepairErrCnt,
|
||||
DirChildrenNumLimit: dirChildrenNumLimit,
|
||||
// Ip: strings.Split(r.RemoteAddr, ":")[0],
|
||||
Ip: iputil.RealIP(r),
|
||||
EbsAddr: m.bStoreAddr,
|
||||
ServicePath: m.servicePath,
|
||||
ClusterUuid: m.cluster.clusterUuid,
|
||||
ClusterUuidEnable: m.cluster.clusterUuidEnable,
|
||||
ClusterEnableSnapshot: m.cluster.cfg.EnableSnapshot,
|
||||
ForbidWriteOpOfProtoVer0: m.cluster.cfg.forbidWriteOpOfProtoVer0,
|
||||
Ip: iputil.RealIP(r),
|
||||
EbsAddr: m.bStoreAddr,
|
||||
ServicePath: m.servicePath,
|
||||
ClusterUuid: m.cluster.clusterUuid,
|
||||
ClusterUuidEnable: m.cluster.clusterUuidEnable,
|
||||
ClusterEnableSnapshot: m.cluster.cfg.EnableSnapshot,
|
||||
}
|
||||
|
||||
sendOkReply(w, r, newSuccessHTTPReply(cInfo))
|
||||
@ -8153,25 +8153,35 @@ func (m *Server) volAddAllowedStorageClass(w http.ResponseWriter, r *http.Reques
|
||||
sendOkReply(w, r, newSuccessHTTPReply("success"))
|
||||
}
|
||||
|
||||
func (m *Server) getVolListForbidWriteOpOfProtoVer0(w http.ResponseWriter, r *http.Request) {
|
||||
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminGetVolListForbidWriteOpOfProtoVer0))
|
||||
defer func() {
|
||||
doStatAndMetric(proto.AdminGetVolListForbidWriteOpOfProtoVer0, metric, nil, nil)
|
||||
}()
|
||||
func (m *Server) getVolListForbidWriteOpOfProtoVer0() (volsForbidWriteOpOfProtoVer0 []string) {
|
||||
volsForbidWriteOpOfProtoVer0 = make([]string, 0)
|
||||
|
||||
volsForbidWriteOpOfProtoVer0 := make([]string, 0)
|
||||
m.cluster.volMutex.RLock()
|
||||
defer m.cluster.volMutex.RUnlock()
|
||||
for _, vol := range m.cluster.vols {
|
||||
if vol.ForbidWriteOpOfProtoVer0.Load() {
|
||||
volsForbidWriteOpOfProtoVer0 = append(volsForbidWriteOpOfProtoVer0, vol.Name)
|
||||
}
|
||||
}
|
||||
m.cluster.volMutex.RUnlock()
|
||||
|
||||
cInfo := &proto.VolListForbidWriteOpOfProtoVer0{
|
||||
VolsForbidWriteOpOfProtoVer0: volsForbidWriteOpOfProtoVer0,
|
||||
return
|
||||
}
|
||||
|
||||
func (m *Server) getUpgradeCompatibleSettings(w http.ResponseWriter, r *http.Request) {
|
||||
metric := exporter.NewTPCnt(apiToMetricsName(proto.AdminGetUpgradeCompatibleSettings))
|
||||
defer func() {
|
||||
doStatAndMetric(proto.AdminGetUpgradeCompatibleSettings, metric, nil, nil)
|
||||
}()
|
||||
|
||||
volsForbidWriteOpOfProtoVer0 := m.getVolListForbidWriteOpOfProtoVer0()
|
||||
|
||||
cInfo := &proto.UpgradeCompatibleSettings{
|
||||
VolsForbidWriteOpOfProtoVer0: volsForbidWriteOpOfProtoVer0,
|
||||
ClusterForbidWriteOpOfProtoVer0: m.cluster.cfg.forbidWriteOpOfProtoVer0,
|
||||
LegacyDataMediaType: m.cluster.legacyDataMediaType,
|
||||
}
|
||||
log.LogInfof("[getVolListForbidWriteOpOfProtoVer0] total %v, VolsForbidWriteOpOfProtoVer0: %v",
|
||||
log.LogInfof("[getUpgradeCompatibleSettings] cluster.legacyDataMediaType(%v), ClusterForbidWriteOpOfProtoVer0(%v), VolsForbidWriteOpOfProtoVer0(total %v): %v",
|
||||
cInfo.LegacyDataMediaType, cInfo.ClusterForbidWriteOpOfProtoVer0,
|
||||
len(cInfo.VolsForbidWriteOpOfProtoVer0), cInfo.VolsForbidWriteOpOfProtoVer0)
|
||||
|
||||
sendOkReply(w, r, newSuccessHTTPReply(cInfo))
|
||||
|
||||
@ -123,6 +123,7 @@ type Cluster struct {
|
||||
fileStatsEnable bool
|
||||
clusterUuidEnable bool
|
||||
authenticate bool
|
||||
legacyDataMediaType uint32
|
||||
|
||||
S3ApiQosQuota *sync.Map // (api,uid,limtType) -> limitQuota
|
||||
QosAcceptLimit *rate.Limiter
|
||||
@ -1307,15 +1308,15 @@ func (c *Cluster) addDataNode(nodeAddr, zoneName string, nodesetId uint64, media
|
||||
nodeAddr, zoneName, nodesetId, mediaType)
|
||||
|
||||
if !proto.IsValidMediaType(mediaType) {
|
||||
if !proto.IsValidMediaType(c.server.config.legacyDataMediaType) {
|
||||
err = fmt.Errorf("invalid mediaType(%v) in req when adding datanode(%v), and conf legacyDataMediaType not set",
|
||||
if !proto.IsValidMediaType(c.legacyDataMediaType) {
|
||||
err = fmt.Errorf("invalid mediaType(%v) in req when adding datanode(%v), and cluster LegacyDataMediaType not set",
|
||||
mediaType, nodeAddr)
|
||||
return
|
||||
}
|
||||
|
||||
mediaType = c.server.config.legacyDataMediaType
|
||||
log.LogWarnf("[addDataNode] adding datanode(%v), set mediaType as conf legacyDataMediaType(%v)",
|
||||
nodeAddr, proto.MediaTypeString(c.server.config.legacyDataMediaType))
|
||||
mediaType = c.legacyDataMediaType
|
||||
log.LogWarnf("[addDataNode] adding datanode(%v), set mediaType as cluster LegacyDataMediaType(%v)",
|
||||
nodeAddr, proto.MediaTypeString(c.legacyDataMediaType))
|
||||
}
|
||||
|
||||
// datanode existed
|
||||
|
||||
@ -379,8 +379,8 @@ func (m *Server) registerAPIRoutes(router *mux.Router) {
|
||||
Path(proto.AdminGetClusterMetaNodes).
|
||||
HandlerFunc(m.getAllMetaNodes)
|
||||
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
|
||||
Path(proto.AdminGetVolListForbidWriteOpOfProtoVer0).
|
||||
HandlerFunc(m.getVolListForbidWriteOpOfProtoVer0)
|
||||
Path(proto.AdminGetUpgradeCompatibleSettings).
|
||||
HandlerFunc(m.getUpgradeCompatibleSettings)
|
||||
|
||||
// volume management APIs
|
||||
router.NewRoute().Methods(http.MethodGet, http.MethodPost).
|
||||
|
||||
@ -134,6 +134,7 @@ func (m *Server) loadMetadata() {
|
||||
var autoUpdatedZones []*Zone
|
||||
var autoUpdatedLegacyVols []*Vol
|
||||
var updatedDataPartitions []*DataPartition
|
||||
var updatedClusterValue bool
|
||||
|
||||
log.LogInfo("action[loadMetadata] begin")
|
||||
syslog.Println("action[loadMetadata] begin")
|
||||
@ -141,7 +142,7 @@ func (m *Server) loadMetadata() {
|
||||
m.restoreIDAlloc()
|
||||
m.cluster.fsm.restore()
|
||||
|
||||
if err = m.cluster.loadClusterValue(); err != nil {
|
||||
if err, updatedClusterValue = m.cluster.loadClusterValue(true); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
@ -267,10 +268,18 @@ func (m *Server) loadMetadata() {
|
||||
}
|
||||
log.LogInfo("action[loadS3QoSInfo] end")
|
||||
|
||||
if updatedClusterValue {
|
||||
log.LogInfof("action[loadMetadata] clusterValue updated, persist it")
|
||||
if err = m.cluster.syncPutCluster(); err != nil {
|
||||
log.LogCriticalf("action[loadMetadata] clusterValue updated, but persist failed: %v", err.Error())
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, dn := range updatedDataNodes {
|
||||
log.LogInfof("action[loadVols] auto updated legacy datanode(%v), persist it", dn.Addr)
|
||||
log.LogInfof("action[loadMetadata] auto updated legacy datanode(%v), persist it", dn.Addr)
|
||||
if err = m.cluster.syncUpdateDataNode(dn); err != nil {
|
||||
log.LogCriticalf("action[loadVols] auto updated legacy datanode(%v), but persist failed: %v", dn.Addr, err.Error())
|
||||
log.LogCriticalf("action[loadMetadata] auto updated legacy datanode(%v), but persist failed: %v", dn.Addr, err.Error())
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
@ -278,23 +287,23 @@ func (m *Server) loadMetadata() {
|
||||
for _, z := range autoUpdatedZones {
|
||||
log.LogInfof("action[loadVols] auto updated legacy zone(%v), persist it", z.name)
|
||||
if err = m.cluster.sycnPutZoneInfo(z); err != nil {
|
||||
log.LogCriticalf("action[loadVols] auto updated legacy zone(%v), but persist failed: %v", z.name, err.Error())
|
||||
log.LogCriticalf("action[loadMetadata] auto updated legacy zone(%v), but persist failed: %v", z.name, err.Error())
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, v := range autoUpdatedLegacyVols {
|
||||
log.LogInfof("action[loadVols] auto updated legacy vol(%v), persist it", v.Name)
|
||||
log.LogInfof("action[loadMetadata] auto updated legacy vol(%v), persist it", v.Name)
|
||||
if err = m.cluster.syncUpdateVol(v); err != nil {
|
||||
log.LogCriticalf("action[loadVols] auto updated legacy vol(%v), but persist failed: %v", v.Name, err.Error())
|
||||
log.LogCriticalf("action[loadMetadata] auto updated legacy vol(%v), but persist failed: %v", v.Name, err.Error())
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, dp := range updatedDataPartitions {
|
||||
log.LogInfof("action[loadVols] auto updated legacy dataPartition(%v), persist it", dp.PartitionID)
|
||||
log.LogInfof("action[loadMetadata] auto updated legacy dataPartition(%v), persist it", dp.PartitionID)
|
||||
if err = m.cluster.syncUpdateDataPartition(dp); err != nil {
|
||||
log.LogCriticalf("action[loadVols] auto updated legacy dataPartition(%v), but persist failed: %v",
|
||||
log.LogCriticalf("action[loadMetadata] auto updated legacy dataPartition(%v), but persist failed: %v",
|
||||
dp.PartitionID, err.Error())
|
||||
panic(err)
|
||||
}
|
||||
|
||||
@ -17,6 +17,7 @@ package master
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
syslog "log"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
@ -70,6 +71,7 @@ type clusterValue struct {
|
||||
AutoDpMetaRepairParallelCnt uint32
|
||||
DataPartitionTimeoutSec int64
|
||||
ForbidWriteOpOfProtoVer0 bool
|
||||
LegacyDataMediaType uint32
|
||||
}
|
||||
|
||||
func newClusterValue(c *Cluster) (cv *clusterValue) {
|
||||
@ -109,6 +111,7 @@ func newClusterValue(c *Cluster) (cv *clusterValue) {
|
||||
AutoDpMetaRepairParallelCnt: c.AutoDpMetaRepairParallelCnt.Load(),
|
||||
DataPartitionTimeoutSec: c.getDataPartitionTimeoutSec(),
|
||||
ForbidWriteOpOfProtoVer0: c.cfg.forbidWriteOpOfProtoVer0,
|
||||
LegacyDataMediaType: c.legacyDataMediaType,
|
||||
}
|
||||
return cv
|
||||
}
|
||||
@ -1173,8 +1176,8 @@ func (c *Cluster) checkSetMediaTypeForLegacyZones() (updatedZones []*Zone) {
|
||||
}
|
||||
|
||||
if _, exists := zonesHasDatanode[zone.name]; exists {
|
||||
zone.SetDataMediaType(c.server.config.legacyDataMediaType)
|
||||
log.LogWarnf("[checkSetMediaTypeForLegacyZones] set mediaType(%v) by config legacyDataMediaType for legacy zone(%v)",
|
||||
zone.SetDataMediaType(c.legacyDataMediaType)
|
||||
log.LogWarnf("[checkSetMediaTypeForLegacyZones] set mediaType(%v) by config LegacyDataMediaType for legacy zone(%v)",
|
||||
proto.MediaTypeString(zone.dataMediaType), zone.name)
|
||||
updatedZones = append(updatedZones, zone)
|
||||
}
|
||||
@ -1226,21 +1229,51 @@ func (c *Cluster) checkPersistClusterValue() {
|
||||
log.LogInfo("action[checkPersistClusterValue] add cluster value record")
|
||||
}
|
||||
|
||||
func (c *Cluster) loadClusterValue() (err error) {
|
||||
func (c *Cluster) updateClusterLegacyDataMediaType() (updatedClusterValue bool) {
|
||||
if c.legacyDataMediaType == proto.MediaType_Unspecified {
|
||||
if proto.IsValidMediaType(c.cfg.legacyDataMediaType) {
|
||||
c.legacyDataMediaType = c.cfg.legacyDataMediaType
|
||||
updatedClusterValue = true
|
||||
|
||||
msg := fmt.Sprintf("update clusterValue LegacyDataMediaType as config: %v",
|
||||
proto.MediaTypeString(c.cfg.legacyDataMediaType))
|
||||
log.LogWarnf("[loadClusterValue] %v", msg)
|
||||
syslog.Println(msg)
|
||||
} else {
|
||||
log.LogWarnf("[loadClusterValue] clusterValue LegacyDataMediaType not set, and config LegacyDataMediaType(%v) invalid",
|
||||
c.cfg.legacyDataMediaType)
|
||||
}
|
||||
} else {
|
||||
if c.cfg.legacyDataMediaType != c.legacyDataMediaType {
|
||||
log.LogWarnf("[loadClusterValue] clusterValue persisted LegacyDataMediaType(%v) and config LegacyDataMediaType(%v) different",
|
||||
proto.MediaTypeString(c.legacyDataMediaType), proto.MediaTypeString(c.cfg.legacyDataMediaType))
|
||||
} else {
|
||||
msg := fmt.Sprintf("clusterValue persisted LegacyDataMediaType is the same with config: %v",
|
||||
proto.MediaTypeString(c.cfg.legacyDataMediaType))
|
||||
log.LogInfof("[loadClusterValue] %v", msg)
|
||||
syslog.Println(msg)
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (c *Cluster) loadClusterValue(isStartMaster bool) (err error, updatedClusterValue bool) {
|
||||
result, err := c.fsm.store.SeekForPrefix([]byte(clusterPrefix))
|
||||
if err != nil {
|
||||
err = fmt.Errorf("action[loadClusterValue],err:%v", err.Error())
|
||||
return err
|
||||
return err, false
|
||||
}
|
||||
for _, value := range result {
|
||||
cv := &clusterValue{}
|
||||
if err = json.Unmarshal(value, cv); err != nil {
|
||||
log.LogErrorf("action[loadClusterValue], unmarshal err:%v", err.Error())
|
||||
return err
|
||||
return err, false
|
||||
}
|
||||
|
||||
if cv.Name != c.Name {
|
||||
log.LogErrorf("action[loadClusterValue] loaded cluster value: %+v", cv)
|
||||
log.LogErrorf("action[loadClusterValue] clusterName(%v) not match loaded clusterName(%v), n loaded cluster value: %+v",
|
||||
c.Name, cv.Name, cv)
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
@ -1302,8 +1335,15 @@ func (c *Cluster) loadClusterValue() (err error) {
|
||||
c.updateAutoDpMetaRepairParallelCnt(cv.AutoDpMetaRepairParallelCnt)
|
||||
c.updateDataPartitionTimeoutSec(cv.DataPartitionTimeoutSec)
|
||||
c.cfg.forbidWriteOpOfProtoVer0 = cv.ForbidWriteOpOfProtoVer0
|
||||
c.legacyDataMediaType = cv.LegacyDataMediaType
|
||||
log.LogInfof("action[loadClusterValue] ForbidWriteOpOfProtoVer0(%v)", cv.ForbidWriteOpOfProtoVer0)
|
||||
}
|
||||
|
||||
if isStartMaster {
|
||||
// check need update c.LegacyDataMediaType and do persist
|
||||
updatedClusterValue = c.updateClusterLegacyDataMediaType()
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
@ -1516,9 +1556,9 @@ func (c *Cluster) loadDataNodes() (err error, updatedDataNodes []*DataNode) {
|
||||
}
|
||||
|
||||
if dnv.MediaType == proto.MediaType_Unspecified {
|
||||
dnv.MediaType = c.server.config.legacyDataMediaType
|
||||
dnv.MediaType = c.legacyDataMediaType
|
||||
updated = true
|
||||
log.LogWarnf("legacy datanode(%v), set mediaType(%v) by config legacyDataMediaType",
|
||||
log.LogWarnf("legacy datanode(%v), set mediaType(%v) by cluster LegacyDataMediaType",
|
||||
dnv.Addr, proto.MediaTypeString(dnv.MediaType))
|
||||
}
|
||||
|
||||
@ -1624,16 +1664,16 @@ func (c *Cluster) setStorageClassForLegacyVol(vv *volValue) (err error, updated
|
||||
}
|
||||
|
||||
if proto.IsHot(vv.VolType) {
|
||||
if !proto.IsValidMediaType(c.server.config.legacyDataMediaType) {
|
||||
err = fmt.Errorf("try to set hot vol(%v) volStorageClass, but config legacyDataMediaType not set", vv.Name)
|
||||
if !proto.IsValidMediaType(c.legacyDataMediaType) {
|
||||
err = fmt.Errorf("try to set hot vol(%v) volStorageClass, but cluster LegacyDataMediaType not set", vv.Name)
|
||||
return
|
||||
}
|
||||
|
||||
vv.VolStorageClass = proto.GetStorageClassByMediaType(c.server.config.legacyDataMediaType)
|
||||
vv.VolStorageClass = proto.GetStorageClassByMediaType(c.legacyDataMediaType)
|
||||
vv.AllowedStorageClass = append(vv.AllowedStorageClass, vv.VolStorageClass)
|
||||
vv.CacheDpStorageClass = proto.StorageClass_Unspecified
|
||||
updated = true
|
||||
log.LogWarnf("legacy vol(%v), set volStorageClass(%v) by config legacyDataMediaType",
|
||||
log.LogWarnf("legacy vol(%v), set volStorageClass(%v) by cluster LegacyDataMediaType",
|
||||
vv.Name, proto.StorageClassString(vv.VolStorageClass))
|
||||
} else {
|
||||
vv.VolStorageClass = proto.StorageClass_BlobStore
|
||||
@ -1645,12 +1685,12 @@ func (c *Cluster) setStorageClassForLegacyVol(vv *volValue) (err error, updated
|
||||
log.LogWarnf("legacy cold vol(%v) cacheCapacity is 0, set cacheDpStorageClass(%v)",
|
||||
vv.Name, proto.StorageClassString(vv.CacheDpStorageClass))
|
||||
} else {
|
||||
if !proto.IsValidMediaType(c.server.config.legacyDataMediaType) {
|
||||
err = fmt.Errorf("try to set cold vol(%v) CacheDpStorageClass, but config legacyDataMediaType not set", vv.Name)
|
||||
if !proto.IsValidMediaType(c.legacyDataMediaType) {
|
||||
err = fmt.Errorf("try to set cold vol(%v) CacheDpStorageClass, but cluster LegacyDataMediaType not set", vv.Name)
|
||||
return
|
||||
}
|
||||
vv.CacheDpStorageClass = proto.GetStorageClassByMediaType(c.server.config.legacyDataMediaType)
|
||||
log.LogWarnf("legacy cold vol(%v), set cacheDpStorageClass(%v) by config legacyDataMediaType",
|
||||
vv.CacheDpStorageClass = proto.GetStorageClassByMediaType(c.legacyDataMediaType)
|
||||
log.LogWarnf("legacy cold vol(%v), set cacheDpStorageClass(%v) by cluster LegacyDataMediaType",
|
||||
vv.Name, proto.StorageClassString(vv.CacheDpStorageClass))
|
||||
}
|
||||
}
|
||||
@ -1798,9 +1838,9 @@ func (c *Cluster) loadDataPartitions() (err error, updatedDataPartitions []*Data
|
||||
}
|
||||
|
||||
if dpv.MediaType == proto.MediaType_Unspecified {
|
||||
dpv.MediaType = c.server.config.legacyDataMediaType
|
||||
dpv.MediaType = c.legacyDataMediaType
|
||||
updated = true
|
||||
log.LogWarnf("legacy dataPartition(id:%v), set mediaType(%v) by config legacyDataMediaType",
|
||||
log.LogWarnf("legacy dataPartition(id:%v), set mediaType(%v) by cluster LegacyDataMediaType",
|
||||
dpv.PartitionID, proto.MediaTypeString(dpv.MediaType))
|
||||
}
|
||||
|
||||
|
||||
@ -254,8 +254,6 @@ const (
|
||||
|
||||
metaNodeDeleteBatchCountKey = "batchCount"
|
||||
configNameResolveInterval = "nameResolveInterval" // int
|
||||
|
||||
cfgLegacyStorageClass = "legacyReplicaStorageClass"
|
||||
)
|
||||
|
||||
const (
|
||||
|
||||
@ -342,24 +342,6 @@ func (m *MetaNode) parseConfig(cfg *config.Config) (err error) {
|
||||
return err
|
||||
}
|
||||
|
||||
var legacyStorageClassMsg string
|
||||
if !cfg.HasKey(cfgLegacyStorageClass) {
|
||||
legacyReplicaStorageClass = proto.StorageClass_Unspecified
|
||||
legacyStorageClassMsg = fmt.Sprintf("parseConfig: [%v] not set", cfgLegacyStorageClass)
|
||||
} else {
|
||||
err, legacyReplicaStorageClass = cfg.GetUint32(cfgLegacyStorageClass)
|
||||
if err != nil || !proto.IsValidStorageClass(legacyReplicaStorageClass) {
|
||||
err = fmt.Errorf("config [%v] invalid value: %v", cfgLegacyStorageClass, legacyReplicaStorageClass)
|
||||
log.LogErrorf("parseConfig: err:%v", err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
legacyStorageClassMsg = fmt.Sprintf("parseConfig: config[%v]: %v",
|
||||
cfgLegacyStorageClass, proto.StorageClassString(legacyReplicaStorageClass))
|
||||
}
|
||||
syslog.Println(legacyStorageClassMsg)
|
||||
log.LogInfof(legacyStorageClassMsg)
|
||||
|
||||
err = m.validConfig()
|
||||
return
|
||||
}
|
||||
@ -475,6 +457,7 @@ func (m *MetaNode) register() (err error) {
|
||||
var nodeAddress string
|
||||
var volsForbidWriteOpVerMsg string
|
||||
var nodeForbidWriteOpOfProtoVerMsg string
|
||||
var legacyReplicaStorageClassMsg string
|
||||
|
||||
for {
|
||||
tryCnt++
|
||||
@ -493,25 +476,28 @@ func (m *MetaNode) register() (err error) {
|
||||
clusterEnableSnapshot = m.clusterEnableSnapshot
|
||||
m.clusterId = gClusterInfo.Cluster
|
||||
nodeAddress = m.localAddr + ":" + m.listen
|
||||
m.nodeForbidWriteOpOfProtoVer0 = gClusterInfo.ForbidWriteOpOfProtoVer0
|
||||
nodeForbidWriteOpOfProtoVerMsg = fmt.Sprintf("[register] from master, node forbid write Operate Of proto version-0: %v",
|
||||
m.nodeForbidWriteOpOfProtoVer0)
|
||||
|
||||
volListForbidWriteOpOfProtoVer0 := make([]string, 0)
|
||||
var volListForbidFromMaster *proto.VolListForbidWriteOpOfProtoVer0
|
||||
if volListForbidFromMaster, err = getVolListForbiddenWriteOpOfProtoVer0(); err != nil {
|
||||
var settingsFromMaster *proto.UpgradeCompatibleSettings
|
||||
if settingsFromMaster, err = getUpgradeCompatibleSettings(); err != nil {
|
||||
if strings.Contains(err.Error(), proto.KeyWordInHttpApiNotSupportErr) {
|
||||
// master may be lower version and has no this API
|
||||
volsForbidWriteOpVerMsg = fmt.Sprintf("[register] master version has no api GetVolListForbiddenWriteOpOfProtoVer0, ues default value(false)")
|
||||
volsForbidWriteOpVerMsg = fmt.Sprintf("[register] master version has no api GetUpgradeCompatibleSettings, ues default value(false)")
|
||||
} else {
|
||||
log.LogErrorf("[register] tryCnt(%v), failed to get volume list forbidden write op of proto version-0 from master, err: %v", tryCnt, err)
|
||||
log.LogErrorf("[register] tryCnt(%v), GetUpgradeCompatibleSettings from master err: %v", tryCnt, err)
|
||||
time.Sleep(3 * time.Second)
|
||||
continue
|
||||
}
|
||||
} else {
|
||||
volListForbidWriteOpOfProtoVer0 = volListForbidFromMaster.VolsForbidWriteOpOfProtoVer0
|
||||
volListForbidWriteOpOfProtoVer0 = settingsFromMaster.VolsForbidWriteOpOfProtoVer0
|
||||
volsForbidWriteOpVerMsg = fmt.Sprintf("[register] from master, volumes forbid write operate of proto version-0: %v",
|
||||
volListForbidWriteOpOfProtoVer0)
|
||||
|
||||
m.nodeForbidWriteOpOfProtoVer0 = settingsFromMaster.ClusterForbidWriteOpOfProtoVer0
|
||||
nodeForbidWriteOpOfProtoVerMsg = fmt.Sprintf("[register] from master, cluster node forbid write Operate Of proto version-0: %v",
|
||||
m.nodeForbidWriteOpOfProtoVer0)
|
||||
|
||||
legacyReplicaStorageClass = proto.GetMediaTypeByStorageClass(settingsFromMaster.LegacyDataMediaType)
|
||||
}
|
||||
volMapForbidWriteOpOfProtoVer0 := make(map[string]struct{})
|
||||
for _, vol := range volListForbidWriteOpOfProtoVer0 {
|
||||
@ -529,6 +515,18 @@ func (m *MetaNode) register() (err error) {
|
||||
}
|
||||
m.nodeId = nodeID
|
||||
|
||||
if proto.IsStorageClassReplica(legacyReplicaStorageClass) {
|
||||
legacyReplicaStorageClassMsg = fmt.Sprintf("[register] from master, legacyReplicaStorageClass(%v)",
|
||||
proto.StorageClassString(legacyReplicaStorageClass))
|
||||
log.LogInfo(legacyReplicaStorageClassMsg)
|
||||
} else {
|
||||
legacyReplicaStorageClassMsg = fmt.Sprintf("[register] from master, invalid legacyReplicaStorageClass(%v)",
|
||||
settingsFromMaster.LegacyDataMediaType)
|
||||
legacyReplicaStorageClass = proto.StorageClass_Unspecified
|
||||
log.LogWarn(legacyReplicaStorageClassMsg)
|
||||
}
|
||||
syslog.Printf("%v \n", legacyReplicaStorageClassMsg)
|
||||
|
||||
log.LogInfo(nodeForbidWriteOpOfProtoVerMsg)
|
||||
syslog.Printf("%v \n", nodeForbidWriteOpOfProtoVerMsg)
|
||||
log.LogInfo(volsForbidWriteOpVerMsg)
|
||||
@ -558,7 +556,7 @@ func (m *MetaNode) RemoveConnection() {
|
||||
atomic.AddInt64(&m.connectionCnt, -1)
|
||||
}
|
||||
|
||||
func getVolListForbiddenWriteOpOfProtoVer0() (volListForbidWriteOpOfProtoVer0 *proto.VolListForbidWriteOpOfProtoVer0, err error) {
|
||||
volListForbidWriteOpOfProtoVer0, err = masterClient.AdminAPI().GetVolListForbiddenWriteOpOfProtoVer0()
|
||||
func getUpgradeCompatibleSettings() (volListForbidWriteOpOfProtoVer0 *proto.UpgradeCompatibleSettings, err error) {
|
||||
volListForbidWriteOpOfProtoVer0, err = masterClient.AdminAPI().GetUpgradeCompatibleSettings()
|
||||
return
|
||||
}
|
||||
|
||||
@ -120,7 +120,7 @@ const (
|
||||
AdminSetDiskBrokenThreshold = "/admin/setDiskBrokenThreshold"
|
||||
AdminQueryDiskBrokenThreshold = "/admin/queryDiskBrokenThreshold"
|
||||
|
||||
AdminGetVolListForbidWriteOpOfProtoVer0 = "/admin/getVolListForbidWriteOpOfProtoVer0"
|
||||
AdminGetUpgradeCompatibleSettings = "/admin/getUpgradeCompatibleSettings"
|
||||
|
||||
// graphql coonsole api
|
||||
ConsoleIQL = "/iql"
|
||||
@ -568,7 +568,6 @@ type ClusterInfo struct {
|
||||
ClusterUuid string
|
||||
ClusterUuidEnable bool
|
||||
ClusterEnableSnapshot bool
|
||||
ForbidWriteOpOfProtoVer0 bool
|
||||
}
|
||||
|
||||
// CreateDataPartitionRequest defines the request to create a data partition.
|
||||
|
||||
@ -154,6 +154,7 @@ type ClusterView struct {
|
||||
StatOfStorageClass []*StatOfStorageClass
|
||||
StatMigrateStorageClass []*StatOfStorageClass
|
||||
ForbidWriteOpOfProtoVer0 bool
|
||||
LegacyDataMediaType uint32
|
||||
}
|
||||
|
||||
// ClusterNode defines the structure of a cluster node
|
||||
@ -449,8 +450,10 @@ type DecommissionTokenStatus struct {
|
||||
RunningDisk []string
|
||||
}
|
||||
|
||||
type VolListForbidWriteOpOfProtoVer0 struct {
|
||||
VolsForbidWriteOpOfProtoVer0 []string
|
||||
type UpgradeCompatibleSettings struct {
|
||||
VolsForbidWriteOpOfProtoVer0 []string
|
||||
ClusterForbidWriteOpOfProtoVer0 bool
|
||||
LegacyDataMediaType uint32
|
||||
}
|
||||
|
||||
type VolVersionInfo struct {
|
||||
|
||||
@ -887,8 +887,8 @@ func (api *AdminAPI) ResetDataPartitionRestoreStatus(dpId uint64) (ok bool, err
|
||||
return
|
||||
}
|
||||
|
||||
func (api *AdminAPI) GetVolListForbiddenWriteOpOfProtoVer0() (volListForbidWriteOpOfProtoVer0 *proto.VolListForbidWriteOpOfProtoVer0, err error) {
|
||||
volListForbidWriteOpOfProtoVer0 = &proto.VolListForbidWriteOpOfProtoVer0{}
|
||||
err = api.mc.requestWith(volListForbidWriteOpOfProtoVer0, newRequest(get, proto.AdminGetVolListForbidWriteOpOfProtoVer0).Header(api.h))
|
||||
func (api *AdminAPI) GetUpgradeCompatibleSettings() (upgradeCompatibleSettings *proto.UpgradeCompatibleSettings, err error) {
|
||||
upgradeCompatibleSettings = &proto.UpgradeCompatibleSettings{}
|
||||
err = api.mc.requestWith(upgradeCompatibleSettings, newRequest(get, proto.AdminGetUpgradeCompatibleSettings).Header(api.h))
|
||||
return
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user