feat(shardnode): fulfill shard node broken disk process and disk reopen process

with #22357426

Signed-off-by: Cloudstriff <chenjiongwendao@qq.com>
This commit is contained in:
Cloudstriff 2024-07-31 10:12:25 +08:00 committed by slasher
parent f3a83ecb47
commit 6b84ed7f32
9 changed files with 420 additions and 75 deletions

View File

@ -219,7 +219,7 @@ type (
MaxTableFileSize int `json:"max_table_file_size,omitempty"`
AllowCompaction bool `json:"allow_compaction,omitempty"`
}
HandleError func(err error)
HandleError func(ctx context.Context, err error)
readOpts struct {
opt ReadOption

View File

@ -424,7 +424,7 @@ func (lr *listReader) ReadNext() (key KeyGetter, val ValueGetter, err error) {
lr.iterator.Next()
}
if err = lr.iterator.Err(); err != nil {
lr.handleError(err)
lr.handleError(context.TODO(), err)
return nil, nil, err
}
if !lr.iterator.Valid() {
@ -445,7 +445,7 @@ func (lr *listReader) ReadNext() (key KeyGetter, val ValueGetter, err error) {
func (lr *listReader) ReadNextCopy() (key []byte, value []byte, err error) {
kg, vg, err := lr.ReadNext()
if err != nil {
lr.handleError(err)
lr.handleError(context.TODO(), err)
return nil, nil, err
}
if kg != nil && vg != nil {
@ -466,7 +466,7 @@ func (lr *listReader) ReadPrev() (key KeyGetter, val ValueGetter, err error) {
lr.iterator.Prev()
}
if err = lr.iterator.Err(); err != nil {
lr.handleError(err)
lr.handleError(context.TODO(), err)
return nil, nil, err
}
if !lr.iterator.Valid() {
@ -487,7 +487,7 @@ func (lr *listReader) ReadPrev() (key KeyGetter, val ValueGetter, err error) {
func (lr *listReader) ReadPrevCopy() (key []byte, value []byte, err error) {
kg, vg, err := lr.ReadPrev()
if err != nil {
lr.handleError(err)
lr.handleError(context.TODO(), err)
return nil, nil, err
}
if kg != nil && vg != nil {
@ -539,7 +539,7 @@ func (lr *listReader) SeekForPrev(key []byte) (err error) {
}
for {
if err = lr.iterator.Err(); err != nil {
lr.handleError(err)
lr.handleError(context.TODO(), err)
return
}
if !lr.iterator.Valid() {
@ -878,7 +878,7 @@ func (s *rocksdb) Read(ctx context.Context, cols []CF, keys [][]byte, opts ...Re
func (s *rocksdb) FlushCF(ctx context.Context, col CF) error {
cf := s.getColumnFamily(col)
if err := s.db.FlushCF(s.fo, cf); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return err
}
return nil
@ -1036,7 +1036,7 @@ func (s *rocksdb) get(ctx context.Context, col CF, key []byte, readOpt ReadOptio
ro = readOpt.(*readOption).opt
}
if v, err = s.db.GetCF(ro, cf, key); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return nil, err
}
if !v.Exists() {
@ -1054,7 +1054,7 @@ func (s *rocksdb) getRaw(ctx context.Context, col CF, key []byte, readOpt ReadOp
ro = readOpt.(*readOption).opt
}
if v, err = s.db.GetCF(ro, cf, key); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return nil, err
}
if !v.Exists() {
@ -1077,7 +1077,7 @@ func (s *rocksdb) read(ctx context.Context, cols []CF, keys [][]byte, readOpt Re
}
_values, err := s.db.MultiGetCFMultiCF(ro, cfhs, keys)
if err != nil {
s.handleError(err)
s.handleError(ctx, err)
return nil, err
}
values = make([]ValueGetter, len(_values))
@ -1099,7 +1099,7 @@ func (s *rocksdb) multiGet(ctx context.Context, col CF, keys [][]byte, readOpt R
cfh := s.getColumnFamily(col)
_values, err := s.db.MultiGetCF(ro, cfh, keys...)
if err != nil {
s.handleError(err)
s.handleError(ctx, err)
return nil, err
}
values = make([]ValueGetter, len(_values))
@ -1120,7 +1120,7 @@ func (s *rocksdb) set(ctx context.Context, col CF, key []byte, value []byte, wri
wo = writeOpt.(*writeOption).opt
}
if err := s.db.PutCF(wo, cf, key, value); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return err
}
return nil
@ -1133,7 +1133,7 @@ func (s *rocksdb) delete(ctx context.Context, col CF, key []byte, writeOpt Write
wo = writeOpt.(*writeOption).opt
}
if err := s.db.DeleteCF(wo, cf, key); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return err
}
return nil
@ -1148,7 +1148,7 @@ func (s *rocksdb) deleteRange(ctx context.Context, col CF, start, end []byte, wr
b := rdb.NewWriteBatch()
b.DeleteRangeCF(cf, start, end)
if err := s.db.Write(wo, b); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return err
}
return nil
@ -1161,7 +1161,7 @@ func (s *rocksdb) write(ctx context.Context, batch WriteBatch, writeOpt WriteOpt
}
_batch := batch.(*writeBatch)
if err := s.db.Write(wo, _batch.batch); err != nil {
s.handleError(err)
s.handleError(ctx, err)
return err
}
return nil

View File

@ -17,7 +17,11 @@ package shardnode
import (
"context"
"fmt"
"sync"
"time"
"github.com/cubefs/cubefs/blobstore/shardnode/storage/store"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
"github.com/cubefs/cubefs/blobstore/common/proto"
@ -61,7 +65,7 @@ func (s *service) initDisks(ctx context.Context) error {
// load disk from local
disks := make([]*storage.Disk, 0, len(s.cfg.DisksConfig.Disks))
for _, diskPath := range s.cfg.DisksConfig.Disks {
disk := storage.OpenDisk(ctx, storage.DiskConfig{
disk, err := storage.OpenDisk(ctx, storage.DiskConfig{
ClusterID: s.cfg.ClusterID,
NodeID: s.transport.NodeID(),
DiskPath: diskPath,
@ -70,7 +74,28 @@ func (s *service) initDisks(ctx context.Context) error {
Transport: s.transport,
RaftConfig: s.cfg.RaftConfig,
ShardBaseConfig: s.cfg.ShardBaseConfig,
HandleEIO: s.handleEIO,
})
// open disk failed, check disk status,
if err != nil {
registeredDisk, ok := registerDiskPathsMap[diskPath]
// fatal abort when disk is not registered,
if !ok {
span.Fatalf("open disk[%s] failed: %s", diskPath, err)
}
// skip open disk when disk path status is not normal
if registeredDisk.Status != proto.DiskStatusNormal {
continue
}
// handleEIO when disk path status is normal and err is EIO
if store.IsEIO(err) {
s.handleEIO(ctx, registeredDisk.DiskID, err)
continue
}
// other situation, do fatal log
span.Fatalf("open disk[%s] failed: %s", diskPath, err)
}
disks = append(disks, disk)
}
// compare local disk and remote disk info, alloc new disk id and register new disk
@ -104,7 +129,7 @@ func (s *service) initDisks(ctx context.Context) error {
disk.SetDiskID(diskID)
}
// save disk meta
if err := disk.SaveDiskInfo(); err != nil {
if err := disk.SaveDiskInfo(ctx); err != nil {
return errors.Newf("save disk info[%+v] failed: %s", disk, err)
}
diskInfo := disk.GetDiskInfo()
@ -134,6 +159,185 @@ func (s *service) initDisks(ctx context.Context) error {
return nil
}
func (s *service) handleEIO(ctx context.Context, diskID proto.DiskID, err error) {
span := trace.SpanFromContextSafe(ctx)
span.Warnf("handle eio from storage layer, disk[%d], err: %s", diskID, err)
disk, err := s.getDisk(diskID)
if err != nil {
span.Warnf("get disk failed: %s, maybe has been removed", err)
return
}
s.groupRun.Do(fmt.Sprintf("disk-%d", diskID), func() (interface{}, error) {
// Note: there is another goroutine set disk broken,
// just return and do not do any progress below
if !disk.SetBroken() {
return nil, nil
}
for {
err := s.transport.SetDiskBroken(ctx, diskID)
if err == nil {
break
}
span.Errorf("set Disk[%d] broken to cm failed", diskID)
time.Sleep(5 * time.Second)
}
// wait for disk repairing
go s.waitRepairCloseDisk(ctx, disk)
return nil, nil
})
}
func (s *service) waitRepairCloseDisk(ctx context.Context, disk *storage.Disk) {
span := trace.SpanFromContextSafe(ctx)
diskInfo := disk.GetDiskInfo()
diskID := diskInfo.DiskID
func() {
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-s.closer.Done():
span.Warnf("service is closed. return")
return
case <-ticker.C:
}
info, err := s.transport.GetDisk(ctx, diskID)
if err != nil {
span.Errorf("get disk info from clustermgr failed. disk[%d], err:%+v", diskID, err)
continue
}
if info.Status >= proto.DiskStatusRepairing {
span.Infof("disk:%d path:%s status:%v", diskID, info.Path, info.Status)
break
}
}
// after the repair is triggered, the handle can be safely removed
span.Infof("Delete %d from the map table of the service", diskID)
s.lock.Lock()
delete(s.disks, diskID)
s.lock.Unlock()
disk.ResetShards()
span.Infof("disk %d will be gc close", diskID)
}()
s.waitReOpenDisk(ctx, diskInfo)
}
func (s *service) waitReOpenDisk(ctx context.Context, diskInfo clustermgr.ShardNodeDiskInfo) {
span := trace.SpanFromContextSafe(ctx)
span.Infof("start to wait for disk[%+v] reopen", diskInfo)
diskID := diskInfo.DiskID
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
// wait for disk reopen
for {
select {
case <-s.closer.Done():
span.Warnf("service is closed. return")
return
case <-ticker.C:
// check old path disk has been repaired or not
info, err := s.transport.GetDisk(ctx, diskID)
if err != nil {
span.Warnf("get disk from cm failed: %s", err)
continue
}
if info.Status != proto.DiskStatusRepaired {
span.Warnf("disk[%d] is not repaired", diskID)
continue
}
// check disk path empty or not
empty, err := storage.IsEmptyDisk(diskInfo.Path)
if err != nil || !empty {
span.Errorf("disk path(%s) is not empty. err: %s", diskInfo.Path, err)
continue
}
ok := func() bool {
success := false
// register new disk
disk, err := storage.OpenDisk(ctx, storage.DiskConfig{
ClusterID: s.cfg.ClusterID,
NodeID: s.transport.NodeID(),
DiskPath: diskInfo.Path,
CheckMountPoint: s.cfg.DisksConfig.CheckMountPoint,
StoreConfig: s.cfg.StoreConfig,
Transport: s.transport,
RaftConfig: s.cfg.RaftConfig,
ShardBaseConfig: s.cfg.ShardBaseConfig,
HandleEIO: s.handleEIO,
})
if err != nil {
span.Errorf("open disk[%s] failed: %s", diskInfo.Path, err)
return false
}
defer func() {
if !success {
disk.Close()
}
}()
// alloc new disk id
if disk.DiskID() == 0 {
diskID, err := s.transport.AllocDiskID(ctx)
if err != nil {
span.Errorf("alloc disk id failed: %s", err)
return false
}
// save disk id
disk.SetDiskID(diskID)
}
// save disk meta
if err := disk.SaveDiskInfo(ctx); err != nil {
span.Errorf("save disk info[%+v] failed: %s", disk, err)
return false
}
diskInfo := disk.GetDiskInfo()
// register disk
if err := s.transport.RegisterDisk(ctx, &diskInfo); err != nil {
span.Errorf("register new disk[%+v] failed: %s", disk, err)
return false
}
if err := disk.Load(ctx); err != nil {
span.Errorf("load disk failed: %s", err)
return false
}
success = true
s.addDisk(disk)
span.Infof("reopen disk[%d] success", diskID)
return true
}()
if ok {
return
}
}
}
}
func initConfig(cfg *Config) {
if cfg.NodeConfig.RaftHost == "" || cfg.NodeConfig.Host == "" {
log.Panicf("invalid node[%+v] config port", cfg.NodeConfig)

View File

@ -18,6 +18,9 @@ import (
"context"
"io"
"os"
"path/filepath"
"regexp"
"runtime"
"sync"
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
@ -48,10 +51,11 @@ type (
RaftConfig raft.Config
Transport base.Transport
ShardBaseConfig ShardBaseConfig
HandleEIO func(ctx context.Context, diskID proto.DiskID, err error)
}
)
func OpenDisk(ctx context.Context, cfg DiskConfig) *Disk {
func OpenDisk(ctx context.Context, cfg DiskConfig) (*Disk, error) {
span := trace.SpanFromContext(ctx)
if cfg.CheckMountPoint {
@ -60,27 +64,34 @@ func OpenDisk(ctx context.Context, cfg DiskConfig) *Disk {
}
}
success := false
disk := &Disk{}
cfg.StoreConfig.Path = cfg.DiskPath
cfg.StoreConfig.KVOption.ColumnFamily = []kvstore.CF{dataCF}
cfg.StoreConfig.RaftOption.ColumnFamily = []kvstore.CF{raftWalCF}
cfg.StoreConfig.HandleEIO = func(err error) {
span.Warnf("handle eio from store layer: %s", err)
if err := cfg.Transport.SetDiskBroken(ctx, disk.diskInfo.DiskID); err != nil {
span.Errorf("set Disk[%d] broken failed", disk.diskInfo.DiskID)
}
cfg.StoreConfig.HandleEIO = func(ctx context.Context, err error) {
cfg.HandleEIO(ctx, disk.DiskID(), err)
}
store, err := store.NewStore(ctx, &cfg.StoreConfig)
if err != nil {
span.Panicf("new store instance failed: %s", errors.Detail(err))
span.Errorf("new store instance failed: %s", errors.Detail(err))
return nil, err
}
// close store engine when open disk failed
defer func() {
if !success {
store.Close()
}
}()
// load Disk meta info
stats, err := store.Stats()
if err != nil {
span.Panicf("stats store info failed: %s", err)
span.Errorf("stats store info failed: %s", err)
return nil, err
}
diskInfo := clustermgr.ShardNodeDiskInfo{
DiskInfo: clustermgr.DiskInfo{
@ -95,17 +106,20 @@ func OpenDisk(ctx context.Context, cfg DiskConfig) *Disk {
},
}
rawFS := store.NewRawFS(sysRawFSPath)
f, err := rawFS.OpenRawFile(diskMetaFile)
f, err := rawFS.OpenRawFile(ctx, diskMetaFile)
if err != nil && !os.IsNotExist(err) {
span.Panicf("open Disk meta file failed : %s", err)
span.Errorf("open Disk meta file failed : %s", err)
return nil, err
}
if err == nil {
b, err := io.ReadAll(f)
if err != nil {
span.Panicf("read Disk meta file failed: %s", err)
span.Errorf("read Disk meta file failed: %s", err)
return nil, err
}
if err := diskInfo.Unmarshal(b); err != nil {
span.Panicf("unmarshal Disk meta failed: %s, raw: %v", err, b)
span.Errorf("unmarshal Disk meta failed: %s, raw: %v", err, b)
return nil, err
}
}
@ -114,7 +128,46 @@ func OpenDisk(ctx context.Context, cfg DiskConfig) *Disk {
disk.store = store
disk.shardsMu.shards = make(map[proto.Suid]*shard)
return disk
// disk will be gc by finalizer
runtime.SetFinalizer(disk, func(disk *Disk) {
disk.Close()
})
success = true
return disk, nil
}
func IsEmptyDisk(path string) (bool, error) {
absPath, err := filepath.Abs(path)
if err != nil {
return false, err
}
safePattern := `.*`
if match, _ := regexp.MatchString(safePattern, absPath); !match {
return false, errors.New("file path is invalid")
}
fis, err := os.ReadDir(absPath)
if err != nil {
return false, err
}
if len(fis) == 0 {
return true, nil
}
sysInitDir := map[string]bool{
"lost+found": true,
}
for _, fi := range fis {
if !sysInitDir[fi.Name()] {
return false, nil
}
}
return true, nil
}
type Disk struct {
@ -196,6 +249,10 @@ func (d *Disk) AddShard(ctx context.Context, suid proto.Suid,
) error {
span := trace.SpanFromContext(ctx)
if err := d.prepRWCheck(); err != nil {
return err
}
d.shardsMu.Lock()
defer d.shardsMu.Unlock()
@ -240,7 +297,11 @@ func (d *Disk) AddShard(ctx context.Context, suid proto.Suid,
}
func (d *Disk) UpdateShard(ctx context.Context, suid proto.Suid, op proto.ShardUpdateType, node clustermgr.ShardUnit) error {
shard, err := d.GetShard(suid)
if err := d.prepRWCheck(); err != nil {
return err
}
shard, err := d.getShard(suid)
if err != nil {
return err
}
@ -253,18 +314,19 @@ func (d *Disk) UpdateShard(ctx context.Context, suid proto.Suid, op proto.ShardU
return shard.UpdateShard(ctx, op, node, nodeHost.String())
}
func (d *Disk) GetShard(suid proto.Suid) (*shard, error) {
d.shardsMu.RLock()
s := d.shardsMu.shards[suid]
d.shardsMu.RUnlock()
if s == nil {
return nil, apierr.ErrShardDoesNotExist
func (d *Disk) GetShard(suid proto.Suid) (ShardHandler, error) {
if err := d.prepRWCheck(); err != nil {
return nil, err
}
return s, nil
return d.getShard(suid)
}
func (d *Disk) DeleteShard(ctx context.Context, suid proto.Suid) error {
if err := d.prepRWCheck(); err != nil {
return err
}
d.shardsMu.RLock()
shard := d.shardsMu.shards[suid]
d.shardsMu.RUnlock()
@ -300,6 +362,10 @@ func (d *Disk) DeleteShard(ctx context.Context, suid proto.Suid) error {
}
func (d *Disk) RangeShard(f func(s ShardHandler) bool) {
if err := d.prepRWCheck(); err != nil {
return
}
d.shardsMu.RLock()
for _, shard := range d.shardsMu.shards {
if !f(shard) {
@ -323,9 +389,13 @@ func (d *Disk) GetShardCnt() int {
return ret
}
func (d *Disk) SaveDiskInfo() error {
func (d *Disk) SaveDiskInfo(ctx context.Context) error {
if err := d.prepRWCheck(); err != nil {
return err
}
rawFS := d.store.NewRawFS(sysRawFSPath)
f, err := rawFS.CreateRawFile(diskMetaFile)
f, err := rawFS.CreateRawFile(ctx, diskMetaFile)
if err != nil {
return err
}
@ -360,3 +430,53 @@ func (d *Disk) DiskID() proto.DiskID {
func (d *Disk) SetDiskID(diskID proto.DiskID) {
d.diskInfo.DiskID = diskID
}
func (d *Disk) SetBroken() bool {
d.lock.Lock()
defer d.lock.Unlock()
if d.diskInfo.Status == proto.DiskStatusNormal {
d.diskInfo.Status = proto.DiskStatusBroken
return true
}
return false
}
func (d *Disk) ResetShards() {
d.lock.Lock()
d.shardsMu.shards = make(map[proto.Suid]*shard)
d.lock.Unlock()
}
func (d *Disk) Close() {
d.raftManager.Close()
d.store.Close()
}
func (d *Disk) getShard(suid proto.Suid) (*shard, error) {
d.shardsMu.RLock()
s := d.shardsMu.shards[suid]
d.shardsMu.RUnlock()
if s == nil {
return nil, apierr.ErrShardDoesNotExist
}
return s, nil
}
func (d *Disk) prepRWCheck() error {
if !d.isWritable() {
return apierr.ErrDiskBroken
}
return nil
}
func (d *Disk) isWritable() bool {
d.lock.RLock()
status := d.diskInfo.Status
d.lock.RUnlock()
return status == proto.DiskStatusNormal
}

View File

@ -15,7 +15,6 @@
package storage
import (
"errors"
"fmt"
"math/rand"
"os"
@ -63,8 +62,11 @@ func newMockDisk(tb testing.TB) (*mockDisk, func()) {
tp.EXPECT().GetNode(A, A).Return(&clustermgr.ShardNodeInfo{}, nil).AnyTimes()
tp.EXPECT().GetDisk(A, A).Return(&clustermgr.ShardNodeDiskInfo{}, nil).AnyTimes()
disk := OpenDisk(ctx, cfg)
disk, err := OpenDisk(ctx, cfg)
require.NoError(tb, err)
disk.diskInfo.DiskID = 1
disk.diskInfo.Status = proto.DiskStatusNormal
require.NoError(tb, disk.Load(ctx))
return &mockDisk{d: disk, tp: tp}, func() {
time.Sleep(time.Second)
@ -87,29 +89,33 @@ func TestServerDisk_Open(t *testing.T) {
cfg.CheckMountPoint = true
_panic()
cfg.CheckMountPoint = false
_panic()
_, err := OpenDisk(ctx, cfg)
require.Errorf(t, err, "")
cfg.StoreConfig.KVOption.CreateIfMissing = true
cfg.StoreConfig.RaftOption.CreateIfMissing = true
cfg.Transport = tp
disk := OpenDisk(ctx, cfg)
disk, err := OpenDisk(ctx, cfg)
require.NoError(t, err)
disk.diskInfo.Status = proto.DiskStatusNormal
tp.EXPECT().SetDiskBroken(A, A).Return(errors.New("set Disk broken"))
disk.cfg.StoreConfig.HandleEIO(nil)
require.NoError(t, disk.SaveDiskInfo())
require.NoError(t, disk.SaveDiskInfo(ctx))
t.Logf("Disk info: %+v", disk.GetDiskInfo())
disk.store.KVStore().Close()
disk.store.RaftStore().Close()
disk = OpenDisk(ctx, cfg)
disk, err = OpenDisk(ctx, cfg)
require.NoError(t, err)
disk.store.KVStore().Close()
disk.store.RaftStore().Close()
f, err := disk.store.NewRawFS(sysRawFSPath).CreateRawFile(diskMetaFile)
f, err := disk.store.NewRawFS(sysRawFSPath).CreateRawFile(ctx, diskMetaFile)
require.NoError(t, err)
f.Write([]byte{'0'})
_panic()
_, err = OpenDisk(ctx, cfg)
require.Errorf(t, err, "")
}
func TestServerDisk_Shard(t *testing.T) {

View File

@ -76,6 +76,7 @@ type (
store *store.Store
raftManager raft.Manager
addrResolver raft.AddressResolver
disk *Disk
}
shardStatus uint8
@ -93,7 +94,8 @@ func newShard(ctx context.Context, cfg shardConfig) (s *shard, err error) {
shardKeys: &shardKeysGenerator{
suid: cfg.suid,
},
cfg: cfg.ShardBaseConfig,
disk: cfg.disk,
cfg: cfg.ShardBaseConfig,
}
s.shardInfoMu.shardInfo = cfg.shardInfo
@ -152,6 +154,8 @@ type shard struct {
lastTruncatedIndex uint64
}
// add disk ref for finalizer gc
disk *Disk
shardKeys *shardKeysGenerator
store *store.Store
raftGroup raft.Group

View File

@ -15,15 +15,16 @@
package store
import (
"context"
"os"
"path/filepath"
)
type (
RawFS interface {
CreateRawFile(name string) (RawFile, error)
OpenRawFile(name string) (RawFile, error)
ReadDir(dir string) ([]string, error)
CreateRawFile(ctx context.Context, name string) (RawFile, error)
OpenRawFile(ctx context.Context, name string) (RawFile, error)
ReadDir(ctx context.Context, dir string) ([]string, error)
}
RawFile interface {
Read(p []byte) (n int, err error)
@ -34,28 +35,28 @@ type (
type posixRawFS struct {
path string
handleError func(err error)
handleError func(ctx context.Context, err error)
}
func (r *posixRawFS) CreateRawFile(name string) (RawFile, error) {
func (r *posixRawFS) CreateRawFile(ctx context.Context, name string) (RawFile, error) {
filePath := r.path + "/" + name
f, err := os.OpenFile(r.path+"/"+name, os.O_CREATE|os.O_RDWR, 0o755)
if err != nil {
if !os.IsNotExist(err) {
r.handleError(err)
r.handleError(ctx, err)
return nil, err
}
dir := filepath.Dir(filePath)
if err = os.MkdirAll(dir, 0o755); err != nil {
r.handleError(err)
r.handleError(ctx, err)
return nil, err
}
f, err = os.OpenFile(r.path+"/"+name, os.O_CREATE|os.O_RDWR, 0o755)
if err != nil {
r.handleError(err)
r.handleError(ctx, err)
return nil, err
}
}
@ -63,19 +64,19 @@ func (r *posixRawFS) CreateRawFile(name string) (RawFile, error) {
return &posixRawFile{f: f, handleError: r.handleError}, nil
}
func (r *posixRawFS) OpenRawFile(name string) (RawFile, error) {
func (r *posixRawFS) OpenRawFile(ctx context.Context, name string) (RawFile, error) {
f, err := os.OpenFile(r.path+"/"+name, os.O_RDONLY, 0o755)
if err != nil {
r.handleError(err)
r.handleError(ctx, err)
return nil, err
}
return &posixRawFile{f: f, handleError: r.handleError}, nil
}
func (r *posixRawFS) ReadDir(dir string) ([]string, error) {
func (r *posixRawFS) ReadDir(ctx context.Context, dir string) ([]string, error) {
entries, err := os.ReadDir(r.path + "/" + dir)
if err != nil {
r.handleError(err)
r.handleError(ctx, err)
return nil, err
}
@ -89,13 +90,13 @@ func (r *posixRawFS) ReadDir(dir string) ([]string, error) {
type posixRawFile struct {
f *os.File
handleError func(err error)
handleError func(ctx context.Context, err error)
}
func (pf *posixRawFile) Read(p []byte) (n int, err error) {
n, err = pf.f.Read(p)
if err != nil {
pf.handleError(err)
pf.handleError(context.TODO(), err)
}
return n, err
}
@ -103,7 +104,7 @@ func (pf *posixRawFile) Read(p []byte) (n int, err error) {
func (pf *posixRawFile) Write(p []byte) (n int, err error) {
n, err = pf.f.Write(p)
if err != nil {
pf.handleError(err)
pf.handleError(context.TODO(), err)
}
return n, err
}

View File

@ -23,28 +23,28 @@ import (
)
type Config struct {
KVOption kvstore.Option `json:"kv_option"`
RaftOption kvstore.Option `json:"raft_option"`
Path string `json:"-"`
HandleEIO func(err error) `json:"-"`
KVOption kvstore.Option `json:"kv_option"`
RaftOption kvstore.Option `json:"raft_option"`
Path string `json:"-"`
HandleEIO func(ctx context.Context, err error) `json:"-"`
}
type Store struct {
kvStore kvstore.Store
raftStore kvstore.Store
defaultRawFS RawFS
handleError func(err error)
handleError func(ctx context.Context, err error)
cfg *Config
}
func NewStore(ctx context.Context, cfg *Config) (*Store, error) {
handleError := func(err error) {
handleError := func(ctx context.Context, err error) {
if err == nil {
return
}
if IsEIO(err) && cfg.HandleEIO != nil {
cfg.HandleEIO(err)
cfg.HandleEIO(ctx, err)
}
}
@ -61,6 +61,7 @@ func NewStore(ctx context.Context, cfg *Config) (*Store, error) {
cfg.RaftOption.HandleError = handleError
raftStore, err := kvstore.NewKVStore(ctx, raftStorePath, kvstore.RocksdbLsmKVType, &cfg.RaftOption)
if err != nil {
kvStore.Close()
return nil, errors.Info(err, "open raft store failed")
}
@ -92,3 +93,8 @@ func (s *Store) DefaultRawFS() RawFS {
func (s *Store) Stats() (Stats, error) {
return StatFS(s.cfg.Path)
}
func (s *Store) Close() {
s.kvStore.Close()
s.raftStore.Close()
}

View File

@ -16,6 +16,8 @@
package shardnode
import (
"context"
"golang.org/x/sync/singleflight"
"sync"
"github.com/cubefs/cubefs/blobstore/util/taskpool"
@ -44,6 +46,7 @@ type Config struct {
NodeConfig clustermgr.ShardNodeInfo `json:"node_config"`
RaftConfig raft.Config `json:"raft_config"`
ShardBaseConfig storage.ShardBaseConfig `json:"shard_base_config"`
HandleIOError func(ctx context.Context)
}
func newService() *service {
@ -55,6 +58,7 @@ type service struct {
disks map[proto.DiskID]*storage.Disk
transport base.Transport
taskPool taskpool.TaskPool
groupRun singleflight.Group
cfg Config
lock sync.RWMutex