diff --git a/blobstore/cli/rpc2.go b/blobstore/cli/rpc2.go index a75d13b68..ee90ed2a8 100644 --- a/blobstore/cli/rpc2.go +++ b/blobstore/cli/rpc2.go @@ -33,6 +33,12 @@ var types = map[string]func() rpc2.Codec{ "shardnode.ShardnodeTaskStatsArgs": func() rpc2.Codec { return new(shardnode.ShardnodeTaskStatsArgs) }, "shardnode.ShardnodeTaskStatsRet": func() rpc2.Codec { return new(shardnode.ShardnodeTaskStatsRet) }, + "shardnode.ListShardArgs": func() rpc2.Codec { return new(shardnode.ListShardArgs) }, + "shardnode.ListShardRet": func() rpc2.Codec { return new(shardnode.ListShardRet) }, + + "shardnode.ListVolumeArgs": func() rpc2.Codec { return new(shardnode.ListVolumeArgs) }, + "shardnode.ListVolumeRet": func() rpc2.Codec { return new(shardnode.ListVolumeRet) }, + "nil": func() rpc2.Codec { return nil }, "rpc2.NoParameter": func() rpc2.Codec { return rpc2.NoParameter }, } diff --git a/blobstore/cli/shardnode/recover.go b/blobstore/cli/shardnode/recover.go new file mode 100644 index 000000000..4bc8a4fb3 --- /dev/null +++ b/blobstore/cli/shardnode/recover.go @@ -0,0 +1,525 @@ +// Copyright 2026 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 shardnode + +import ( + "bytes" + "encoding/json" + "math" + "os" + "time" + + "github.com/desertbit/grumble" + + "github.com/cubefs/cubefs/blobstore/api/clustermgr" + "github.com/cubefs/cubefs/blobstore/api/shardnode" + "github.com/cubefs/cubefs/blobstore/cli/common" + "github.com/cubefs/cubefs/blobstore/cli/common/args" + "github.com/cubefs/cubefs/blobstore/cli/common/fmt" + kvstore "github.com/cubefs/cubefs/blobstore/common/kvstorev2" + "github.com/cubefs/cubefs/blobstore/common/proto" + "github.com/cubefs/cubefs/blobstore/common/raft" + "github.com/cubefs/cubefs/blobstore/common/rpc2" + "github.com/cubefs/cubefs/blobstore/shardnode/storage" + "github.com/cubefs/cubefs/blobstore/util/errors" +) + +func addCmdRecover(cmd *grumble.Command) { + recoverCommand := &grumble.Command{ + Name: "recover", + Help: "recover shardnode, directly recover from db", + } + cmd.AddCommand(recoverCommand) + + // update shard info + recoverCommand.AddCommand(&grumble.Command{ + Name: "updateShardInfo", + Help: "update shard info in storage, should stop shardnode first", + Args: func(a *grumble.Args) { + args.SuidRegister(a) + }, + Flags: func(f *grumble.Flags) { + f.StringL("path", "", "raft storage path") + f.StringL("json", "", "shardInfo json") + }, + Run: cmdUpdateShardInfo, + }) + + // shard data backup + recoverCommand.AddCommand(&grumble.Command{ + Name: "backupShardData", + Help: "backup shard data from storage, should stop shardnode first", + Args: func(a *grumble.Args) { + args.SuidRegister(a) + }, + Flags: func(f *grumble.Flags) { + f.StringL("path", "", "origin storage path") + }, + Run: cmdShardDataBackUp, + }) + + // shard data recover + recoverCommand.AddCommand(&grumble.Command{ + Name: "recoverShardData", + Help: "recover shard data from storage, should stop shardnode first", + Args: func(a *grumble.Args) { + args.SuidRegister(a) + }, + Flags: func(f *grumble.Flags) { + f.StringL("source_path", "", "origin storage path") + f.StringL("dest_path", "", "dest storage path") + }, + Run: cmdShardDataRecover, + }) + + // shard data clear + recoverCommand.AddCommand(&grumble.Command{ + Name: "clearShardData", + Help: "clear shard data in storage", + Args: func(a *grumble.Args) { + args.SuidRegister(a) + }, + Flags: func(f *grumble.Flags) { + f.StringL("path", "", "origin storage path") + }, + Run: cmdShardDataClear, + }) + + // recover disk shards: clear shardnode data and recover by diskID and shardID + recoverCommand.AddCommand(&grumble.Command{ + Name: "recoverDiskShard", + Help: "recover disk shard after clear shard data in storage", + Args: func(a *grumble.Args) { + args.DiskIDRegister(a) + }, + Run: cmdRecoverDiskShard, + }) +} + +func cmdUpdateShardInfo(c *grumble.Context) error { + ctx := common.CmdContext() + path := c.Flags.String("path") + suid := args.Suid(c.Args) + jsonInfo := c.Flags.String("json") + col := kvstore.CF("data") + + store, err := kvstore.NewKVStore(ctx, path, kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{col}}) + if err != nil { + return err + } + defer store.Close() + + g := storage.NewShardKeysGenerator(suid) + infoKey := g.EncodeShardInfoKey() + raw, err := store.GetRaw(ctx, col, infoKey) + if err != nil && !errors.Is(err, kvstore.ErrNotFound) { + return err + } + if err == nil { + sd := clustermgr.Shard{} + err = sd.Unmarshal(raw) + if err != nil { + return err + } + fmt.Println("old shard info: ", common.Readable(sd)) + } + + // update + if len(jsonInfo) == 0 { + return nil + } + newInfo := &clustermgr.Shard{} + if err = json.Unmarshal([]byte(jsonInfo), newInfo); err != nil { + return err + } + fmt.Println("new shard info: ", common.Readable(newInfo)) + + _raw, err := newInfo.Marshal() + if err != nil { + return err + } + err = store.SetRaw(ctx, col, infoKey, _raw) + if err != nil { + return err + } + store.FlushCF(ctx, col) + return nil +} + +func cmdShardDataBackUp(c *grumble.Context) error { + ctx := common.CmdContext() + path := c.Flags.String("path") + suid := args.Suid(c.Args) + + colData := kvstore.CF("data") + colRaft := kvstore.CF("raft-wal") + + backPath := path + "/shard_" + suid.ShardID().ToString() + "_backup" + if err := os.Mkdir(backPath, 0o755); err != nil { + return err + } + fmt.Println("create backup dir: ", backPath) + + backKVStore, err := kvstore.NewKVStore(ctx, backPath+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ + ColumnFamily: []kvstore.CF{colData}, + CreateIfMissing: true, + }) + if err != nil { + return err + } + defer backKVStore.Close() + + backRaftStore, err := kvstore.NewKVStore(ctx, backPath+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ + ColumnFamily: []kvstore.CF{colRaft}, + CreateIfMissing: true, + }) + if err != nil { + return err + } + defer backRaftStore.Close() + fmt.Println("backup store opened successfully") + + fmt.Println("start backup data") + originKVStore, err := kvstore.NewKVStore(ctx, path+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) + if err != nil { + return err + } + defer originKVStore.Close() + fmt.Println("origin store opened successfully") + + g := storage.NewShardKeysGenerator(suid) + shardDataPrefix := g.EncodeShardDataPrefix() + dataList := originKVStore.List(ctx, colData, shardDataPrefix, nil, nil) + kvCount := 0 + for { + kg, vg, err := dataList.ReadNext() + if err != nil { + return err + } + if kg == nil || vg == nil { + fmt.Println("list to end") + break + } + if !bytes.HasPrefix(kg.Key(), shardDataPrefix) { + return errors.New("key prefix not match") + } + if err = backKVStore.SetRaw(ctx, colData, kg.Key(), vg.Value()); err != nil { + return err + } + kg.Close() + vg.Close() + kvCount++ + } + fmt.Println("shard kv data backup done, num: ", kvCount) + + originRaftStore, err := kvstore.NewKVStore(ctx, path+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) + if err != nil { + return err + } + defer originRaftStore.Close() + + raftLogPrefix := raft.EncodeIndexLogKeyPrefix(uint64(suid.ShardID())) + raftLogList := originRaftStore.List(ctx, colRaft, raftLogPrefix, nil, nil) + logCount := 0 + for { + kg, vg, err := raftLogList.ReadNext() + if err != nil { + return err + } + if kg == nil || vg == nil { + fmt.Println("list to end") + break + } + if !bytes.HasPrefix(kg.Key(), raftLogPrefix) { + return errors.New("key prefix not match") + } + if err = backRaftStore.SetRaw(ctx, colRaft, kg.Key(), vg.Value()); err != nil { + return err + } + kg.Close() + vg.Close() + logCount++ + } + fmt.Println("shard raft log backup done, num: ", logCount) + + hardStateKey := raft.EncodeHardStateKey(uint64(suid.ShardID())) + hsRaw, err := originRaftStore.GetRaw(ctx, colRaft, hardStateKey) + if err != nil { + return err + } + if err = backRaftStore.SetRaw(ctx, colRaft, hardStateKey, hsRaw); err != nil { + return err + } + fmt.Println("hardState backup done") + return nil +} + +func cmdShardDataRecover(c *grumble.Context) error { + ctx := common.CmdContext() + srcPath := c.Flags.String("source_path") + destPath := c.Flags.String("dest_path") + suid := args.Suid(c.Args) + + colData := kvstore.CF("data") + colRaft := kvstore.CF("raft-wal") + + kvStore, err := kvstore.NewKVStore(ctx, srcPath+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) + if err != nil { + return err + } + defer kvStore.Close() + + destKVStore, err := kvstore.NewKVStore(ctx, destPath+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ + ColumnFamily: []kvstore.CF{colData}, + CreateIfMissing: true, + }) + if err != nil { + return err + } + defer destKVStore.Close() + + fmt.Println("start recover shard kv data") + g := storage.NewShardKeysGenerator(suid) + shardDataPrefix := g.EncodeShardDataPrefix() + dataList := kvStore.List(ctx, colData, shardDataPrefix, nil, nil) + kvCount := 0 + for { + kg, vg, err := dataList.ReadNext() + if err != nil { + return err + } + if kg == nil || vg == nil { + fmt.Println("list to end") + break + } + if !bytes.HasPrefix(kg.Key(), shardDataPrefix) { + return errors.New("key prefix not match") + } + if err = destKVStore.SetRaw(ctx, colData, kg.Key(), vg.Value()); err != nil { + return err + } + kg.Close() + vg.Close() + kvCount++ + } + fmt.Println("shard kv data recover done, num: ", kvCount) + + raftStore, err := kvstore.NewKVStore(ctx, srcPath+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) + if err != nil { + return err + } + defer raftStore.Close() + + destRaftStore, err := kvstore.NewKVStore(ctx, destPath+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ + ColumnFamily: []kvstore.CF{colRaft}, + CreateIfMissing: true, + }) + if err != nil { + return err + } + defer destRaftStore.Close() + fmt.Println("start recover shard raft data") + + raftLogPrefix := raft.EncodeIndexLogKeyPrefix(uint64(suid.ShardID())) + raftLogList := raftStore.List(ctx, colRaft, raftLogPrefix, nil, nil) + logCount := 0 + for { + kg, vg, err := raftLogList.ReadNext() + if err != nil { + return err + } + if kg == nil || vg == nil { + fmt.Println("list to end") + break + } + if !bytes.HasPrefix(kg.Key(), raftLogPrefix) { + return errors.New("key prefix not match") + } + if err = destRaftStore.SetRaw(ctx, colRaft, kg.Key(), vg.Value()); err != nil { + return err + } + kg.Close() + vg.Close() + logCount++ + } + fmt.Println("shard raft log recover, num: ", logCount) + + hardStateKey := raft.EncodeHardStateKey(uint64(suid.ShardID())) + hsRaw, err := raftStore.GetRaw(ctx, colRaft, hardStateKey) + if err != nil { + return err + } + if err = destRaftStore.SetRaw(ctx, colRaft, hardStateKey, hsRaw); err != nil { + return err + } + fmt.Println("hardState recover done") + return nil +} + +func cmdShardDataClear(c *grumble.Context) error { + ctx := common.CmdContext() + path := c.Flags.String("path") + suid := args.Suid(c.Args) + + colData := kvstore.CF("data") + colRaft := kvstore.CF("raft-wal") + + dataStore, err := kvstore.NewKVStore(ctx, path+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) + if err != nil { + return err + } + defer dataStore.Close() + g := storage.NewShardKeysGenerator(suid) + shardInfoKey := g.EncodeShardInfoKey() + if err = dataStore.Delete(ctx, colData, shardInfoKey); err != nil { + return err + } + fmt.Println("shard info deleted") + + shardDataPrefix := g.EncodeShardDataPrefix() + shardDataMaxPrefix := g.EncodeShardDataMaxPrefix() + if err = dataStore.DeleteRange(ctx, colData, shardDataPrefix, shardDataMaxPrefix); err != nil { + return err + } + dataStore.FlushCF(ctx, colData) + fmt.Println("shard data deleted") + + raftStore, err := kvstore.NewKVStore(ctx, path+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) + if err != nil { + return err + } + + hardStateKey := raft.EncodeHardStateKey(uint64(suid.ShardID())) + if err = raftStore.Delete(ctx, colRaft, hardStateKey); err != nil { + return err + } + fmt.Println("hardState deleted") + + snapShotMetaKey := raft.EncodeSnapshotMetaKey(uint64(suid.ShardID())) + if err = raftStore.Delete(ctx, colRaft, snapShotMetaKey); err != nil { + return err + } + fmt.Println("snapShot meta deleted") + + if err = raftStore.DeleteRange(ctx, colRaft, + raft.EncodeIndexLogKey(uint64(suid.ShardID()), 0), + raft.EncodeIndexLogKey(uint64(suid.ShardID()), math.MaxUint64)); err != nil { + return err + } + raftStore.FlushCF(ctx, colRaft) + fmt.Println("raft log deleted") + return nil +} + +func cmdRecoverDiskShard(c *grumble.Context) error { + ctx := common.CmdContext() + cmClient := newCMClient(c.Flags) + diskID := args.DiskID(c.Args) + shardID := args.ShardID(c.Args) + var shards []clustermgr.ShardUnitInfo + var err error + if shardID != proto.InvalidShardID { + info, err := cmClient.GetShardInfo(ctx, &clustermgr.GetShardArgs{ + ShardID: proto.ShardID(shardID), + }) + if err != nil { + return errors.Info(err, "get single shard failed") + } + shards = []clustermgr.ShardUnitInfo{ + { + Suid: info.Units[0].Suid, + DiskID: info.Units[0].DiskID, + Host: info.Units[0].Host, + Learner: info.Units[0].Learner, + }, + } + } else { + shards, err = cmClient.ListShardUnit(ctx, &clustermgr.ListShardUnitArgs{ + DiskID: diskID, + }) + if err != nil { + return errors.Info(err, "list shard unit failed") + } + } + + snClient := shardnode.New(rpc2.Client{}) + for _, sd := range shards { + shard, err := cmClient.GetShardInfo(ctx, &clustermgr.GetShardArgs{ + ShardID: sd.Suid.ShardID(), + }) + if err != nil { + return errors.Info(err, "get shard failed") + } + + var _err error + var host string + for _, u := range shard.Units { + if u.Suid == sd.Suid { + host = u.Host + continue + } + // remove from raft group + if _err = snClient.UpdateShard(ctx, u.Host, shardnode.UpdateShardArgs{ + DiskID: u.DiskID, + Suid: u.Suid, + ShardUpdateType: proto.ShardUpdateTypeRemoveMember, + Unit: clustermgr.ShardUnit{ + Suid: sd.Suid, + DiskID: sd.DiskID, + Learner: sd.Learner, + }, + }); _err != nil { + continue + } + break + } + if _err != nil { + return errors.Info(_err, "remove shard failed") + } + + time.Sleep(5 * time.Second) + if _err = snClient.AddShard(ctx, host, shardnode.AddShardArgs{ + DiskID: sd.DiskID, + Suid: sd.Suid, + Range: sd.Range, + Units: shard.Units, + RouteVersion: sd.RouteVersion, + }); _err != nil { + return errors.Info(err, "add shard failed") + } + for _, u := range shard.Units { + if u.Suid == sd.Suid { + continue + } + if _err = snClient.UpdateShard(ctx, u.Host, shardnode.UpdateShardArgs{ + DiskID: u.DiskID, + Suid: u.Suid, + ShardUpdateType: proto.ShardUpdateTypeAddMember, + Unit: clustermgr.ShardUnit{ + Suid: sd.Suid, + DiskID: sd.DiskID, + }, + }); _err != nil { + continue + } + break + } + if _err != nil { + return errors.Info(_err, "add shard failed") + } + fmt.Println("add shard success, suid:", sd.Suid) + } + return nil +} diff --git a/blobstore/cli/shardnode/shard.go b/blobstore/cli/shardnode/shard.go index b274b614e..02aaed3dd 100644 --- a/blobstore/cli/shardnode/shard.go +++ b/blobstore/cli/shardnode/shard.go @@ -15,26 +15,15 @@ package shardnode import ( - "bytes" - "encoding/json" - "math" - "os" - "time" - - "github.com/cubefs/cubefs/blobstore/shardnode/storage" - "github.com/desertbit/grumble" "go.etcd.io/etcd/raft/v3/raftpb" - "github.com/cubefs/cubefs/blobstore/api/clustermgr" "github.com/cubefs/cubefs/blobstore/api/shardnode" "github.com/cubefs/cubefs/blobstore/cli/common" "github.com/cubefs/cubefs/blobstore/cli/common/args" "github.com/cubefs/cubefs/blobstore/cli/common/fmt" kvstore "github.com/cubefs/cubefs/blobstore/common/kvstorev2" - "github.com/cubefs/cubefs/blobstore/common/proto" "github.com/cubefs/cubefs/blobstore/common/raft" - "github.com/cubefs/cubefs/blobstore/common/rpc" "github.com/cubefs/cubefs/blobstore/common/rpc2" "github.com/cubefs/cubefs/blobstore/util/errors" ) @@ -43,67 +32,31 @@ func addCmdShard(cmd *grumble.Command) { shardCommand := &grumble.Command{ Name: "shard", Help: "shard tools", - LongHelp: "shard tools for shardNode", + LongHelp: "shard tools for shardnode", } cmd.AddCommand(shardCommand) // get shard stats shardCommand.AddCommand(&grumble.Command{ Name: "get", - Help: "get shard form shardNode", + Help: "get shard form shardnode", Args: func(a *grumble.Args) { - args.NodeHostRegister(a) args.DiskIDRegister(a) args.SuidRegister(a) }, + Flags: func(f *grumble.Flags) { + clusterFlags(f) + }, Run: cmdGetShard, }) - // list shard - shardCommand.AddCommand(&grumble.Command{ - Name: "list", - Help: "list shards form shardNode", - Args: func(a *grumble.Args) { - args.NodeHostRegister(a) - args.DiskIDRegister(a) - a.Uint64("shardID", "shardID") - a.Uint64("count", "list shard count") - }, - Run: cmdListShard, - }) - - // transfer shard leader - shardCommand.AddCommand(&grumble.Command{ - Name: "transferLeader", - Help: "transfer shard leader", - Args: func(a *grumble.Args) { - args.NodeHostRegister(a) - args.DiskIDRegister(a) - args.SuidRegister(a) - a.Uint("targetDiskID", "target leader diskID") - }, - Run: cmdTransferShardLeader, - }) - - // add shard - shardCommand.AddCommand(&grumble.Command{ - Name: "addShard", - Help: "add shard", - Args: func(a *grumble.Args) { - args.NodeHostRegister(a) - args.DiskIDRegister(a) - a.String("json", "add shard args json") - }, - Run: cmdAddShard, - }) - // raft state shardCommand.AddCommand(&grumble.Command{ Name: "hardState", - Help: "get hardState", - Args: func(a *grumble.Args) { - a.String("path", "raft storage path") - a.Uint64("shard", "shardID") + Help: "get hardState from storage, should stop shardnode first", + Flags: func(f *grumble.Flags) { + f.StringL("path", "", "raft storage path") + f.UintL("shard_id", 0, "shardID") }, Run: cmdRaftHardStat, }) @@ -111,169 +64,45 @@ func addCmdShard(cmd *grumble.Command) { // raft log shardCommand.AddCommand(&grumble.Command{ Name: "raftLog", - Help: "get raft log", - Args: func(a *grumble.Args) { - a.String("path", "raft storage path") - a.Uint64("shard", "shardID") - a.Uint64("index", "log index") + Help: "get raft log from storage, should stop shardnode first", + Flags: func(f *grumble.Flags) { + f.StringL("path", "", "raft storage path") + f.UintL("shard_id", 0, "shardID") + f.UintL("index", 0, "log index") }, Run: cmdRaftLog, }) - - // updateShardInfo - shardCommand.AddCommand(&grumble.Command{ - Name: "updateShardInfo", - Help: "get shard info from db", - Args: func(a *grumble.Args) { - a.String("path", "raft storage path") - args.SuidRegister(a) - a.String("json", "shardInfo json") - }, - Run: cmdUpdateShardInfo, - }) - - // shard data backup - shardCommand.AddCommand(&grumble.Command{ - Name: "backupShard", - Help: "backup shard from db", - Args: func(a *grumble.Args) { - a.String("path", "origin storage path") - args.SuidRegister(a) - }, - Run: cmdShardDataBackUp, - }) - - // shard data backup - shardCommand.AddCommand(&grumble.Command{ - Name: "recoverShardData", - Help: "recover shard data from db", - Args: func(a *grumble.Args) { - a.String("path", "origin storage path") - a.String("destPath", "dest storage path") - args.SuidRegister(a) - }, - Run: cmdShardDataRecover, - }) - - // shard data clear - shardCommand.AddCommand(&grumble.Command{ - Name: "clearShardData", - Help: "clear shard data from db", - Args: func(a *grumble.Args) { - a.String("path", "origin storage path") - args.SuidRegister(a) - }, - Run: cmdShardDataClear, - }) - - // recover disk shards - shardCommand.AddCommand(&grumble.Command{ - Name: "recoverDiskShard", - Help: "recover disk shard after clear shard data", - Args: func(a *grumble.Args) { - a.String("cm", "clusterMgr address") - args.NodeHostRegister(a) - args.DiskIDRegister(a) - }, - Run: cmdRecoverDiskShard, - }) - - // recover single disk shard - shardCommand.AddCommand(&grumble.Command{ - Name: "recoverSingleShard", - Help: "recover single disk shard", - Args: func(a *grumble.Args) { - a.String("cm", "clusterMgr address") - args.NodeHostRegister(a) - args.DiskIDRegister(a) - args.SuidRegister(a) - }, - Run: cmdRecoverSingleShard, - }) } func cmdGetShard(c *grumble.Context) error { ctx := common.CmdContext() - host := args.NodeHost(c.Args) + cmClient := newCMClient(c.Flags) diskID := args.DiskID(c.Args) suid := args.Suid(c.Args) - cli := shardnode.New(rpc2.Client{}) - ret, err := cli.GetShardStats(ctx, host, shardnode.GetShardArgs{ + diskInfo, err := cmClient.ShardNodeDiskInfo(ctx, diskID) + if err != nil { + return errors.Info(err, "get shardnode disk info failed") + } + + snClient := shardnode.New(rpc2.Client{}) + + ret, err := snClient.GetShardStats(ctx, diskInfo.Host, shardnode.GetShardArgs{ DiskID: diskID, Suid: suid, }) if err != nil { return err } + fmt.Println(common.Readable(ret)) return nil } -func cmdListShard(c *grumble.Context) error { - ctx := common.CmdContext() - host := args.NodeHost(c.Args) - diskID := args.DiskID(c.Args) - shardID := c.Args.Uint64("shardID") - count := c.Args.Uint64("count") - - cli := shardnode.New(rpc2.Client{}) - ret, err := cli.ListShards(ctx, host, shardnode.ListShardArgs{ - DiskID: diskID, - ShardID: proto.ShardID(shardID), - Count: count, - }) - if err != nil { - return err - } - for _, shard := range ret.Shards { - fmt.Println(common.Readable(shard)) - } - return nil -} - -func cmdTransferShardLeader(c *grumble.Context) error { - ctx := common.CmdContext() - host := args.NodeHost(c.Args) - diskID := args.DiskID(c.Args) - suid := args.Suid(c.Args) - targetDiskID := c.Args.Uint("targetDiskID") - - cli := shardnode.New(rpc2.Client{}) - err := cli.TransferShardLeader(ctx, host, shardnode.TransferShardLeaderArgs{ - DiskID: diskID, - Suid: suid, - DestDiskID: proto.DiskID(targetDiskID), - }) - if err != nil { - return err - } - return nil -} - -func cmdAddShard(c *grumble.Context) error { - ctx := common.CmdContext() - host := args.NodeHost(c.Args) - jsonStr := c.Args.String("json") - - req := shardnode.AddShardArgs{} - err := json.Unmarshal([]byte(jsonStr), &req) - if err != nil { - return err - } - - cli := shardnode.New(rpc2.Client{}) - err = cli.AddShard(ctx, host, req) - if err != nil { - return err - } - return nil -} - func cmdRaftHardStat(c *grumble.Context) error { ctx := common.CmdContext() - path := c.Args.String("path") - shard := c.Args.Uint64("shard") + path := c.Flags.String("path") + shard := c.Flags.Uint64("shard_id") store, err := kvstore.NewKVStore(ctx, path, kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{"raft-wal"}}) if err != nil { @@ -297,9 +126,9 @@ func cmdRaftHardStat(c *grumble.Context) error { func cmdRaftLog(c *grumble.Context) error { ctx := common.CmdContext() - path := c.Args.String("path") - shard := c.Args.Uint64("shard") - index := c.Args.Uint64("index") + path := c.Flags.String("path") + shard := c.Flags.Uint64("shard_id") + index := c.Flags.Uint64("index") store, err := kvstore.NewKVStore(ctx, path, kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{"raft-wal"}}) if err != nil { @@ -317,516 +146,3 @@ func cmdRaftLog(c *grumble.Context) error { } return nil } - -func cmdRecoverDiskShard(c *grumble.Context) error { - ctx := common.CmdContext() - cmHost := c.Args.String("cm") - host := args.NodeHost(c.Args) - diskID := args.DiskID(c.Args) - cfg := &clustermgr.Config{LbConfig: rpc.LbConfig{Hosts: []string{cmHost}}} - cmClient := clustermgr.New(cfg) - ret, err := cmClient.ListShardUnit(ctx, &clustermgr.ListShardUnitArgs{ - DiskID: diskID, - }) - if err != nil { - return errors.Info(err, "list shard unit failed") - } - - // fetch disk info - disks, err := cmClient.ListShardNodeDisk(ctx, &clustermgr.ListOptionArgs{Count: 100}) - if err != nil { - return err - } - diskHostMap := make(map[proto.DiskID]*clustermgr.ShardNodeDiskInfo) - for _, info := range disks.Disks { - diskHostMap[info.DiskID] = info - } - fmt.Println("load disk info success") - - snClient := shardnode.New(rpc2.Client{}) - for _, unit := range ret { - shardInfo, err := cmClient.ListShard(ctx, &clustermgr.ListShardArgs{ - Marker: unit.Suid.ShardID() - 1, - Count: 1, - }) - if err != nil { - return errors.Info(err, "list shard failed") - } - if len(shardInfo.Shards) != 1 { - return errors.New(fmt.Sprintf("get shard failed, shards: %+v", shardInfo.Shards)) - } - shard := shardInfo.Shards[0] - if shard.ShardID != unit.Suid.ShardID() { - return errors.New(fmt.Sprintf("get wrong shard: %+v", shard)) - } - units := make([]clustermgr.ShardUnit, len(shard.Units)) - for i, u := range shard.Units { - units[i] = clustermgr.ShardUnit{ - Suid: u.Suid, - DiskID: u.DiskID, - Host: u.Host, - Learner: u.Learner, - } - if u.Suid != unit.Suid { - // remove from raft group - info, ok := diskHostMap[u.DiskID] - if !ok { - return errors.New("load disk info in map failed") - } - if _err := snClient.UpdateShard(ctx, info.Host, shardnode.UpdateShardArgs{ - DiskID: u.DiskID, - Suid: u.Suid, - ShardUpdateType: proto.ShardUpdateTypeRemoveMember, - Unit: clustermgr.ShardUnit{ - Suid: unit.Suid, - DiskID: unit.DiskID, - Learner: unit.Learner, - }, - }); _err != nil { - fmt.Println("remove shard failed") - return _err - } - } - } - time.Sleep(100 * time.Millisecond) - if err = snClient.AddShard(ctx, host, shardnode.AddShardArgs{ - DiskID: unit.DiskID, - Suid: unit.Suid, - Range: unit.Range, - Units: units, - RouteVersion: unit.RouteVersion, - }); err != nil { - return errors.Info(err, "add shard failed") - } - for _, u := range shard.Units { - if u.Suid != unit.Suid { - // add to raft group - info, ok := diskHostMap[u.DiskID] - if !ok { - return errors.New("load disk info in map failed") - } - if _err := snClient.UpdateShard(ctx, info.Host, shardnode.UpdateShardArgs{ - DiskID: u.DiskID, - Suid: u.Suid, - ShardUpdateType: proto.ShardUpdateTypeAddMember, - Unit: clustermgr.ShardUnit{ - Suid: unit.Suid, - DiskID: unit.DiskID, - }, - }); _err != nil { - fmt.Println("add shard failed") - return _err - } - } - } - fmt.Println("add shard success, suid:", unit.Suid) - } - return nil -} - -func cmdRecoverSingleShard(c *grumble.Context) error { - ctx := common.CmdContext() - cmHost := c.Args.String("cm") - host := args.NodeHost(c.Args) - diskID := args.DiskID(c.Args) - suid := args.Suid(c.Args) - cfg := &clustermgr.Config{LbConfig: rpc.LbConfig{Hosts: []string{cmHost}}} - cmClient := clustermgr.New(cfg) - - // fetch disk info - disks, err := cmClient.ListShardNodeDisk(ctx, &clustermgr.ListOptionArgs{Count: 100}) - if err != nil { - return err - } - diskHostMap := make(map[proto.DiskID]*clustermgr.ShardNodeDiskInfo) - for _, info := range disks.Disks { - diskHostMap[info.DiskID] = info - } - fmt.Println("load disk info success") - - shardInfo, err := cmClient.ListShard(ctx, &clustermgr.ListShardArgs{ - Marker: suid.ShardID() - 1, - Count: 1, - }) - if err != nil || len(shardInfo.Shards) != 1 { - fmt.Println("get shard info failed") - return err - } - shard := shardInfo.Shards[0] - units := make([]clustermgr.ShardUnit, len(shard.Units)+1) - snClient := shardnode.New(rpc2.Client{}) - for i, u := range shard.Units { - // remove from raft group - fmt.Println("remove form disk: %d", u.DiskID) - info, ok := diskHostMap[u.DiskID] - if !ok { - return errors.New("load disk info in map failed") - } - if _err := snClient.UpdateShard(ctx, info.Host, shardnode.UpdateShardArgs{ - DiskID: u.DiskID, - Suid: u.Suid, - ShardUpdateType: proto.ShardUpdateTypeRemoveMember, - Unit: clustermgr.ShardUnit{ - Suid: suid, - DiskID: diskID, - }, - }); _err != nil { - fmt.Println("remove shard failed") - return _err - } - units[i] = clustermgr.ShardUnit{ - Suid: u.Suid, - DiskID: u.DiskID, - Host: u.Host, - Learner: u.Learner, - } - } - units[len(shard.Units)] = clustermgr.ShardUnit{ - Suid: suid, - DiskID: diskID, - Host: host, - Learner: true, - } - time.Sleep(100 * time.Millisecond) - if err = snClient.AddShard(ctx, host, shardnode.AddShardArgs{ - DiskID: diskID, - Suid: suid, - Range: shard.Range, - Units: units, - RouteVersion: shard.RouteVersion, - }); err != nil { - return errors.Info(err, "add shard failed") - } - for _, u := range shard.Units { - fmt.Println("add to disk: %d", u.DiskID) - // add to raft group - info, ok := diskHostMap[u.DiskID] - if !ok { - return errors.New("load disk info in map failed") - } - if _err := snClient.UpdateShard(ctx, info.Host, shardnode.UpdateShardArgs{ - DiskID: u.DiskID, - Suid: u.Suid, - ShardUpdateType: proto.ShardUpdateTypeAddMember, - Unit: clustermgr.ShardUnit{ - Suid: suid, - DiskID: diskID, - }, - }); _err != nil { - fmt.Println("add shard failed") - return _err - } - } - fmt.Println("add shard success, suid:", suid) - return nil -} - -func cmdUpdateShardInfo(c *grumble.Context) error { - ctx := common.CmdContext() - path := c.Args.String("path") - suid := args.Suid(c.Args) - jsonInfo := c.Args.String("json") - col := kvstore.CF("data") - - store, err := kvstore.NewKVStore(ctx, path, kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{col}}) - if err != nil { - return err - } - defer store.Close() - - g := storage.NewShardKeysGenerator(suid) - infoKey := g.EncodeShardInfoKey() - raw, err := store.GetRaw(ctx, col, infoKey) - if err != nil && !errors.Is(err, kvstore.ErrNotFound) { - return err - } - if err == nil { - sd := clustermgr.Shard{} - err = sd.Unmarshal(raw) - if err != nil { - return err - } - fmt.Println(common.Readable(sd)) - } - - // update - if len(jsonInfo) == 0 { - return nil - } - newInfo := &clustermgr.Shard{} - if err = json.Unmarshal([]byte(jsonInfo), newInfo); err != nil { - return err - } - fmt.Println("new info:") - fmt.Println(common.Readable(newInfo)) - - _raw, err := newInfo.Marshal() - if err != nil { - return err - } - err = store.SetRaw(ctx, col, infoKey, _raw) - if err != nil { - return err - } - store.FlushCF(ctx, col) - return nil -} - -func cmdShardDataBackUp(c *grumble.Context) error { - ctx := common.CmdContext() - path := c.Args.String("path") - suid := args.Suid(c.Args) - - colData := kvstore.CF("data") - colRaft := kvstore.CF("raft-wal") - - backPath := path + "/" + suid.ShardID().ToString() - if err := os.Mkdir(backPath, 0o755); err != nil { - return err - } - fmt.Println("create backup dir: ", backPath) - - backKVStore, err := kvstore.NewKVStore(ctx, backPath+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ - ColumnFamily: []kvstore.CF{colData}, - CreateIfMissing: true, - }) - if err != nil { - return err - } - defer backKVStore.Close() - - backRaftStore, err := kvstore.NewKVStore(ctx, backPath+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ - ColumnFamily: []kvstore.CF{colRaft}, - CreateIfMissing: true, - }) - if err != nil { - return err - } - defer backRaftStore.Close() - fmt.Println("backup store opened") - - fmt.Println("start backup data") - originKVStore, err := kvstore.NewKVStore(ctx, path+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) - if err != nil { - return err - } - defer originKVStore.Close() - fmt.Println("origin store open") - - g := storage.NewShardKeysGenerator(suid) - shardDataPrefix := g.EncodeShardDataPrefix() - dataList := originKVStore.List(ctx, colData, shardDataPrefix, nil, nil) - kvCount := 0 - for { - kg, vg, err := dataList.ReadNext() - if err != nil { - return err - } - if kg == nil || vg == nil { - fmt.Println("list to end") - break - } - if !bytes.HasPrefix(kg.Key(), shardDataPrefix) { - return errors.New("key prefix not match") - } - if err = backKVStore.SetRaw(ctx, colData, kg.Key(), vg.Value()); err != nil { - return err - } - kg.Close() - vg.Close() - kvCount++ - } - fmt.Println("shard kv data backup done, num: ", kvCount) - - originRaftStore, err := kvstore.NewKVStore(ctx, path+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) - if err != nil { - return err - } - defer originRaftStore.Close() - - raftLogPrefix := raft.EncodeIndexLogKeyPrefix(uint64(suid.ShardID())) - raftLogList := originRaftStore.List(ctx, colRaft, raftLogPrefix, nil, nil) - logCount := 0 - for { - kg, vg, err := raftLogList.ReadNext() - if err != nil { - return err - } - if kg == nil || vg == nil { - fmt.Println("list to end") - break - } - if !bytes.HasPrefix(kg.Key(), raftLogPrefix) { - return errors.New("key prefix not match") - } - if err = backRaftStore.SetRaw(ctx, colRaft, kg.Key(), vg.Value()); err != nil { - return err - } - kg.Close() - vg.Close() - logCount++ - } - fmt.Println("shard raft log backup done, num: ", logCount) - - hardStateKey := raft.EncodeHardStateKey(uint64(suid.ShardID())) - hsRaw, err := originRaftStore.GetRaw(ctx, colRaft, hardStateKey) - if err != nil { - return err - } - if err = backRaftStore.SetRaw(ctx, colRaft, hardStateKey, hsRaw); err != nil { - return err - } - fmt.Println("hardState backup done") - return nil -} - -func cmdShardDataRecover(c *grumble.Context) error { - ctx := common.CmdContext() - path := c.Args.String("path") - destPath := c.Args.String("destPath") - suid := args.Suid(c.Args) - - colData := kvstore.CF("data") - colRaft := kvstore.CF("raft-wal") - - kvStore, err := kvstore.NewKVStore(ctx, path+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) - if err != nil { - return err - } - defer kvStore.Close() - - destKVStore, err := kvstore.NewKVStore(ctx, destPath+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) - if err != nil { - return err - } - defer destKVStore.Close() - - fmt.Println("start recover shard kv data") - g := storage.NewShardKeysGenerator(suid) - shardDataPrefix := g.EncodeShardDataPrefix() - dataList := kvStore.List(ctx, colData, shardDataPrefix, nil, nil) - kvCount := 0 - for { - kg, vg, err := dataList.ReadNext() - if err != nil { - return err - } - if kg == nil || vg == nil { - fmt.Println("list to end") - break - } - if !bytes.HasPrefix(kg.Key(), shardDataPrefix) { - return errors.New("key prefix not match") - } - if err = destKVStore.SetRaw(ctx, colData, kg.Key(), vg.Value()); err != nil { - return err - } - kg.Close() - vg.Close() - kvCount++ - } - fmt.Println("shard kv data recover done, num: ", kvCount) - - raftStore, err := kvstore.NewKVStore(ctx, path+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) - if err != nil { - return err - } - defer raftStore.Close() - - destRaftStore, err := kvstore.NewKVStore(ctx, destPath+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) - if err != nil { - return err - } - defer destRaftStore.Close() - fmt.Println("start recover shard raft data") - - raftLogPrefix := raft.EncodeIndexLogKeyPrefix(uint64(suid.ShardID())) - raftLogList := raftStore.List(ctx, colRaft, raftLogPrefix, nil, nil) - logCount := 0 - for { - kg, vg, err := raftLogList.ReadNext() - if err != nil { - return err - } - if kg == nil || vg == nil { - fmt.Println("list to end") - break - } - if !bytes.HasPrefix(kg.Key(), raftLogPrefix) { - return errors.New("key prefix not match") - } - if err = destRaftStore.SetRaw(ctx, colRaft, kg.Key(), vg.Value()); err != nil { - return err - } - kg.Close() - vg.Close() - logCount++ - } - fmt.Println("shard raft log recover, num: ", logCount) - - hardStateKey := raft.EncodeHardStateKey(uint64(suid.ShardID())) - hsRaw, err := raftStore.GetRaw(ctx, colRaft, hardStateKey) - if err != nil { - return err - } - if err = destRaftStore.SetRaw(ctx, colRaft, hardStateKey, hsRaw); err != nil { - return err - } - fmt.Println("hardState recover done") - return nil -} - -func cmdShardDataClear(c *grumble.Context) error { - ctx := common.CmdContext() - path := c.Args.String("path") - suid := args.Suid(c.Args) - - colData := kvstore.CF("data") - colRaft := kvstore.CF("raft-wal") - - dataStore, err := kvstore.NewKVStore(ctx, path+"/kv", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colData}}) - if err != nil { - return err - } - defer dataStore.Close() - g := storage.NewShardKeysGenerator(suid) - shardInfoKey := g.EncodeShardInfoKey() - if err = dataStore.Delete(ctx, colData, shardInfoKey); err != nil { - return err - } - fmt.Println("shard info deleted") - - shardDataPrefix := g.EncodeShardDataPrefix() - shardDataMaxPrefix := g.EncodeShardDataMaxPrefix() - if err = dataStore.DeleteRange(ctx, colData, shardDataPrefix, shardDataMaxPrefix); err != nil { - return err - } - dataStore.FlushCF(ctx, colData) - fmt.Println("shard data deleted") - - raftStore, err := kvstore.NewKVStore(ctx, path+"/raft", kvstore.RocksdbLsmKVType, &kvstore.Option{ColumnFamily: []kvstore.CF{colRaft}}) - if err != nil { - return err - } - - hardStateKey := raft.EncodeHardStateKey(uint64(suid.ShardID())) - if err = raftStore.Delete(ctx, colRaft, hardStateKey); err != nil { - return err - } - fmt.Println("hardState deleted") - - snapShotMetaKey := raft.EncodeSnapshotMetaKey(uint64(suid.ShardID())) - if err = raftStore.Delete(ctx, colRaft, snapShotMetaKey); err != nil { - return err - } - fmt.Println("snapShot meta deleted") - - if err = raftStore.DeleteRange(ctx, colRaft, - raft.EncodeIndexLogKey(uint64(suid.ShardID()), 0), - raft.EncodeIndexLogKey(uint64(suid.ShardID()), math.MaxUint64)); err != nil { - return err - } - raftStore.FlushCF(ctx, colRaft) - fmt.Println("raft log deleted") - return nil -} diff --git a/blobstore/cli/shardnode/shardnode.go b/blobstore/cli/shardnode/shardnode.go index c23a27b88..fc67b7bb3 100644 --- a/blobstore/cli/shardnode/shardnode.go +++ b/blobstore/cli/shardnode/shardnode.go @@ -15,7 +15,13 @@ package shardnode import ( + "strings" + "github.com/desertbit/grumble" + + "github.com/cubefs/cubefs/blobstore/api/clustermgr" + "github.com/cubefs/cubefs/blobstore/cli/common/fmt" + "github.com/cubefs/cubefs/blobstore/cli/config" ) func Register(app *grumble.App) { @@ -24,7 +30,25 @@ func Register(app *grumble.App) { Help: "shardnode manager tools", } app.AddCommand(snCommand) - addCmdVol(snCommand) addCmdShard(snCommand) + addCmdRecover(snCommand) addCmdTCMalloc(snCommand) } + +func newCMClient(f grumble.FlagMap) *clustermgr.Client { + clusterID := f.String("cluster_id") + if clusterID == "" { + clusterID = fmt.Sprintf("%d", config.DefaultClusterID()) + } + var hosts []string + if str := strings.TrimSpace(f.String("hosts")); str != "" { + hosts = strings.Split(str, " ") + } + return config.NewCluster(clusterID, hosts, f.String("secret")) +} + +func clusterFlags(f *grumble.Flags) { + f.StringL("cluster_id", "", "specific clustermgr cluster id") + f.StringL("secret", "", "specific clustermgr secret") + f.StringL("hosts", "", "specific clustermgr hosts") +} diff --git a/blobstore/cli/shardnode/tcmalloc.go b/blobstore/cli/shardnode/tcmalloc.go index d2470a273..a79281b99 100644 --- a/blobstore/cli/shardnode/tcmalloc.go +++ b/blobstore/cli/shardnode/tcmalloc.go @@ -17,12 +17,13 @@ package shardnode import ( "errors" + "github.com/desertbit/grumble" + "github.com/cubefs/cubefs/blobstore/api/shardnode" "github.com/cubefs/cubefs/blobstore/cli/common" "github.com/cubefs/cubefs/blobstore/cli/common/args" "github.com/cubefs/cubefs/blobstore/cli/common/fmt" "github.com/cubefs/cubefs/blobstore/common/rpc2" - "github.com/desertbit/grumble" ) func addCmdTCMalloc(cmd *grumble.Command) { diff --git a/blobstore/cli/shardnode/volume.go b/blobstore/cli/shardnode/volume.go deleted file mode 100644 index 1814cb4c3..000000000 --- a/blobstore/cli/shardnode/volume.go +++ /dev/null @@ -1,60 +0,0 @@ -// Copyright 2024 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 shardnode - -import ( - "github.com/desertbit/grumble" - - "github.com/cubefs/cubefs/blobstore/api/shardnode" - "github.com/cubefs/cubefs/blobstore/cli/common" - "github.com/cubefs/cubefs/blobstore/cli/common/args" - "github.com/cubefs/cubefs/blobstore/cli/common/fmt" - "github.com/cubefs/cubefs/blobstore/common/rpc2" -) - -func addCmdVol(cmd *grumble.Command) { - volumeCommand := &grumble.Command{ - Name: "volume", - Help: "volume tools", - LongHelp: "volume tools for shardNode", - } - cmd.AddCommand(volumeCommand) - - volumeCommand.AddCommand(&grumble.Command{ - Name: "list", - Help: "list all volume shardNode allocated", - Args: func(a *grumble.Args) { - args.NodeHostRegister(a) - args.CodeModeRegister(a) - }, - Run: cmdListVolume, - }) -} - -func cmdListVolume(c *grumble.Context) error { - ctx := common.CmdContext() - host := args.NodeHost(c.Args) - mode := args.CodeMode(c.Args) - - cli := shardnode.New(rpc2.Client{}) - value, err := cli.ListVolume(ctx, host, shardnode.ListVolumeArgs{ - CodeMode: mode, - }) - if err != nil { - return err - } - fmt.Println(common.Readable(value)) - return nil -}