Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 36 additions & 6 deletions src/control/server/ctl_storage_rpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -802,21 +802,50 @@ func (cs *ControlService) StorageScan(ctx context.Context, req *ctlpb.StorageSca
return resp, nil
}

func (cs *ControlService) formatMetadata(instances []Engine, reformat bool) (bool, error) {
func (cs *ControlService) formatMetadata(instances []Engine, reformat, replace bool) (bool, error) {
// Format control metadata first, if needed
if needs, err := cs.storage.ControlMetadataNeedsFormat(); err != nil {
return false, errors.Wrap(err, "detecting if metadata format is needed")
} else if needs || reformat {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

From my understanding, the engine-selection logic in the new replace branch relies solely on eng.GetStorage().ScmNeedsFormat(), which is purely an SCM/tmpfs mount check and has no relation to whether that engine's own control_metadata subdirectory is actually intact.

If I am correct, that means: if engine X's control-metadata subdirectory is corrupted but its SCM still reads as mounted/formatted, ScmNeedsFormat() returns false for it, it's excluded from needFormatIdxs, and its corrupted metadata is never wiped via --replace — silently, with no error or indication anything was skipped.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I've refactor the PR please review whether this comment is still relevant thanks

// Full format needed
engineIdxs := make([]uint, len(instances))
for i, eng := range instances {
engineIdxs[i] = uint(eng.Index())
}

cs.log.Debug("formatting control metadata storage")
cs.log.Debug("formatting control metadata storage (all engines)")
if err := cs.storage.FormatControlMetadata(engineIdxs); err != nil {
return false, errors.Wrap(err, "formatting control metadata storage")
}

return true, nil
} else if replace {
// Selective format: only format engines with missing metadata directories
var needFormatIdxs []uint
for _, eng := range instances {
engIdx := uint(eng.Index())

// Check if engine's metadata directory is missing
needsFormat, err := eng.GetStorage().ControlMetadataEngineNeedsFormat()
if err != nil {
return false, errors.Wrapf(err, "checking if engine %d metadata needs format", engIdx)
}

if needsFormat {
cs.log.Debugf("engine %d metadata needs format", engIdx)
needFormatIdxs = append(needFormatIdxs, engIdx)
}
}

if len(needFormatIdxs) == 0 {
return false, errors.New("format replace option only valid if at least one " +
"engine requires metadata format but currently no engines need metadata format")
}

cs.log.Debugf("formatting control metadata storage for engines %v (--replace)", needFormatIdxs)
if err := cs.storage.FormatControlMetadata(needFormatIdxs); err != nil {
return false, errors.Wrap(err, "formatting control metadata storage")
}
return true, nil
}

Expand Down Expand Up @@ -1059,10 +1088,11 @@ func (cs *ControlService) StorageFormat(ctx context.Context, req *ctlpb.StorageF
return resp, nil
}

// DAOS-15947: control_metadata format is valid in --replace case where multiple engines
// require replacement or format on the same host. No need to handle independently for
// individual engine as if control_metadata is missing then it needs to be created.
mdFormatted, err := cs.formatMetadata(instances, req.Reformat)
// DAOS-15947, DAOS-19385: control_metadata format is required in --replace case
// to ensure old rank metadata is cleared. Only engines with missing metadata
// directories will have their control_metadata subdirectories reformatted,
// preserving healthy engines. SCM formatting is handled separately.
mdFormatted, err := cs.formatMetadata(instances, req.Reformat, req.Replace)
if err != nil {
return nil, err
}
Expand Down
36 changes: 23 additions & 13 deletions src/control/server/storage/metadata/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -198,26 +198,36 @@ func (p *Provider) isUsableFS(fs *system.FsType, path string) bool {
func (p *Provider) setupDataDir(req storage.MetadataFormatRequest) error {
perms := os.FileMode(0775)

if err := p.sys.RemoveAll(req.DataPath); err != nil {
return errors.Wrap(err, "removing old control metadata subdirectory")
}

if err := p.sys.Mkdir(req.DataPath, perms); err != nil {
return errors.Wrap(err, "creating control metadata subdirectory")
// Helper to create directory with ownership
createDirWithOwner := func(path string) error {
if err := p.sys.Mkdir(path, perms); err != nil {
return errors.Wrapf(err, "creating directory %s", path)
}
if err := p.sys.Chown(path, req.OwnerUID, req.OwnerGID); err != nil {
return errors.Wrapf(err, "setting ownership of %s to %d/%d", path, req.OwnerUID, req.OwnerGID)
}
return nil
}

if err := p.sys.Chown(req.DataPath, req.OwnerUID, req.OwnerGID); err != nil {
return errors.Wrapf(err, "setting ownership of control metadata subdirectory to %d/%d", req.OwnerUID, req.OwnerGID)
// Ensure DataPath exists
if _, err := p.sys.Stat(req.DataPath); os.IsNotExist(err) {
p.log.Debugf("creating control metadata subdirectory %q", req.DataPath)
if err := createDirWithOwner(req.DataPath); err != nil {
return errors.Wrap(err, "creating control metadata subdirectory")
}
} else if err != nil {
return errors.Wrap(err, "checking control metadata subdirectory")
}

// Selectively remove and recreate engine directories
p.log.Debugf("formatting control metadata for engines %v", req.EngineIdxs)
for _, idx := range req.EngineIdxs {
engPath := storage.ControlMetadataEngineDir(req.DataPath, idx)
if err := p.sys.Mkdir(engPath, perms); err != nil {
return errors.Wrapf(err, "creating control metadata engine %d subdirectory", idx)
if err := p.sys.RemoveAll(engPath); err != nil {
return errors.Wrapf(err, "removing control metadata for engine %d", idx)
}

if err := p.sys.Chown(engPath, req.OwnerUID, req.OwnerGID); err != nil {
return errors.Wrapf(err, "setting ownership of control metadata engine %d subdirectory to %d/%d", idx, req.OwnerUID, req.OwnerGID)
if err := createDirWithOwner(engPath); err != nil {
return errors.Wrapf(err, "control metadata engine %d subdirectory", idx)
}
}

Expand Down
170 changes: 167 additions & 3 deletions src/control/server/storage/metadata/provider_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -159,7 +159,7 @@ func TestMetadata_Provider_Format(t *testing.T) {
expMkfsOpts: []string{"-q"},
expMkfs: true,
},
"remove old data dir fails": {
"remove old data dir fails with permission denied": {
req: deviceReq,
setup: func(t *testing.T, root string) func() {
t.Helper()
Expand All @@ -175,7 +175,7 @@ func TestMetadata_Provider_Format(t *testing.T) {
}
}
},
expErr: errors.New("removing old control metadata subdirectory"),
expErr: errors.New("permission denied"),
expMkfsOpts: []string{"-q"},
expMkfs: true,
},
Expand Down Expand Up @@ -231,7 +231,7 @@ func TestMetadata_Provider_Format(t *testing.T) {
expMkfsOpts: []string{"-q", "-L", "old_label"},
expMkfs: true,
},
"path only doesn't attempt device format": {
"path only; doesn't attempt device format": {
req: pathReq,
sysCfg: &system.MockSysConfig{
MkfsErr: errors.New("mkfs was called!"),
Expand Down Expand Up @@ -646,3 +646,167 @@ func TestMetadata_Provider_Unmount(t *testing.T) {
})
}
}

// TestMetadata_Provider_setupDataDir_SelectiveEngineDelete tests the new selective
// engine directory deletion feature for DAOS-19385.
func TestMetadata_Provider_setupDataDir_SelectiveEngineDelete(t *testing.T) {
for name, tc := range map[string]struct {
setupExisting func(t *testing.T, dataPath string)
req storage.MetadataFormatRequest
sysCfg *system.MockSysConfig
expErr error
expExist []uint
expNotExist []uint
}{
"selective delete: engines 0 and 2 when datapath exists": {
setupExisting: func(t *testing.T, dataPath string) {
// Create the data directory and subdirectories for engines 0, 1, 2
if err := os.MkdirAll(dataPath, 0775); err != nil {
t.Fatal(err)
}
for _, idx := range []uint{0, 1, 2} {
engPath := storage.ControlMetadataEngineDir(dataPath, idx)
if err := os.MkdirAll(engPath, 0775); err != nil {
t.Fatal(err)
}
testFile := filepath.Join(engPath, "test.dat")
if err := os.WriteFile(testFile, []byte("test"), 0644); err != nil {
t.Fatal(err)
}
}
},
req: storage.MetadataFormatRequest{
RootPath: "/test_root",
DataPath: "/test_root/data",
OwnerUID: 100,
OwnerGID: 200,
EngineIdxs: []uint{0, 2},
},
expExist: []uint{1},
expNotExist: []uint{},
},
"selective delete: single engine when others exist": {
setupExisting: func(t *testing.T, dataPath string) {
if err := os.MkdirAll(dataPath, 0775); err != nil {
t.Fatal(err)
}
for _, idx := range []uint{0, 1, 2, 3} {
engPath := storage.ControlMetadataEngineDir(dataPath, idx)
if err := os.MkdirAll(engPath, 0775); err != nil {
t.Fatal(err)
}
}
},
req: storage.MetadataFormatRequest{
RootPath: "/test_root",
DataPath: "/test_root/data",
OwnerUID: 100,
OwnerGID: 200,
EngineIdxs: []uint{1},
},
expExist: []uint{0, 2, 3},
expNotExist: []uint{},
},
"datapath doesn't exist: create with specific engines": {
setupExisting: func(t *testing.T, dataPath string) {
// Don't create anything
},
req: storage.MetadataFormatRequest{
RootPath: "/test_root",
DataPath: "/test_root/data",
OwnerUID: 100,
OwnerGID: 200,
EngineIdxs: []uint{0, 2},
},
expExist: []uint{},
expNotExist: []uint{1, 3},
},
"stat error on datapath": {
setupExisting: func(t *testing.T, dataPath string) {
// Setup doesn't matter
},
req: storage.MetadataFormatRequest{
RootPath: "/test_root",
DataPath: "/test_root/data",
OwnerUID: 100,
OwnerGID: 200,
EngineIdxs: []uint{0, 1},
},
sysCfg: &system.MockSysConfig{
StatErrors: map[string]error{
"/test_root/data": errors.New("mock stat error"),
},
},
expErr: errors.New("mock stat error"),
},
} {
t.Run(name, func(t *testing.T) {
log, buf := logging.NewTestLogger(t.Name())
defer test.ShowBufferOnFailure(t, buf)

testDir, cleanupTestDir := test.CreateTestDir(t)
defer cleanupTestDir()

oldDataPath := tc.req.DataPath
tc.req.RootPath = filepath.Join(testDir, tc.req.RootPath)
tc.req.DataPath = filepath.Join(testDir, tc.req.DataPath)

if tc.sysCfg == nil {
tc.sysCfg = &system.MockSysConfig{}
}
// Enable real file operations for testing (unless StatErrors configured)
if tc.sysCfg.StatErrors == nil {
tc.sysCfg.RealStat = true
}
tc.sysCfg.RealMkdir = true
tc.sysCfg.RealRemoveAll = true

if tc.sysCfg.StatErrors != nil {
if statErr, exists := tc.sysCfg.StatErrors[oldDataPath]; exists {
tc.sysCfg.StatErrors[tc.req.DataPath] = statErr
}
}

if tc.setupExisting != nil {
tc.setupExisting(t, tc.req.DataPath)
}

// Ensure parent directory exists for tests where DataPath doesn't exist
if tc.req.RootPath != "" {
if err := os.MkdirAll(tc.req.RootPath, 0775); err != nil && !os.IsExist(err) {
t.Fatal(err)
}
}

p := NewProvider(log, system.NewMockSysProvider(log, tc.sysCfg), nil)

err := p.setupDataDir(tc.req)

test.CmpErr(t, tc.expErr, err)
if tc.expErr != nil {
return
}

for _, idx := range tc.expExist {
engPath := storage.ControlMetadataEngineDir(tc.req.DataPath, idx)
if _, err := os.Stat(engPath); os.IsNotExist(err) {
t.Errorf("expected engine %d directory to exist at %s", idx, engPath)
}
}

for _, idx := range tc.expNotExist {
engPath := storage.ControlMetadataEngineDir(tc.req.DataPath, idx)
if _, err := os.Stat(engPath); err == nil {
t.Errorf("expected engine %d directory NOT to exist at %s", idx, engPath)
}
}

for _, idx := range tc.req.EngineIdxs {
engPath := storage.ControlMetadataEngineDir(tc.req.DataPath, idx)
if _, err := os.Stat(engPath); os.IsNotExist(err) {
t.Errorf("expected engine %d directory created at %s", idx, engPath)
}
}
})
}
}
32 changes: 31 additions & 1 deletion src/control/server/storage/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"context"
"fmt"
"os"
"path/filepath"
"sync"

"github.com/dustin/go-humanize"
Expand All @@ -27,7 +28,8 @@ import (

const defaultMetadataPath = "/mnt/daos"

// SystemProvider provides operating system capabilities.
// SystemProvider provides a limited set of operating system capabilities. Not all of the
// capabilities in src/control/provider/system are exposed via the storage provider.
type SystemProvider interface {
system.IsMountedProvider
GetfsUsage(string) (uint64, uint64, error)
Expand Down Expand Up @@ -115,6 +117,30 @@ func (p *Provider) ControlMetadataPathConfigured() bool {
return false
}

// ControlMetadataEngineNeedsFormat checks if this engine's superblock exists.
// This is distinct from ControlMetadataNeedsFormat which checks the host-level DataPath.
func (p *Provider) ControlMetadataEngineNeedsFormat() (bool, error) {
if p == nil {
return false, errors.New("nil provider")
}

if !p.engineStorage.ControlMetadata.HasPath() {
// No metadata section defined, metadata stored on SCM
return false, nil
}

superblockPath := filepath.Join(p.ControlMetadataEnginePath(), "superblock")

if _, err := p.Sys.ReadFile(superblockPath); os.IsNotExist(err) {
p.log.Debugf("engine %d superblock missing: %s", p.engineIndex, superblockPath)
return true, nil
} else if err != nil {
return false, errors.Wrapf(err, "checking engine %d superblock", p.engineIndex)
}

return false, nil
}

// ControlMetadataPath returns the path where control plane metadata is stored.
func (p *Provider) ControlMetadataPath() string {
if p == nil {
Expand Down Expand Up @@ -316,6 +342,10 @@ func (p *Provider) MountScm() error {

// UnmountTmpfs unmounts SCM based on provider config.
func (p *Provider) UnmountTmpfs() error {
if p == nil {
return errors.New("nil provider")
}

cfg, err := p.GetScmConfig()
if err != nil {
return err
Expand Down
Loading
Loading