mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 14:45:51 +00:00
Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
feef8ed40e | ||
|
|
8e04f961f0 | ||
|
|
5400fdca3a | ||
|
|
d3d4e354bf | ||
|
|
df6f325574 | ||
|
|
0a6a76b71e | ||
|
|
93eff8fa45 | ||
|
|
cc53bb17ca | ||
|
|
e1f2e7e357 | ||
|
|
6e715f14d5 | ||
|
|
b09f4cc3e7 | ||
|
|
a8fab5503c |
@@ -20,6 +20,8 @@ jobs:
|
||||
include:
|
||||
- worker: erasure_coding
|
||||
path: test/plugin_workers/erasure_coding
|
||||
- worker: ec_repair
|
||||
path: test/plugin_workers/ec_repair
|
||||
- worker: vacuum
|
||||
path: test/plugin_workers/vacuum
|
||||
- worker: volume_balance
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
package ec_repair_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
pluginworkers "github.com/seaweedfs/seaweedfs/test/plugin_workers"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
|
||||
pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker"
|
||||
"github.com/stretchr/testify/require"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
)
|
||||
|
||||
func TestEcRepairDetectionFindsCandidates(t *testing.T) {
|
||||
const (
|
||||
volumeID = uint32(100)
|
||||
collection = "ec-test"
|
||||
diskType = "hdd"
|
||||
)
|
||||
|
||||
node1 := nodeSpec{
|
||||
id: "node1",
|
||||
address: "127.0.0.1:11001",
|
||||
diskType: diskType,
|
||||
diskID: 0,
|
||||
ecShards: []*master_pb.VolumeEcShardInformationMessage{
|
||||
buildEcShardInfo(volumeID, collection, diskType, 0, map[uint32]int64{
|
||||
0: 100,
|
||||
1: 100,
|
||||
2: 100,
|
||||
3: 100,
|
||||
4: 100,
|
||||
}),
|
||||
},
|
||||
}
|
||||
|
||||
node2 := nodeSpec{
|
||||
id: "node2",
|
||||
address: "127.0.0.1:11002",
|
||||
diskType: diskType,
|
||||
diskID: 0,
|
||||
ecShards: []*master_pb.VolumeEcShardInformationMessage{
|
||||
buildEcShardInfo(volumeID, collection, diskType, 0, map[uint32]int64{
|
||||
0: 50,
|
||||
5: 100,
|
||||
6: 100,
|
||||
7: 100,
|
||||
8: 100,
|
||||
9: 100,
|
||||
}),
|
||||
},
|
||||
}
|
||||
|
||||
topo := buildTopology([]nodeSpec{node1, node2})
|
||||
response := &master_pb.VolumeListResponse{TopologyInfo: topo}
|
||||
master := pluginworkers.NewMasterServer(t, response)
|
||||
|
||||
dialOption := grpc.WithTransportCredentials(insecure.NewCredentials())
|
||||
handler := pluginworker.NewEcRepairHandler(dialOption)
|
||||
harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{
|
||||
WorkerOptions: pluginworker.WorkerOptions{GrpcDialOption: dialOption},
|
||||
Handlers: []pluginworker.JobHandler{handler},
|
||||
})
|
||||
harness.WaitForJobType("ec_repair")
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
proposals, err := harness.Plugin().RunDetection(ctx, "ec_repair", &plugin_pb.ClusterContext{
|
||||
MasterGrpcAddresses: []string{master.Address()},
|
||||
}, 10)
|
||||
require.NoError(t, err)
|
||||
require.NotEmpty(t, proposals)
|
||||
require.Equal(t, "ec_repair", proposals[0].JobType)
|
||||
}
|
||||
@@ -0,0 +1,131 @@
|
||||
package ec_repair_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
pluginworkers "github.com/seaweedfs/seaweedfs/test/plugin_workers"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
|
||||
pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker"
|
||||
"github.com/stretchr/testify/require"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
)
|
||||
|
||||
func TestEcRepairExecutionRepairsShards(t *testing.T) {
|
||||
const (
|
||||
volumeID = uint32(200)
|
||||
collection = "ec-repair"
|
||||
diskType = "hdd"
|
||||
)
|
||||
|
||||
volumeServers := []*pluginworkers.VolumeServer{
|
||||
pluginworkers.NewVolumeServer(t, ""),
|
||||
pluginworkers.NewVolumeServer(t, ""),
|
||||
pluginworkers.NewVolumeServer(t, ""),
|
||||
pluginworkers.NewVolumeServer(t, ""),
|
||||
}
|
||||
|
||||
nodes := []nodeSpec{
|
||||
{
|
||||
id: "node1",
|
||||
address: volumeServers[0].Address(),
|
||||
diskType: diskType,
|
||||
diskID: 0,
|
||||
ecShards: []*master_pb.VolumeEcShardInformationMessage{
|
||||
buildEcShardInfo(volumeID, collection, diskType, 0, map[uint32]int64{
|
||||
0: 100,
|
||||
1: 100,
|
||||
2: 100,
|
||||
3: 100,
|
||||
4: 100,
|
||||
}),
|
||||
},
|
||||
},
|
||||
{
|
||||
id: "node2",
|
||||
address: volumeServers[1].Address(),
|
||||
diskType: diskType,
|
||||
diskID: 0,
|
||||
ecShards: []*master_pb.VolumeEcShardInformationMessage{
|
||||
buildEcShardInfo(volumeID, collection, diskType, 0, map[uint32]int64{
|
||||
0: 50,
|
||||
5: 100,
|
||||
6: 100,
|
||||
7: 100,
|
||||
8: 100,
|
||||
9: 100,
|
||||
}),
|
||||
},
|
||||
},
|
||||
{
|
||||
id: "node3",
|
||||
address: volumeServers[2].Address(),
|
||||
diskType: diskType,
|
||||
diskID: 0,
|
||||
},
|
||||
{
|
||||
id: "node4",
|
||||
address: volumeServers[3].Address(),
|
||||
diskType: diskType,
|
||||
diskID: 0,
|
||||
},
|
||||
}
|
||||
|
||||
topo := buildTopology(nodes)
|
||||
response := &master_pb.VolumeListResponse{TopologyInfo: topo}
|
||||
master := pluginworkers.NewMasterServer(t, response)
|
||||
|
||||
dialOption := grpc.WithTransportCredentials(insecure.NewCredentials())
|
||||
handler := pluginworker.NewEcRepairHandler(dialOption)
|
||||
harness := pluginworkers.NewHarness(t, pluginworkers.HarnessConfig{
|
||||
WorkerOptions: pluginworker.WorkerOptions{GrpcDialOption: dialOption},
|
||||
Handlers: []pluginworker.JobHandler{handler},
|
||||
})
|
||||
harness.WaitForJobType("ec_repair")
|
||||
|
||||
job := &plugin_pb.JobSpec{
|
||||
JobId: fmt.Sprintf("ec-repair-%d", volumeID),
|
||||
JobType: "ec_repair",
|
||||
Parameters: map[string]*plugin_pb.ConfigValue{
|
||||
"volume_id": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(volumeID)},
|
||||
},
|
||||
"collection": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: collection},
|
||||
},
|
||||
"disk_type": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: diskType},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
result, err := harness.Plugin().ExecuteJob(ctx, job, &plugin_pb.ClusterContext{
|
||||
MasterGrpcAddresses: []string{master.Address()},
|
||||
}, 1)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, result)
|
||||
require.True(t, result.Success)
|
||||
|
||||
rebuildCalls := 0
|
||||
copyCalls := 0
|
||||
deleteCalls := 0
|
||||
unmountCalls := 0
|
||||
for _, vs := range volumeServers {
|
||||
rebuildCalls += len(vs.RebuildRequests())
|
||||
copyCalls += len(vs.CopyRequests())
|
||||
deleteCalls += len(vs.EcDeleteRequests())
|
||||
unmountCalls += len(vs.UnmountRequests())
|
||||
}
|
||||
|
||||
require.Greater(t, rebuildCalls, 0)
|
||||
require.Greater(t, copyCalls, 0)
|
||||
require.Greater(t, deleteCalls, 0)
|
||||
require.Greater(t, unmountCalls, 0)
|
||||
}
|
||||
@@ -0,0 +1,69 @@
|
||||
package ec_repair_test
|
||||
|
||||
import (
|
||||
"sort"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
)
|
||||
|
||||
type nodeSpec struct {
|
||||
id string
|
||||
address string
|
||||
diskType string
|
||||
diskID uint32
|
||||
ecShards []*master_pb.VolumeEcShardInformationMessage
|
||||
}
|
||||
|
||||
func buildTopology(nodes []nodeSpec) *master_pb.TopologyInfo {
|
||||
dataNodes := make([]*master_pb.DataNodeInfo, 0, len(nodes))
|
||||
for _, node := range nodes {
|
||||
diskInfo := &master_pb.DiskInfo{
|
||||
DiskId: node.diskID,
|
||||
MaxVolumeCount: 100,
|
||||
VolumeCount: 10,
|
||||
EcShardInfos: node.ecShards,
|
||||
}
|
||||
dataNodes = append(dataNodes, &master_pb.DataNodeInfo{
|
||||
Id: node.id,
|
||||
Address: node.address,
|
||||
DiskInfos: map[string]*master_pb.DiskInfo{node.diskType: diskInfo},
|
||||
})
|
||||
}
|
||||
return &master_pb.TopologyInfo{
|
||||
DataCenterInfos: []*master_pb.DataCenterInfo{
|
||||
{
|
||||
Id: "dc1",
|
||||
RackInfos: []*master_pb.RackInfo{
|
||||
{
|
||||
Id: "rack1",
|
||||
DataNodeInfos: dataNodes,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func buildEcShardInfo(volumeID uint32, collection, diskType string, diskID uint32, shardSizes map[uint32]int64) *master_pb.VolumeEcShardInformationMessage {
|
||||
shardIDs := make([]int, 0, len(shardSizes))
|
||||
for shardID := range shardSizes {
|
||||
shardIDs = append(shardIDs, int(shardID))
|
||||
}
|
||||
sort.Ints(shardIDs)
|
||||
|
||||
var bits uint32
|
||||
sizes := make([]int64, 0, len(shardIDs))
|
||||
for _, shardID := range shardIDs {
|
||||
bits |= (1 << shardID)
|
||||
sizes = append(sizes, shardSizes[uint32(shardID)])
|
||||
}
|
||||
|
||||
return &master_pb.VolumeEcShardInformationMessage{
|
||||
Id: volumeID,
|
||||
Collection: collection,
|
||||
EcIndexBits: bits,
|
||||
DiskType: diskType,
|
||||
DiskId: diskID,
|
||||
ShardSizes: sizes,
|
||||
}
|
||||
}
|
||||
@@ -30,6 +30,10 @@ type VolumeServer struct {
|
||||
mu sync.Mutex
|
||||
receivedFiles map[string]uint64
|
||||
mountRequests []*volume_server_pb.VolumeEcShardsMountRequest
|
||||
unmountRequests []*volume_server_pb.VolumeEcShardsUnmountRequest
|
||||
copyRequests []*volume_server_pb.VolumeEcShardsCopyRequest
|
||||
rebuildRequests []*volume_server_pb.VolumeEcShardsRebuildRequest
|
||||
ecDeleteRequests []*volume_server_pb.VolumeEcShardsDeleteRequest
|
||||
deleteRequests []*volume_server_pb.VolumeDeleteRequest
|
||||
markReadonlyCalls int
|
||||
vacuumGarbageRatio float64
|
||||
@@ -134,6 +138,46 @@ func (v *VolumeServer) MountRequests() []*volume_server_pb.VolumeEcShardsMountRe
|
||||
return out
|
||||
}
|
||||
|
||||
// UnmountRequests returns recorded unmount requests.
|
||||
func (v *VolumeServer) UnmountRequests() []*volume_server_pb.VolumeEcShardsUnmountRequest {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
out := make([]*volume_server_pb.VolumeEcShardsUnmountRequest, len(v.unmountRequests))
|
||||
copy(out, v.unmountRequests)
|
||||
return out
|
||||
}
|
||||
|
||||
// CopyRequests returns recorded EC shard copy requests.
|
||||
func (v *VolumeServer) CopyRequests() []*volume_server_pb.VolumeEcShardsCopyRequest {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
out := make([]*volume_server_pb.VolumeEcShardsCopyRequest, len(v.copyRequests))
|
||||
copy(out, v.copyRequests)
|
||||
return out
|
||||
}
|
||||
|
||||
// RebuildRequests returns recorded EC shard rebuild requests.
|
||||
func (v *VolumeServer) RebuildRequests() []*volume_server_pb.VolumeEcShardsRebuildRequest {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
out := make([]*volume_server_pb.VolumeEcShardsRebuildRequest, len(v.rebuildRequests))
|
||||
copy(out, v.rebuildRequests)
|
||||
return out
|
||||
}
|
||||
|
||||
// EcDeleteRequests returns recorded EC shard delete requests.
|
||||
func (v *VolumeServer) EcDeleteRequests() []*volume_server_pb.VolumeEcShardsDeleteRequest {
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
out := make([]*volume_server_pb.VolumeEcShardsDeleteRequest, len(v.ecDeleteRequests))
|
||||
copy(out, v.ecDeleteRequests)
|
||||
return out
|
||||
}
|
||||
|
||||
// DeleteRequests returns recorded delete requests.
|
||||
func (v *VolumeServer) DeleteRequests() []*volume_server_pb.VolumeDeleteRequest {
|
||||
v.mu.Lock()
|
||||
@@ -260,6 +304,34 @@ func (v *VolumeServer) VolumeEcShardsMount(ctx context.Context, req *volume_serv
|
||||
return &volume_server_pb.VolumeEcShardsMountResponse{}, nil
|
||||
}
|
||||
|
||||
func (v *VolumeServer) VolumeEcShardsUnmount(ctx context.Context, req *volume_server_pb.VolumeEcShardsUnmountRequest) (*volume_server_pb.VolumeEcShardsUnmountResponse, error) {
|
||||
v.mu.Lock()
|
||||
v.unmountRequests = append(v.unmountRequests, req)
|
||||
v.mu.Unlock()
|
||||
return &volume_server_pb.VolumeEcShardsUnmountResponse{}, nil
|
||||
}
|
||||
|
||||
func (v *VolumeServer) VolumeEcShardsCopy(ctx context.Context, req *volume_server_pb.VolumeEcShardsCopyRequest) (*volume_server_pb.VolumeEcShardsCopyResponse, error) {
|
||||
v.mu.Lock()
|
||||
v.copyRequests = append(v.copyRequests, req)
|
||||
v.mu.Unlock()
|
||||
return &volume_server_pb.VolumeEcShardsCopyResponse{}, nil
|
||||
}
|
||||
|
||||
func (v *VolumeServer) VolumeEcShardsRebuild(ctx context.Context, req *volume_server_pb.VolumeEcShardsRebuildRequest) (*volume_server_pb.VolumeEcShardsRebuildResponse, error) {
|
||||
v.mu.Lock()
|
||||
v.rebuildRequests = append(v.rebuildRequests, req)
|
||||
v.mu.Unlock()
|
||||
return &volume_server_pb.VolumeEcShardsRebuildResponse{}, nil
|
||||
}
|
||||
|
||||
func (v *VolumeServer) VolumeEcShardsDelete(ctx context.Context, req *volume_server_pb.VolumeEcShardsDeleteRequest) (*volume_server_pb.VolumeEcShardsDeleteResponse, error) {
|
||||
v.mu.Lock()
|
||||
v.ecDeleteRequests = append(v.ecDeleteRequests, req)
|
||||
v.mu.Unlock()
|
||||
return &volume_server_pb.VolumeEcShardsDeleteResponse{}, nil
|
||||
}
|
||||
|
||||
func (v *VolumeServer) VolumeDelete(ctx context.Context, req *volume_server_pb.VolumeDeleteRequest) (*volume_server_pb.VolumeDeleteResponse, error) {
|
||||
v.mu.Lock()
|
||||
v.deleteRequests = append(v.deleteRequests, req)
|
||||
|
||||
@@ -0,0 +1,708 @@
|
||||
package pluginworker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/topology"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/worker_pb"
|
||||
ecrepair "github.com/seaweedfs/seaweedfs/weed/worker/tasks/ec_repair"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
const recentTaskWindowSize = 10
|
||||
|
||||
type ecRepairWorkerConfig struct {
|
||||
MinIntervalSeconds int
|
||||
}
|
||||
|
||||
// EcRepairHandler is the plugin job handler for EC shard repair.
|
||||
type EcRepairHandler struct {
|
||||
grpcDialOption grpc.DialOption
|
||||
}
|
||||
|
||||
func NewEcRepairHandler(grpcDialOption grpc.DialOption) *EcRepairHandler {
|
||||
return &EcRepairHandler{grpcDialOption: grpcDialOption}
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) Capability() *plugin_pb.JobTypeCapability {
|
||||
return &plugin_pb.JobTypeCapability{
|
||||
JobType: "ec_repair",
|
||||
CanDetect: true,
|
||||
CanExecute: true,
|
||||
MaxDetectionConcurrency: 1,
|
||||
MaxExecutionConcurrency: 1,
|
||||
DisplayName: "EC Repair",
|
||||
Description: "Repairs missing or inconsistent EC shards",
|
||||
}
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) Descriptor() *plugin_pb.JobTypeDescriptor {
|
||||
return &plugin_pb.JobTypeDescriptor{
|
||||
JobType: "ec_repair",
|
||||
DisplayName: "EC Repair",
|
||||
Description: "Detect and repair missing or inconsistent erasure coding shards",
|
||||
Icon: "fas fa-toolbox",
|
||||
DescriptorVersion: 1,
|
||||
AdminConfigForm: &plugin_pb.ConfigForm{
|
||||
FormId: "ec-repair-admin",
|
||||
Title: "EC Repair Admin Config",
|
||||
Description: "Admin-side controls for EC repair detection scope.",
|
||||
Sections: []*plugin_pb.ConfigSection{
|
||||
{
|
||||
SectionId: "scope",
|
||||
Title: "Scope",
|
||||
Description: "Optional filters applied before EC repair detection.",
|
||||
Fields: []*plugin_pb.ConfigField{
|
||||
{
|
||||
Name: "collection_filter",
|
||||
Label: "Collection Filter",
|
||||
Description: "Only detect EC repairs for this collection when set.",
|
||||
Placeholder: "all collections",
|
||||
FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_STRING,
|
||||
Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_TEXT,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
DefaultValues: map[string]*plugin_pb.ConfigValue{
|
||||
"collection_filter": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: ""},
|
||||
},
|
||||
},
|
||||
},
|
||||
WorkerConfigForm: &plugin_pb.ConfigForm{
|
||||
FormId: "ec-repair-worker",
|
||||
Title: "EC Repair Worker Config",
|
||||
Description: "Worker-side EC repair controls.",
|
||||
Sections: []*plugin_pb.ConfigSection{
|
||||
{
|
||||
SectionId: "interval",
|
||||
Title: "Detection Interval",
|
||||
Description: "Minimum interval between EC repair scans.",
|
||||
Fields: []*plugin_pb.ConfigField{
|
||||
{
|
||||
Name: "min_interval_seconds",
|
||||
Label: "Minimum Detection Interval (s)",
|
||||
Description: "Skip detection if the last successful run is more recent than this interval.",
|
||||
FieldType: plugin_pb.ConfigFieldType_CONFIG_FIELD_TYPE_INT64,
|
||||
Widget: plugin_pb.ConfigWidget_CONFIG_WIDGET_NUMBER,
|
||||
Required: true,
|
||||
MinValue: &plugin_pb.ConfigValue{Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: 0}},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
DefaultValues: map[string]*plugin_pb.ConfigValue{
|
||||
"min_interval_seconds": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: 300},
|
||||
},
|
||||
},
|
||||
},
|
||||
AdminRuntimeDefaults: &plugin_pb.AdminRuntimeDefaults{
|
||||
Enabled: true,
|
||||
DetectionIntervalSeconds: 10 * 60,
|
||||
DetectionTimeoutSeconds: 300,
|
||||
MaxJobsPerDetection: 500,
|
||||
GlobalExecutionConcurrency: 8,
|
||||
PerWorkerExecutionConcurrency: 2,
|
||||
RetryLimit: 1,
|
||||
RetryBackoffSeconds: 30,
|
||||
},
|
||||
WorkerDefaultValues: map[string]*plugin_pb.ConfigValue{
|
||||
"min_interval_seconds": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: 300},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) Detect(ctx context.Context, request *plugin_pb.RunDetectionRequest, sender DetectionSender) error {
|
||||
if request == nil {
|
||||
return fmt.Errorf("run detection request is nil")
|
||||
}
|
||||
if sender == nil {
|
||||
return fmt.Errorf("detection sender is nil")
|
||||
}
|
||||
if request.JobType != "" && request.JobType != "ec_repair" {
|
||||
return fmt.Errorf("job type %q is not handled by ec_repair worker", request.JobType)
|
||||
}
|
||||
|
||||
workerConfig := deriveEcRepairWorkerConfig(request.GetWorkerConfigValues())
|
||||
if shouldSkipDetectionByInterval(request.GetLastSuccessfulRun(), workerConfig.MinIntervalSeconds) {
|
||||
minInterval := time.Duration(workerConfig.MinIntervalSeconds) * time.Second
|
||||
_ = sender.SendActivity(buildDetectorActivity(
|
||||
"skipped_by_interval",
|
||||
fmt.Sprintf("EC REPAIR: Detection skipped due to min interval (%s)", minInterval),
|
||||
map[string]*plugin_pb.ConfigValue{
|
||||
"min_interval_seconds": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(workerConfig.MinIntervalSeconds)},
|
||||
},
|
||||
},
|
||||
))
|
||||
if err := sender.SendProposals(&plugin_pb.DetectionProposals{
|
||||
JobType: "ec_repair",
|
||||
Proposals: []*plugin_pb.JobProposal{},
|
||||
HasMore: false,
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
return sender.SendComplete(&plugin_pb.DetectionComplete{
|
||||
JobType: "ec_repair",
|
||||
Success: true,
|
||||
TotalProposals: 0,
|
||||
})
|
||||
}
|
||||
|
||||
collectionFilter := strings.TrimSpace(readStringConfig(request.GetAdminConfigValues(), "collection_filter", ""))
|
||||
masters := make([]string, 0)
|
||||
if request.ClusterContext != nil {
|
||||
masters = append(masters, request.ClusterContext.MasterGrpcAddresses...)
|
||||
}
|
||||
|
||||
response, _, err := h.fetchTopology(ctx, masters)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
maxResults := int(request.MaxResults)
|
||||
candidates, hasMore, err := ecrepair.Detect(response.TopologyInfo, collectionFilter, maxResults)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
proposals := make([]*plugin_pb.JobProposal, 0, len(candidates))
|
||||
for _, candidate := range candidates {
|
||||
proposal, proposalErr := buildEcRepairProposal(candidate)
|
||||
if proposalErr != nil {
|
||||
glog.Warningf("Plugin worker skip invalid ec_repair proposal: %v", proposalErr)
|
||||
continue
|
||||
}
|
||||
proposals = append(proposals, proposal)
|
||||
}
|
||||
|
||||
if err := sender.SendProposals(&plugin_pb.DetectionProposals{
|
||||
JobType: "ec_repair",
|
||||
Proposals: proposals,
|
||||
HasMore: hasMore,
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return sender.SendComplete(&plugin_pb.DetectionComplete{
|
||||
JobType: "ec_repair",
|
||||
Success: true,
|
||||
TotalProposals: int32(len(proposals)),
|
||||
})
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) Execute(ctx context.Context, request *plugin_pb.ExecuteJobRequest, sender ExecutionSender) error {
|
||||
if request == nil || request.Job == nil {
|
||||
return fmt.Errorf("execute request/job is nil")
|
||||
}
|
||||
if sender == nil {
|
||||
return fmt.Errorf("execution sender is nil")
|
||||
}
|
||||
if request.Job.JobType != "" && request.Job.JobType != "ec_repair" {
|
||||
return fmt.Errorf("job type %q is not handled by ec_repair worker", request.Job.JobType)
|
||||
}
|
||||
|
||||
volumeID := readInt64Config(request.Job.Parameters, "volume_id", 0)
|
||||
if volumeID <= 0 {
|
||||
return fmt.Errorf("missing volume_id in job parameters")
|
||||
}
|
||||
collection := readStringConfig(request.Job.Parameters, "collection", "")
|
||||
diskType := readStringConfig(request.Job.Parameters, "disk_type", "")
|
||||
|
||||
masters := make([]string, 0)
|
||||
if request.ClusterContext != nil {
|
||||
masters = append(masters, request.ClusterContext.MasterGrpcAddresses...)
|
||||
}
|
||||
|
||||
if err := sendProgress(sender, request.Job, 0, "assigned", "ec repair job accepted"); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
response, activeTopology, err := h.fetchTopology(ctx, masters)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
plan, err := ecrepair.BuildRepairPlan(response.TopologyInfo, activeTopology, uint32(volumeID), collection, diskType)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if len(plan.MissingShards) == 0 && len(plan.DeleteByNode) == 0 {
|
||||
resultSummary := fmt.Sprintf("EC repair skipped for volume %d (no issues found)", volumeID)
|
||||
return sender.SendCompleted(&plugin_pb.JobCompleted{
|
||||
JobId: request.Job.JobId,
|
||||
JobType: request.Job.JobType,
|
||||
Success: true,
|
||||
Result: &plugin_pb.JobResult{
|
||||
Summary: resultSummary,
|
||||
OutputValues: map[string]*plugin_pb.ConfigValue{
|
||||
"volume_id": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: volumeID},
|
||||
},
|
||||
},
|
||||
},
|
||||
Activities: []*plugin_pb.ActivityEvent{
|
||||
buildExecutorActivity("completed", resultSummary),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
if err := h.executeRepairPlan(ctx, plan, request.Job, sender); err != nil {
|
||||
_ = sendProgress(sender, request.Job, 100, "failed", err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
resultSummary := fmt.Sprintf("EC repair completed for volume %d", volumeID)
|
||||
return sender.SendCompleted(&plugin_pb.JobCompleted{
|
||||
JobId: request.Job.JobId,
|
||||
JobType: request.Job.JobType,
|
||||
Success: true,
|
||||
Result: &plugin_pb.JobResult{
|
||||
Summary: resultSummary,
|
||||
OutputValues: map[string]*plugin_pb.ConfigValue{
|
||||
"volume_id": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: volumeID},
|
||||
},
|
||||
"collection": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: plan.Collection},
|
||||
},
|
||||
"disk_type": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: plan.DiskType},
|
||||
},
|
||||
},
|
||||
},
|
||||
Activities: []*plugin_pb.ActivityEvent{
|
||||
buildExecutorActivity("completed", resultSummary),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) executeRepairPlan(ctx context.Context, plan *ecrepair.RepairPlan, job *plugin_pb.JobSpec, sender ExecutionSender) error {
|
||||
if plan == nil {
|
||||
return fmt.Errorf("repair plan is nil")
|
||||
}
|
||||
|
||||
if len(plan.MissingShards) > 0 {
|
||||
if err := sendProgress(sender, job, 10, "copy_existing_shards", "copying EC shards to rebuilder"); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := h.copyShardsToRebuilder(ctx, plan); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := sendProgress(sender, job, 40, "rebuild_missing_shards", "rebuilding missing EC shards"); err != nil {
|
||||
return err
|
||||
}
|
||||
rebuilt, err := h.rebuildMissingShards(ctx, plan)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := sendProgress(sender, job, 60, "distribute_shards", "distributing rebuilt EC shards"); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := h.distributeRebuiltShards(ctx, plan, rebuilt); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := sendProgress(sender, job, 75, "cleanup_rebuilder", "cleaning up rebuilder temporary shards"); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := h.cleanupRebuilder(ctx, plan, rebuilt); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
if len(plan.DeleteByNode) > 0 {
|
||||
if err := sendProgress(sender, job, 85, "delete_extra_shards", "deleting extra or mismatched shards"); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := h.deleteExtraShards(ctx, plan); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return sendProgress(sender, job, 100, "completed", "ec repair completed")
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) copyShardsToRebuilder(ctx context.Context, plan *ecrepair.RepairPlan) error {
|
||||
if plan.Rebuilder.NodeAddress == "" {
|
||||
return fmt.Errorf("rebuilder node is required")
|
||||
}
|
||||
if len(plan.CopySources) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
localShardSet := make(map[uint32]struct{}, len(plan.Rebuilder.LocalShards))
|
||||
for _, shardID := range plan.Rebuilder.LocalShards {
|
||||
localShardSet[shardID] = struct{}{}
|
||||
}
|
||||
|
||||
shardIDs := make([]uint32, 0, len(plan.CopySources))
|
||||
for shardID := range plan.CopySources {
|
||||
shardIDs = append(shardIDs, shardID)
|
||||
}
|
||||
sort.Slice(shardIDs, func(i, j int) bool { return shardIDs[i] < shardIDs[j] })
|
||||
|
||||
copyIndexFiles := true
|
||||
for _, shardID := range shardIDs {
|
||||
source := strings.TrimSpace(plan.CopySources[shardID])
|
||||
if source == "" {
|
||||
continue
|
||||
}
|
||||
if source == plan.Rebuilder.NodeAddress {
|
||||
if _, ok := localShardSet[shardID]; ok {
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
err := operation.WithVolumeServerClient(false, pb.ServerAddress(plan.Rebuilder.NodeAddress), h.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
|
||||
_, err := client.VolumeEcShardsCopy(ctx, &volume_server_pb.VolumeEcShardsCopyRequest{
|
||||
VolumeId: plan.VolumeID,
|
||||
Collection: plan.Collection,
|
||||
ShardIds: []uint32{shardID},
|
||||
CopyEcxFile: copyIndexFiles,
|
||||
CopyEcjFile: copyIndexFiles,
|
||||
CopyVifFile: copyIndexFiles,
|
||||
SourceDataNode: source,
|
||||
DiskId: plan.Rebuilder.DiskID,
|
||||
})
|
||||
return err
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
copyIndexFiles = false
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) rebuildMissingShards(ctx context.Context, plan *ecrepair.RepairPlan) ([]uint32, error) {
|
||||
var rebuilt []uint32
|
||||
if plan.Rebuilder.NodeAddress == "" {
|
||||
return nil, fmt.Errorf("rebuilder node is required")
|
||||
}
|
||||
if err := operation.WithVolumeServerClient(false, pb.ServerAddress(plan.Rebuilder.NodeAddress), h.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
|
||||
resp, err := client.VolumeEcShardsRebuild(ctx, &volume_server_pb.VolumeEcShardsRebuildRequest{
|
||||
VolumeId: plan.VolumeID,
|
||||
Collection: plan.Collection,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rebuilt = append(rebuilt, resp.RebuiltShardIds...)
|
||||
return nil
|
||||
}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(rebuilt) == 0 {
|
||||
if len(plan.MissingShards) > 0 {
|
||||
glog.V(1).Infof("EC Repair: resp.RebuiltShardIds empty; assuming all missing shards rebuilt for volume %d (%d shards)", plan.VolumeID, len(plan.MissingShards))
|
||||
} else {
|
||||
glog.V(1).Infof("EC Repair: resp.RebuiltShardIds empty for volume %d with no missing shards declared", plan.VolumeID)
|
||||
}
|
||||
rebuilt = append(rebuilt, plan.MissingShards...)
|
||||
}
|
||||
return rebuilt, nil
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) distributeRebuiltShards(ctx context.Context, plan *ecrepair.RepairPlan, rebuilt []uint32) error {
|
||||
if len(plan.Targets) == 0 {
|
||||
return nil
|
||||
}
|
||||
if len(rebuilt) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
rebuiltSet := make(map[uint32]struct{}, len(rebuilt))
|
||||
for _, shardID := range rebuilt {
|
||||
rebuiltSet[shardID] = struct{}{}
|
||||
}
|
||||
|
||||
targetCopyIndex := make(map[string]bool)
|
||||
for _, target := range plan.Targets {
|
||||
if target.NodeAddress == "" {
|
||||
continue
|
||||
}
|
||||
var shardIDs []uint32
|
||||
for _, shardID := range target.ShardIDs {
|
||||
if _, ok := rebuiltSet[shardID]; ok {
|
||||
shardIDs = append(shardIDs, shardID)
|
||||
}
|
||||
}
|
||||
if len(shardIDs) == 0 {
|
||||
continue
|
||||
}
|
||||
sort.Slice(shardIDs, func(i, j int) bool { return shardIDs[i] < shardIDs[j] })
|
||||
|
||||
copyIndex := !targetCopyIndex[target.NodeAddress]
|
||||
targetCopyIndex[target.NodeAddress] = true
|
||||
|
||||
err := operation.WithVolumeServerClient(false, pb.ServerAddress(target.NodeAddress), h.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
|
||||
if target.NodeAddress != plan.Rebuilder.NodeAddress {
|
||||
_, err := client.VolumeEcShardsCopy(ctx, &volume_server_pb.VolumeEcShardsCopyRequest{
|
||||
VolumeId: plan.VolumeID,
|
||||
Collection: plan.Collection,
|
||||
ShardIds: shardIDs,
|
||||
CopyEcxFile: copyIndex,
|
||||
CopyEcjFile: copyIndex,
|
||||
CopyVifFile: copyIndex,
|
||||
SourceDataNode: plan.Rebuilder.NodeAddress,
|
||||
DiskId: target.DiskID,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
_, err := client.VolumeEcShardsMount(ctx, &volume_server_pb.VolumeEcShardsMountRequest{
|
||||
VolumeId: plan.VolumeID,
|
||||
Collection: plan.Collection,
|
||||
ShardIds: shardIDs,
|
||||
})
|
||||
return err
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) cleanupRebuilder(ctx context.Context, plan *ecrepair.RepairPlan, rebuilt []uint32) error {
|
||||
if plan.Rebuilder.NodeAddress == "" || len(rebuilt) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
keep := make(map[uint32]struct{})
|
||||
for _, shardID := range plan.Rebuilder.LocalShards {
|
||||
keep[shardID] = struct{}{}
|
||||
}
|
||||
for _, target := range plan.Targets {
|
||||
if target.NodeAddress != plan.Rebuilder.NodeAddress {
|
||||
continue
|
||||
}
|
||||
for _, shardID := range target.ShardIDs {
|
||||
keep[shardID] = struct{}{}
|
||||
}
|
||||
}
|
||||
|
||||
var toDelete []uint32
|
||||
for _, shardID := range rebuilt {
|
||||
if _, ok := keep[shardID]; ok {
|
||||
continue
|
||||
}
|
||||
toDelete = append(toDelete, shardID)
|
||||
}
|
||||
if len(toDelete) == 0 {
|
||||
return nil
|
||||
}
|
||||
return deleteShardIds(ctx, h.grpcDialOption, plan.Rebuilder.NodeAddress, plan.VolumeID, plan.Collection, toDelete)
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) deleteExtraShards(ctx context.Context, plan *ecrepair.RepairPlan) error {
|
||||
for nodeAddress, shardIDs := range plan.DeleteByNode {
|
||||
if len(shardIDs) == 0 {
|
||||
continue
|
||||
}
|
||||
if err := deleteShardIds(ctx, h.grpcDialOption, nodeAddress, plan.VolumeID, plan.Collection, shardIDs); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func deleteShardIds(ctx context.Context, dialOption grpc.DialOption, nodeAddress string, volumeID uint32, collection string, shardIDs []uint32) error {
|
||||
sorted := make([]uint32, len(shardIDs))
|
||||
copy(sorted, shardIDs)
|
||||
sort.Slice(sorted, func(i, j int) bool { return sorted[i] < sorted[j] })
|
||||
|
||||
return operation.WithVolumeServerClient(false, pb.ServerAddress(nodeAddress), dialOption, func(client volume_server_pb.VolumeServerClient) error {
|
||||
_, err := client.VolumeEcShardsUnmount(ctx, &volume_server_pb.VolumeEcShardsUnmountRequest{
|
||||
VolumeId: volumeID,
|
||||
ShardIds: sorted,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = client.VolumeEcShardsDelete(ctx, &volume_server_pb.VolumeEcShardsDeleteRequest{
|
||||
VolumeId: volumeID,
|
||||
Collection: collection,
|
||||
ShardIds: sorted,
|
||||
})
|
||||
return err
|
||||
})
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) fetchTopology(ctx context.Context, masterAddresses []string) (*master_pb.VolumeListResponse, *topology.ActiveTopology, error) {
|
||||
if h.grpcDialOption == nil {
|
||||
return nil, nil, fmt.Errorf("grpc dial option is not configured")
|
||||
}
|
||||
if len(masterAddresses) == 0 {
|
||||
return nil, nil, fmt.Errorf("no master addresses provided in cluster context")
|
||||
}
|
||||
|
||||
for _, masterAddress := range masterAddresses {
|
||||
response, err := h.fetchVolumeList(ctx, masterAddress)
|
||||
if err != nil {
|
||||
glog.Warningf("Plugin worker failed master volume list at %s: %v", masterAddress, err)
|
||||
continue
|
||||
}
|
||||
if response == nil || response.TopologyInfo == nil {
|
||||
continue
|
||||
}
|
||||
activeTopology := topology.NewActiveTopology(recentTaskWindowSize)
|
||||
if err := activeTopology.UpdateTopology(response.TopologyInfo); err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
return response, activeTopology, nil
|
||||
}
|
||||
|
||||
return nil, nil, fmt.Errorf("failed to load topology from all provided masters")
|
||||
}
|
||||
|
||||
func (h *EcRepairHandler) fetchVolumeList(ctx context.Context, address string) (*master_pb.VolumeListResponse, error) {
|
||||
var lastErr error
|
||||
for _, candidate := range masterAddressCandidates(address) {
|
||||
if ctx.Err() != nil {
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
|
||||
dialCtx, cancelDial := context.WithTimeout(ctx, 5*time.Second)
|
||||
conn, err := pb.GrpcDial(dialCtx, candidate, false, h.grpcDialOption)
|
||||
cancelDial()
|
||||
if err != nil {
|
||||
lastErr = err
|
||||
continue
|
||||
}
|
||||
|
||||
client := master_pb.NewSeaweedClient(conn)
|
||||
callCtx, cancelCall := context.WithTimeout(ctx, 10*time.Second)
|
||||
response, callErr := client.VolumeList(callCtx, &master_pb.VolumeListRequest{})
|
||||
cancelCall()
|
||||
_ = conn.Close()
|
||||
|
||||
if callErr == nil {
|
||||
return response, nil
|
||||
}
|
||||
lastErr = callErr
|
||||
}
|
||||
|
||||
if lastErr == nil {
|
||||
lastErr = fmt.Errorf("no valid master address candidate")
|
||||
}
|
||||
return nil, lastErr
|
||||
}
|
||||
|
||||
func deriveEcRepairWorkerConfig(values map[string]*plugin_pb.ConfigValue) *ecRepairWorkerConfig {
|
||||
return &ecRepairWorkerConfig{
|
||||
MinIntervalSeconds: int(readInt64Config(values, "min_interval_seconds", 300)),
|
||||
}
|
||||
}
|
||||
|
||||
func buildEcRepairProposal(candidate *ecrepair.RepairCandidate) (*plugin_pb.JobProposal, error) {
|
||||
if candidate == nil {
|
||||
return nil, fmt.Errorf("repair candidate is nil")
|
||||
}
|
||||
|
||||
params := &worker_pb.TaskParams{
|
||||
VolumeId: candidate.VolumeID,
|
||||
Collection: candidate.Collection,
|
||||
}
|
||||
payload, err := proto.Marshal(params)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal task params: %w", err)
|
||||
}
|
||||
|
||||
proposalID := fmt.Sprintf("ec-repair-%d-%d", candidate.VolumeID, time.Now().UnixNano())
|
||||
dedupeKey := fmt.Sprintf("ec_repair:%d", candidate.VolumeID)
|
||||
if candidate.Collection != "" {
|
||||
dedupeKey = dedupeKey + ":" + candidate.Collection
|
||||
}
|
||||
if candidate.DiskType != "" {
|
||||
dedupeKey = dedupeKey + ":" + candidate.DiskType
|
||||
}
|
||||
|
||||
summary := fmt.Sprintf("Repair EC volume %d", candidate.VolumeID)
|
||||
if candidate.Collection != "" {
|
||||
summary = summary + " (" + candidate.Collection + ")"
|
||||
}
|
||||
|
||||
detail := fmt.Sprintf("missing shards=%d, extra shards=%d, mismatched shards=%d", candidate.MissingShards, candidate.ExtraShards, candidate.MismatchedShards)
|
||||
|
||||
return &plugin_pb.JobProposal{
|
||||
ProposalId: proposalID,
|
||||
DedupeKey: dedupeKey,
|
||||
JobType: "ec_repair",
|
||||
Priority: plugin_pb.JobPriority_JOB_PRIORITY_NORMAL,
|
||||
Summary: summary,
|
||||
Detail: detail,
|
||||
Parameters: map[string]*plugin_pb.ConfigValue{
|
||||
"task_params_pb": {
|
||||
Kind: &plugin_pb.ConfigValue_BytesValue{BytesValue: payload},
|
||||
},
|
||||
"volume_id": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(candidate.VolumeID)},
|
||||
},
|
||||
"collection": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: candidate.Collection},
|
||||
},
|
||||
"disk_type": {
|
||||
Kind: &plugin_pb.ConfigValue_StringValue{StringValue: candidate.DiskType},
|
||||
},
|
||||
"missing_shards": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(candidate.MissingShards)},
|
||||
},
|
||||
"extra_shards": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(candidate.ExtraShards)},
|
||||
},
|
||||
"mismatched_shards": {
|
||||
Kind: &plugin_pb.ConfigValue_Int64Value{Int64Value: int64(candidate.MismatchedShards)},
|
||||
},
|
||||
},
|
||||
Labels: map[string]string{
|
||||
"task_type": "ec_repair",
|
||||
"volume_id": fmt.Sprintf("%d", candidate.VolumeID),
|
||||
"collection": candidate.Collection,
|
||||
"disk_type": candidate.DiskType,
|
||||
"missing_shards": fmt.Sprintf("%d", candidate.MissingShards),
|
||||
"extra_shards": fmt.Sprintf("%d", candidate.ExtraShards),
|
||||
"mismatched_shards": fmt.Sprintf("%d", candidate.MismatchedShards),
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func sendProgress(sender ExecutionSender, job *plugin_pb.JobSpec, percent float64, stage string, message string) error {
|
||||
if sender == nil || job == nil {
|
||||
return nil
|
||||
}
|
||||
return sender.SendProgress(&plugin_pb.JobProgressUpdate{
|
||||
JobId: job.JobId,
|
||||
JobType: job.JobType,
|
||||
State: plugin_pb.JobState_JOB_STATE_RUNNING,
|
||||
ProgressPercent: percent,
|
||||
Stage: stage,
|
||||
Message: message,
|
||||
Activities: []*plugin_pb.ActivityEvent{
|
||||
buildExecutorActivity(stage, message),
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,674 @@
|
||||
package ec_repair
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/topology"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding/placement"
|
||||
"github.com/seaweedfs/seaweedfs/weed/worker/tasks/util"
|
||||
)
|
||||
|
||||
type volumeShardState struct {
|
||||
Key VolumeKey
|
||||
ShardMap map[uint32][]ShardLocation
|
||||
MaxShardID uint32
|
||||
SeenTargets map[string]struct{}
|
||||
DataShards int
|
||||
ParityShards int
|
||||
}
|
||||
|
||||
type shardAnalysis struct {
|
||||
Missing []uint32
|
||||
Keep map[uint32]ShardLocation
|
||||
Delete []ShardLocation
|
||||
Mismatched []ShardLocation
|
||||
TotalShards int
|
||||
}
|
||||
|
||||
func sortVolumeKeys(keys []VolumeKey) {
|
||||
sort.Slice(keys, func(i, j int) bool {
|
||||
if keys[i].VolumeID != keys[j].VolumeID {
|
||||
return keys[i].VolumeID < keys[j].VolumeID
|
||||
}
|
||||
if keys[i].Collection != keys[j].Collection {
|
||||
return keys[i].Collection < keys[j].Collection
|
||||
}
|
||||
return keys[i].DiskType < keys[j].DiskType
|
||||
})
|
||||
}
|
||||
|
||||
// Detect scans topology for EC shard issues and returns repair candidates.
|
||||
func Detect(topoInfo *master_pb.TopologyInfo, collectionFilter string, maxResults int) ([]*RepairCandidate, bool, error) {
|
||||
if topoInfo == nil {
|
||||
return nil, false, fmt.Errorf("topology info is nil")
|
||||
}
|
||||
|
||||
states := collectShardStates(topoInfo, collectionFilter)
|
||||
keys := make([]VolumeKey, 0, len(states))
|
||||
for key := range states {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
sortVolumeKeys(keys)
|
||||
|
||||
if maxResults < 0 {
|
||||
maxResults = 0
|
||||
}
|
||||
|
||||
var candidates []*RepairCandidate
|
||||
hasMore := false
|
||||
|
||||
for idx, key := range keys {
|
||||
state := states[key]
|
||||
if state == nil {
|
||||
continue
|
||||
}
|
||||
analysis := analyzeShardState(state)
|
||||
if len(analysis.Missing) == 0 && len(analysis.Delete) == 0 {
|
||||
continue
|
||||
}
|
||||
if len(analysis.Missing) > 0 && len(analysis.Keep) < state.DataShards {
|
||||
// Not enough shards to rebuild missing shards safely.
|
||||
continue
|
||||
}
|
||||
|
||||
candidates = append(candidates, &RepairCandidate{
|
||||
VolumeID: key.VolumeID,
|
||||
Collection: key.Collection,
|
||||
DiskType: key.DiskType,
|
||||
MissingShards: len(analysis.Missing),
|
||||
ExtraShards: len(analysis.Delete),
|
||||
MismatchedShards: len(analysis.Mismatched),
|
||||
})
|
||||
|
||||
if maxResults > 0 && len(candidates) >= maxResults {
|
||||
hasMore = idx < len(keys)-1
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
return candidates, hasMore, nil
|
||||
}
|
||||
|
||||
// BuildRepairPlan creates a concrete repair plan for one EC volume.
|
||||
func BuildRepairPlan(
|
||||
topoInfo *master_pb.TopologyInfo,
|
||||
activeTopology *topology.ActiveTopology,
|
||||
volumeID uint32,
|
||||
collection string,
|
||||
diskType string,
|
||||
) (*RepairPlan, error) {
|
||||
if topoInfo == nil {
|
||||
return nil, fmt.Errorf("topology info is nil")
|
||||
}
|
||||
states := collectShardStates(topoInfo, "")
|
||||
|
||||
keys := make([]VolumeKey, 0, len(states))
|
||||
for key := range states {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
sortVolumeKeys(keys)
|
||||
|
||||
var state *volumeShardState
|
||||
for _, key := range keys {
|
||||
if key.VolumeID != volumeID {
|
||||
continue
|
||||
}
|
||||
if collection != "" && key.Collection != collection {
|
||||
continue
|
||||
}
|
||||
if diskType != "" && key.DiskType != diskType {
|
||||
continue
|
||||
}
|
||||
state = states[key]
|
||||
break
|
||||
}
|
||||
if state == nil {
|
||||
return nil, fmt.Errorf("ec volume %d not found in topology", volumeID)
|
||||
}
|
||||
|
||||
analysis := analyzeShardState(state)
|
||||
if len(analysis.Missing) == 0 && len(analysis.Delete) == 0 {
|
||||
return &RepairPlan{
|
||||
VolumeID: volumeID,
|
||||
Collection: state.Key.Collection,
|
||||
DiskType: state.Key.DiskType,
|
||||
DeleteByNode: map[string][]uint32{},
|
||||
}, nil
|
||||
}
|
||||
if len(analysis.Missing) > 0 && len(analysis.Keep) < state.DataShards {
|
||||
return nil, fmt.Errorf("ec volume %d has %d shards, need at least %d to rebuild", volumeID, len(analysis.Keep), state.DataShards)
|
||||
}
|
||||
|
||||
deleteByNode := groupDeleteByNode(analysis.Delete, activeTopology)
|
||||
|
||||
plan := &RepairPlan{
|
||||
VolumeID: volumeID,
|
||||
Collection: state.Key.Collection,
|
||||
DiskType: state.Key.DiskType,
|
||||
MissingShards: analysis.Missing,
|
||||
DeleteByNode: deleteByNode,
|
||||
CopySources: make(map[uint32]string),
|
||||
}
|
||||
|
||||
if len(analysis.Keep) > 0 {
|
||||
rebuilder := selectRebuilder(analysis.Keep)
|
||||
plan.Rebuilder = rebuilder
|
||||
plan.CopySources = buildCopySources(analysis.Keep)
|
||||
}
|
||||
|
||||
if len(analysis.Missing) > 0 {
|
||||
if activeTopology == nil {
|
||||
return nil, fmt.Errorf("active topology is required to place missing shards")
|
||||
}
|
||||
usedNodes := collectUsedNodes(analysis.Keep)
|
||||
targets, err := planMissingTargets(activeTopology, analysis.Missing, usedNodes, state.Key.DiskType)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
plan.Targets = targets
|
||||
}
|
||||
|
||||
return plan, nil
|
||||
}
|
||||
|
||||
func collectShardStates(topoInfo *master_pb.TopologyInfo, collectionFilter string) map[VolumeKey]*volumeShardState {
|
||||
states := make(map[VolumeKey]*volumeShardState)
|
||||
filter := strings.TrimSpace(collectionFilter)
|
||||
|
||||
for _, dc := range topoInfo.DataCenterInfos {
|
||||
for _, rack := range dc.RackInfos {
|
||||
for _, node := range rack.DataNodeInfos {
|
||||
nodeAddress := string(pb.NewServerAddressFromDataNode(node))
|
||||
for diskTypeKey, diskInfo := range node.DiskInfos {
|
||||
if diskInfo == nil {
|
||||
continue
|
||||
}
|
||||
for _, shardInfo := range diskInfo.EcShardInfos {
|
||||
if shardInfo == nil {
|
||||
continue
|
||||
}
|
||||
if filter != "" && shardInfo.Collection != filter {
|
||||
continue
|
||||
}
|
||||
|
||||
recordedDiskType := strings.TrimSpace(shardInfo.DiskType)
|
||||
if recordedDiskType == "" {
|
||||
recordedDiskType = diskTypeKey
|
||||
}
|
||||
|
||||
key := VolumeKey{
|
||||
VolumeID: shardInfo.Id,
|
||||
Collection: shardInfo.Collection,
|
||||
DiskType: recordedDiskType,
|
||||
}
|
||||
state := states[key]
|
||||
if state == nil {
|
||||
state = &volumeShardState{
|
||||
Key: key,
|
||||
ShardMap: make(map[uint32][]ShardLocation),
|
||||
SeenTargets: make(map[string]struct{}),
|
||||
DataShards: erasure_coding.DataShardsCount,
|
||||
ParityShards: erasure_coding.ParityShardsCount,
|
||||
}
|
||||
states[key] = state
|
||||
}
|
||||
|
||||
shardsInfo := erasure_coding.ShardsInfoFromVolumeEcShardInformationMessage(shardInfo)
|
||||
for _, shard := range shardsInfo.AsSlice() {
|
||||
shardID := uint32(shard.Id)
|
||||
if shardID > state.MaxShardID {
|
||||
state.MaxShardID = shardID
|
||||
}
|
||||
locationKey := fmt.Sprintf("%s:%d:%d", node.Id, shardInfo.DiskId, shardID)
|
||||
if _, seen := state.SeenTargets[locationKey]; seen {
|
||||
continue
|
||||
}
|
||||
state.SeenTargets[locationKey] = struct{}{}
|
||||
|
||||
state.ShardMap[shardID] = append(state.ShardMap[shardID], ShardLocation{
|
||||
NodeID: node.Id,
|
||||
NodeAddress: nodeAddress,
|
||||
DataCenter: dc.Id,
|
||||
Rack: rack.Id,
|
||||
DiskType: recordedDiskType,
|
||||
DiskID: shardInfo.DiskId,
|
||||
ShardID: shardID,
|
||||
Size: int64(shard.Size),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return states
|
||||
}
|
||||
|
||||
func analyzeShardState(state *volumeShardState) shardAnalysis {
|
||||
analysis := shardAnalysis{
|
||||
Keep: make(map[uint32]ShardLocation),
|
||||
}
|
||||
if state == nil {
|
||||
return analysis
|
||||
}
|
||||
|
||||
expectedTotal := state.DataShards + state.ParityShards
|
||||
if expectedTotal <= 0 {
|
||||
expectedTotal = erasure_coding.TotalShardsCount
|
||||
}
|
||||
if int(state.MaxShardID)+1 > expectedTotal && int(state.MaxShardID)+1 <= erasure_coding.MaxShardCount {
|
||||
expectedTotal = int(state.MaxShardID) + 1
|
||||
}
|
||||
analysis.TotalShards = expectedTotal
|
||||
|
||||
candidates := make(map[uint32][]ShardLocation, expectedTotal)
|
||||
var mismatched []ShardLocation
|
||||
for shardID := 0; shardID < expectedTotal; shardID++ {
|
||||
locations := state.ShardMap[uint32(shardID)]
|
||||
if len(locations) == 0 {
|
||||
analysis.Missing = append(analysis.Missing, uint32(shardID))
|
||||
continue
|
||||
}
|
||||
filtered, outliers := splitByCanonicalSize(locations)
|
||||
mismatched = append(mismatched, outliers...)
|
||||
candidates[uint32(shardID)] = filtered
|
||||
}
|
||||
|
||||
analysis.Mismatched = mismatched
|
||||
keep := selectDiverseLocations(candidates)
|
||||
analysis.Keep = keep
|
||||
|
||||
var deletions []ShardLocation
|
||||
for shardID, locations := range candidates {
|
||||
keptLocation, ok := keep[shardID]
|
||||
for _, location := range locations {
|
||||
if ok && location.NodeID == keptLocation.NodeID && location.DiskID == keptLocation.DiskID {
|
||||
continue
|
||||
}
|
||||
deletions = append(deletions, location)
|
||||
}
|
||||
}
|
||||
deletions = append(deletions, mismatched...)
|
||||
analysis.Delete = dedupeLocations(deletions)
|
||||
|
||||
return analysis
|
||||
}
|
||||
|
||||
func splitByCanonicalSize(locations []ShardLocation) ([]ShardLocation, []ShardLocation) {
|
||||
if len(locations) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
sizeCounts := make(map[int64]int)
|
||||
for _, location := range locations {
|
||||
if location.Size > 0 {
|
||||
sizeCounts[location.Size]++
|
||||
}
|
||||
}
|
||||
|
||||
canonicalSize := int64(0)
|
||||
if len(sizeCounts) > 0 {
|
||||
maxCount := -1
|
||||
for size, count := range sizeCounts {
|
||||
if count > maxCount || (count == maxCount && size > canonicalSize) {
|
||||
maxCount = count
|
||||
canonicalSize = size
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
var filtered []ShardLocation
|
||||
var outliers []ShardLocation
|
||||
for _, location := range locations {
|
||||
if len(sizeCounts) == 0 {
|
||||
filtered = append(filtered, location)
|
||||
continue
|
||||
}
|
||||
if location.Size == canonicalSize {
|
||||
filtered = append(filtered, location)
|
||||
continue
|
||||
}
|
||||
outliers = append(outliers, location)
|
||||
}
|
||||
if len(filtered) == 0 {
|
||||
filtered = append(filtered, locations...)
|
||||
outliers = nil
|
||||
}
|
||||
return filtered, outliers
|
||||
}
|
||||
|
||||
func selectDiverseLocations(candidates map[uint32][]ShardLocation) map[uint32]ShardLocation {
|
||||
keep := make(map[uint32]ShardLocation, len(candidates))
|
||||
if len(candidates) == 0 {
|
||||
return keep
|
||||
}
|
||||
|
||||
shardIDs := make([]uint32, 0, len(candidates))
|
||||
for shardID := range candidates {
|
||||
shardIDs = append(shardIDs, shardID)
|
||||
}
|
||||
sort.Slice(shardIDs, func(i, j int) bool {
|
||||
li := len(candidates[shardIDs[i]])
|
||||
lj := len(candidates[shardIDs[j]])
|
||||
if li != lj {
|
||||
return li < lj
|
||||
}
|
||||
return shardIDs[i] < shardIDs[j]
|
||||
})
|
||||
|
||||
dcCount := make(map[string]int)
|
||||
rackCount := make(map[string]int)
|
||||
nodeCount := make(map[string]int)
|
||||
|
||||
for _, shardID := range shardIDs {
|
||||
locations := candidates[shardID]
|
||||
if len(locations) == 0 {
|
||||
continue
|
||||
}
|
||||
best := locations[0]
|
||||
bestScore := scoreLocation(best, dcCount, rackCount, nodeCount)
|
||||
for _, location := range locations[1:] {
|
||||
score := scoreLocation(location, dcCount, rackCount, nodeCount)
|
||||
if scoreLess(score, bestScore) {
|
||||
best = location
|
||||
bestScore = score
|
||||
}
|
||||
}
|
||||
keep[shardID] = best
|
||||
|
||||
dcCount[best.DataCenter]++
|
||||
rackKey := best.DataCenter + ":" + best.Rack
|
||||
rackCount[rackKey]++
|
||||
nodeCount[best.NodeID]++
|
||||
}
|
||||
|
||||
return keep
|
||||
}
|
||||
|
||||
type locationScore struct {
|
||||
dcCount int
|
||||
rackCount int
|
||||
nodeCount int
|
||||
nodeID string
|
||||
diskID uint32
|
||||
}
|
||||
|
||||
func scoreLocation(loc ShardLocation, dcCount, rackCount, nodeCount map[string]int) locationScore {
|
||||
rackKey := loc.DataCenter + ":" + loc.Rack
|
||||
return locationScore{
|
||||
dcCount: dcCount[loc.DataCenter],
|
||||
rackCount: rackCount[rackKey],
|
||||
nodeCount: nodeCount[loc.NodeID],
|
||||
nodeID: loc.NodeID,
|
||||
diskID: loc.DiskID,
|
||||
}
|
||||
}
|
||||
|
||||
func scoreLess(a, b locationScore) bool {
|
||||
if a.dcCount != b.dcCount {
|
||||
return a.dcCount < b.dcCount
|
||||
}
|
||||
if a.rackCount != b.rackCount {
|
||||
return a.rackCount < b.rackCount
|
||||
}
|
||||
if a.nodeCount != b.nodeCount {
|
||||
return a.nodeCount < b.nodeCount
|
||||
}
|
||||
if a.nodeID != b.nodeID {
|
||||
return a.nodeID < b.nodeID
|
||||
}
|
||||
return a.diskID < b.diskID
|
||||
}
|
||||
|
||||
func dedupeLocations(locations []ShardLocation) []ShardLocation {
|
||||
if len(locations) == 0 {
|
||||
return nil
|
||||
}
|
||||
seen := make(map[string]struct{}, len(locations))
|
||||
out := make([]ShardLocation, 0, len(locations))
|
||||
for _, loc := range locations {
|
||||
key := fmt.Sprintf("%s:%d:%d", loc.NodeID, loc.DiskID, loc.ShardID)
|
||||
if _, ok := seen[key]; ok {
|
||||
continue
|
||||
}
|
||||
seen[key] = struct{}{}
|
||||
out = append(out, loc)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func groupDeleteByNode(locations []ShardLocation, activeTopology *topology.ActiveTopology) map[string][]uint32 {
|
||||
result := make(map[string]map[uint32]struct{})
|
||||
for _, loc := range locations {
|
||||
nodeAddress := strings.TrimSpace(loc.NodeAddress)
|
||||
if nodeAddress == "" && activeTopology != nil {
|
||||
resolved, err := util.ResolveServerAddress(loc.NodeID, activeTopology)
|
||||
if err == nil {
|
||||
nodeAddress = resolved
|
||||
}
|
||||
}
|
||||
if nodeAddress == "" {
|
||||
glog.Warningf("EC Repair plan: unable to resolve node address for shard %d on node %s diskType=%s diskID=%d", loc.ShardID, loc.NodeID, loc.DiskType, loc.DiskID)
|
||||
continue
|
||||
}
|
||||
if result[nodeAddress] == nil {
|
||||
result[nodeAddress] = make(map[uint32]struct{})
|
||||
}
|
||||
result[nodeAddress][loc.ShardID] = struct{}{}
|
||||
}
|
||||
|
||||
final := make(map[string][]uint32, len(result))
|
||||
for nodeAddress, shardSet := range result {
|
||||
shards := make([]uint32, 0, len(shardSet))
|
||||
for shardID := range shardSet {
|
||||
shards = append(shards, shardID)
|
||||
}
|
||||
sort.Slice(shards, func(i, j int) bool { return shards[i] < shards[j] })
|
||||
final[nodeAddress] = shards
|
||||
}
|
||||
return final
|
||||
}
|
||||
|
||||
func selectRebuilder(keep map[uint32]ShardLocation) RebuilderPlan {
|
||||
var rebuilder RebuilderPlan
|
||||
if len(keep) == 0 {
|
||||
return rebuilder
|
||||
}
|
||||
|
||||
nodeCounts := make(map[string]int)
|
||||
diskCounts := make(map[string]map[uint32]int)
|
||||
localShards := make(map[string][]uint32)
|
||||
for shardID, loc := range keep {
|
||||
nodeCounts[loc.NodeAddress]++
|
||||
if diskCounts[loc.NodeAddress] == nil {
|
||||
diskCounts[loc.NodeAddress] = make(map[uint32]int)
|
||||
}
|
||||
diskCounts[loc.NodeAddress][loc.DiskID]++
|
||||
localShards[loc.NodeAddress] = append(localShards[loc.NodeAddress], shardID)
|
||||
}
|
||||
|
||||
var selectedAddr string
|
||||
maxCount := -1
|
||||
for addr, count := range nodeCounts {
|
||||
if count > maxCount || (count == maxCount && addr < selectedAddr) {
|
||||
selectedAddr = addr
|
||||
maxCount = count
|
||||
}
|
||||
}
|
||||
|
||||
var selectedDisk uint32
|
||||
maxDiskCount := -1
|
||||
for diskID, count := range diskCounts[selectedAddr] {
|
||||
if count > maxDiskCount || (count == maxDiskCount && diskID < selectedDisk) {
|
||||
selectedDisk = diskID
|
||||
maxDiskCount = count
|
||||
}
|
||||
}
|
||||
|
||||
shards := localShards[selectedAddr]
|
||||
sort.Slice(shards, func(i, j int) bool { return shards[i] < shards[j] })
|
||||
|
||||
rebuilder.NodeAddress = selectedAddr
|
||||
rebuilder.DiskID = selectedDisk
|
||||
rebuilder.LocalShards = shards
|
||||
return rebuilder
|
||||
}
|
||||
|
||||
func buildCopySources(keep map[uint32]ShardLocation) map[uint32]string {
|
||||
copySources := make(map[uint32]string, len(keep))
|
||||
for shardID, loc := range keep {
|
||||
copySources[shardID] = loc.NodeAddress
|
||||
}
|
||||
return copySources
|
||||
}
|
||||
|
||||
func collectUsedNodes(keep map[uint32]ShardLocation) map[string]bool {
|
||||
used := make(map[string]bool)
|
||||
for _, loc := range keep {
|
||||
if loc.NodeID == "" {
|
||||
continue
|
||||
}
|
||||
used[loc.NodeID] = true
|
||||
}
|
||||
return used
|
||||
}
|
||||
|
||||
func planMissingTargets(activeTopology *topology.ActiveTopology, missing []uint32, usedNodes map[string]bool, diskType string) ([]ShardTarget, error) {
|
||||
if activeTopology == nil {
|
||||
return nil, fmt.Errorf("active topology is nil")
|
||||
}
|
||||
if len(missing) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
selected, err := selectTargetDisks(activeTopology, len(missing), usedNodes, diskType)
|
||||
if err != nil {
|
||||
selected, err = selectTargetDisks(activeTopology, len(missing), nil, diskType)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
missingIDs := make([]uint32, len(missing))
|
||||
copy(missingIDs, missing)
|
||||
sort.Slice(missingIDs, func(i, j int) bool { return missingIDs[i] < missingIDs[j] })
|
||||
|
||||
targetMap := make(map[string]*ShardTarget)
|
||||
for i, shardID := range missingIDs {
|
||||
disk := selected[i]
|
||||
address, err := util.ResolveServerAddress(disk.NodeID, activeTopology)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
key := fmt.Sprintf("%s:%d", disk.NodeID, disk.DiskID)
|
||||
target := targetMap[key]
|
||||
if target == nil {
|
||||
target = &ShardTarget{NodeAddress: address, DiskID: disk.DiskID}
|
||||
targetMap[key] = target
|
||||
}
|
||||
target.ShardIDs = append(target.ShardIDs, shardID)
|
||||
}
|
||||
|
||||
targets := make([]ShardTarget, 0, len(targetMap))
|
||||
for _, target := range targetMap {
|
||||
sort.Slice(target.ShardIDs, func(i, j int) bool { return target.ShardIDs[i] < target.ShardIDs[j] })
|
||||
targets = append(targets, *target)
|
||||
}
|
||||
sort.Slice(targets, func(i, j int) bool {
|
||||
if targets[i].NodeAddress != targets[j].NodeAddress {
|
||||
return targets[i].NodeAddress < targets[j].NodeAddress
|
||||
}
|
||||
return targets[i].DiskID < targets[j].DiskID
|
||||
})
|
||||
|
||||
return targets, nil
|
||||
}
|
||||
|
||||
func selectTargetDisks(activeTopology *topology.ActiveTopology, shardCount int, usedNodes map[string]bool, diskType string) ([]*placement.DiskCandidate, error) {
|
||||
disks := activeTopology.GetDisksWithEffectiveCapacity(topology.TaskTypeErasureCoding, "", 0)
|
||||
if len(disks) == 0 {
|
||||
return nil, fmt.Errorf("no disks available for EC repair placement")
|
||||
}
|
||||
|
||||
candidates := make([]*topology.DiskInfo, 0, len(disks))
|
||||
for _, disk := range disks {
|
||||
if disk == nil || disk.DiskInfo == nil {
|
||||
continue
|
||||
}
|
||||
if diskType != "" && !strings.EqualFold(disk.DiskType, diskType) {
|
||||
continue
|
||||
}
|
||||
if usedNodes != nil && usedNodes[disk.NodeID] {
|
||||
continue
|
||||
}
|
||||
candidates = append(candidates, disk)
|
||||
}
|
||||
if len(candidates) == 0 {
|
||||
return nil, fmt.Errorf("no disks available for EC repair placement after filtering")
|
||||
}
|
||||
|
||||
diskCandidates := diskInfosToCandidates(candidates)
|
||||
if len(diskCandidates) == 0 {
|
||||
return nil, fmt.Errorf("no candidate disks available for EC repair placement")
|
||||
}
|
||||
|
||||
config := placement.PlacementRequest{
|
||||
ShardsNeeded: shardCount,
|
||||
MaxShardsPerServer: 0,
|
||||
MaxShardsPerRack: 0,
|
||||
MaxTaskLoad: topology.MaxTaskLoadForECPlacement,
|
||||
PreferDifferentServers: true,
|
||||
PreferDifferentRacks: true,
|
||||
}
|
||||
|
||||
result, err := placement.SelectDestinations(diskCandidates, config)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(result.SelectedDisks) < shardCount {
|
||||
return nil, fmt.Errorf("found %d destination disks, need %d", len(result.SelectedDisks), shardCount)
|
||||
}
|
||||
return result.SelectedDisks, nil
|
||||
}
|
||||
|
||||
func diskInfosToCandidates(disks []*topology.DiskInfo) []*placement.DiskCandidate {
|
||||
candidates := make([]*placement.DiskCandidate, 0, len(disks))
|
||||
for _, disk := range disks {
|
||||
if disk == nil || disk.DiskInfo == nil {
|
||||
continue
|
||||
}
|
||||
freeSlots := int(disk.DiskInfo.MaxVolumeCount - disk.DiskInfo.VolumeCount)
|
||||
if freeSlots < 0 {
|
||||
freeSlots = 0
|
||||
}
|
||||
|
||||
ecShardCount := 0
|
||||
if disk.DiskInfo.EcShardInfos != nil {
|
||||
for _, shardInfo := range disk.DiskInfo.EcShardInfos {
|
||||
if shardInfo.DiskId == disk.DiskID {
|
||||
ecShardCount += erasure_coding.GetShardCount(shardInfo)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
candidates = append(candidates, &placement.DiskCandidate{
|
||||
NodeID: disk.NodeID,
|
||||
DiskID: disk.DiskID,
|
||||
DataCenter: disk.DataCenter,
|
||||
Rack: disk.Rack,
|
||||
FreeSlots: freeSlots,
|
||||
MaxVolumeCount: disk.DiskInfo.MaxVolumeCount,
|
||||
VolumeCount: disk.DiskInfo.VolumeCount,
|
||||
ShardCount: ecShardCount,
|
||||
LoadCount: disk.LoadCount,
|
||||
})
|
||||
}
|
||||
return candidates
|
||||
}
|
||||
@@ -0,0 +1,304 @@
|
||||
package ec_repair
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/topology"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
const (
|
||||
planTestVolumeID = uint32(999)
|
||||
planTestCollection = "test-collection"
|
||||
recentTaskWindowSize = 10
|
||||
)
|
||||
|
||||
type planTestNode struct {
|
||||
id string
|
||||
address string
|
||||
diskType string
|
||||
diskID uint32
|
||||
shards map[uint32]int64
|
||||
}
|
||||
|
||||
type multiVolumeSpec struct {
|
||||
nodeID string
|
||||
address string
|
||||
diskType string
|
||||
diskID uint32
|
||||
volumeID uint32
|
||||
collection string
|
||||
shards map[uint32]int64
|
||||
}
|
||||
|
||||
func TestDetectReturnsExtraShardCandidates(t *testing.T) {
|
||||
nodes := []planTestNode{
|
||||
{
|
||||
id: "nodeA",
|
||||
address: "127.0.0.1:9100",
|
||||
diskType: "hdd",
|
||||
diskID: 0,
|
||||
shards: makeShardMap(0, 9, 100),
|
||||
},
|
||||
{
|
||||
id: "nodeB",
|
||||
address: "127.0.0.1:9101",
|
||||
diskType: "hdd",
|
||||
diskID: 0,
|
||||
shards: map[uint32]int64{
|
||||
0: 90, // mismatched size
|
||||
2: 100,
|
||||
3: 100,
|
||||
4: 100,
|
||||
5: 100,
|
||||
6: 100,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
topo := buildPlanTopology(nodes, planTestVolumeID, planTestCollection)
|
||||
candidates, hasMore, err := Detect(topo, "", 0)
|
||||
require.NoError(t, err)
|
||||
require.False(t, hasMore)
|
||||
require.Len(t, candidates, 1)
|
||||
|
||||
candidate := candidates[0]
|
||||
require.Greater(t, candidate.ExtraShards, 0)
|
||||
require.Greater(t, candidate.MismatchedShards, 0)
|
||||
}
|
||||
|
||||
func TestDetectHonorsMaxResults(t *testing.T) {
|
||||
specs := []multiVolumeSpec{
|
||||
{
|
||||
nodeID: "volume100-node1",
|
||||
address: "127.0.0.1:9200",
|
||||
diskType: "hdd",
|
||||
diskID: 0,
|
||||
volumeID: 100,
|
||||
collection: planTestCollection,
|
||||
shards: makeShardMap(0, 9, 100),
|
||||
},
|
||||
{
|
||||
nodeID: "volume100-node2",
|
||||
address: "127.0.0.1:9201",
|
||||
diskType: "hdd",
|
||||
diskID: 1,
|
||||
volumeID: 100,
|
||||
collection: planTestCollection,
|
||||
shards: map[uint32]int64{
|
||||
0: 50,
|
||||
},
|
||||
},
|
||||
{
|
||||
nodeID: "volume101-node1",
|
||||
address: "127.0.0.1:9300",
|
||||
diskType: "hdd",
|
||||
diskID: 0,
|
||||
volumeID: 101,
|
||||
collection: planTestCollection,
|
||||
shards: makeShardMap(0, 9, 100),
|
||||
},
|
||||
{
|
||||
nodeID: "volume101-node2",
|
||||
address: "127.0.0.1:9301",
|
||||
diskType: "hdd",
|
||||
diskID: 1,
|
||||
volumeID: 101,
|
||||
collection: planTestCollection,
|
||||
shards: map[uint32]int64{
|
||||
1: 90,
|
||||
},
|
||||
},
|
||||
{
|
||||
nodeID: "volume102-node1",
|
||||
address: "127.0.0.1:9400",
|
||||
diskType: "hdd",
|
||||
diskID: 0,
|
||||
volumeID: 102,
|
||||
collection: planTestCollection,
|
||||
shards: makeShardMap(0, 9, 100),
|
||||
},
|
||||
{
|
||||
nodeID: "volume102-node2",
|
||||
address: "127.0.0.1:9401",
|
||||
diskType: "hdd",
|
||||
diskID: 1,
|
||||
volumeID: 102,
|
||||
collection: planTestCollection,
|
||||
shards: map[uint32]int64{
|
||||
2: 50,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
topo := buildMultiVolumeTopology(specs)
|
||||
|
||||
candidates, hasMore, err := Detect(topo, "", 2)
|
||||
require.NoError(t, err)
|
||||
require.True(t, hasMore)
|
||||
require.Len(t, candidates, 2)
|
||||
}
|
||||
|
||||
func TestBuildRepairPlanRequiresEnoughShards(t *testing.T) {
|
||||
nodes := []planTestNode{
|
||||
{id: "nodeA", address: "n1", diskType: "hdd", diskID: 0, shards: makeShardMap(0, 4, 100)},
|
||||
}
|
||||
topo := buildPlanTopology(nodes, planTestVolumeID, planTestCollection)
|
||||
|
||||
_, err := BuildRepairPlan(topo, nil, planTestVolumeID, planTestCollection, "hdd")
|
||||
require.Error(t, err)
|
||||
require.Contains(t, err.Error(), fmt.Sprintf("need at least %d", erasure_coding.DataShardsCount))
|
||||
}
|
||||
|
||||
func TestBuildRepairPlanIncludesTargetsAndDeletes(t *testing.T) {
|
||||
nodes := []planTestNode{
|
||||
{id: "nodeA", address: "n1", diskType: "hdd", diskID: 0, shards: makeShardMap(0, 9, 100)},
|
||||
{id: "nodeB", address: "n2", diskType: "hdd", diskID: 0, shards: map[uint32]int64{0: 50}},
|
||||
}
|
||||
topo := buildPlanTopology(nodes, planTestVolumeID, planTestCollection)
|
||||
|
||||
activeTopo := buildActiveTopology(t, []string{"n3", "n4", "n5", "n6"})
|
||||
plan, err := BuildRepairPlan(topo, activeTopo, planTestVolumeID, planTestCollection, "hdd")
|
||||
require.NoError(t, err)
|
||||
require.NotEmpty(t, plan.MissingShards)
|
||||
require.Contains(t, plan.MissingShards, uint32(10))
|
||||
require.NotEmpty(t, plan.Targets)
|
||||
require.NotEmpty(t, plan.DeleteByNode)
|
||||
}
|
||||
|
||||
func makeShardMap(from, to int, size int64) map[uint32]int64 {
|
||||
shards := make(map[uint32]int64)
|
||||
for i := from; i <= to; i++ {
|
||||
shards[uint32(i)] = size
|
||||
}
|
||||
return shards
|
||||
}
|
||||
|
||||
func buildPlanTopology(nodes []planTestNode, volumeID uint32, collection string) *master_pb.TopologyInfo {
|
||||
dataNodes := make([]*master_pb.DataNodeInfo, 0, len(nodes))
|
||||
for _, node := range nodes {
|
||||
diskInfo := &master_pb.DiskInfo{
|
||||
DiskId: node.diskID,
|
||||
MaxVolumeCount: 100,
|
||||
VolumeCount: 10,
|
||||
}
|
||||
if len(node.shards) > 0 {
|
||||
diskInfo.EcShardInfos = []*master_pb.VolumeEcShardInformationMessage{
|
||||
buildEcShardInfo(volumeID, collection, node.diskType, node.diskID, node.shards),
|
||||
}
|
||||
}
|
||||
dataNodes = append(dataNodes, &master_pb.DataNodeInfo{
|
||||
Id: node.id,
|
||||
Address: node.address,
|
||||
DiskInfos: map[string]*master_pb.DiskInfo{node.diskType: diskInfo},
|
||||
})
|
||||
}
|
||||
return &master_pb.TopologyInfo{
|
||||
DataCenterInfos: []*master_pb.DataCenterInfo{
|
||||
{
|
||||
Id: "dc1",
|
||||
RackInfos: []*master_pb.RackInfo{
|
||||
{
|
||||
Id: "rack1",
|
||||
DataNodeInfos: dataNodes,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func buildMultiVolumeTopology(specs []multiVolumeSpec) *master_pb.TopologyInfo {
|
||||
dataNodes := make([]*master_pb.DataNodeInfo, 0, len(specs))
|
||||
for _, spec := range specs {
|
||||
diskInfo := &master_pb.DiskInfo{
|
||||
DiskId: spec.diskID,
|
||||
MaxVolumeCount: 100,
|
||||
VolumeCount: 10,
|
||||
}
|
||||
if len(spec.shards) > 0 {
|
||||
diskInfo.EcShardInfos = []*master_pb.VolumeEcShardInformationMessage{
|
||||
buildEcShardInfo(spec.volumeID, spec.collection, spec.diskType, spec.diskID, spec.shards),
|
||||
}
|
||||
}
|
||||
dataNodes = append(dataNodes, &master_pb.DataNodeInfo{
|
||||
Id: spec.nodeID,
|
||||
Address: spec.address,
|
||||
DiskInfos: map[string]*master_pb.DiskInfo{spec.diskType: diskInfo},
|
||||
})
|
||||
}
|
||||
return &master_pb.TopologyInfo{
|
||||
DataCenterInfos: []*master_pb.DataCenterInfo{
|
||||
{
|
||||
Id: "dc1",
|
||||
RackInfos: []*master_pb.RackInfo{
|
||||
{
|
||||
Id: "rack1",
|
||||
DataNodeInfos: dataNodes,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func buildActiveTopology(t *testing.T, nodeIDs []string) *topology.ActiveTopology {
|
||||
t.Helper()
|
||||
info := &master_pb.TopologyInfo{
|
||||
DataCenterInfos: []*master_pb.DataCenterInfo{
|
||||
{
|
||||
Id: "dc1",
|
||||
RackInfos: []*master_pb.RackInfo{
|
||||
{
|
||||
Id: "rack1",
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
rack := info.DataCenterInfos[0].RackInfos[0]
|
||||
for _, id := range nodeIDs {
|
||||
rack.DataNodeInfos = append(rack.DataNodeInfos, &master_pb.DataNodeInfo{
|
||||
Id: id,
|
||||
Address: id,
|
||||
DiskInfos: map[string]*master_pb.DiskInfo{
|
||||
"hdd": {
|
||||
DiskId: 0,
|
||||
MaxVolumeCount: 200,
|
||||
VolumeCount: 50,
|
||||
},
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
active := topology.NewActiveTopology(recentTaskWindowSize)
|
||||
require.NoError(t, active.UpdateTopology(info))
|
||||
return active
|
||||
}
|
||||
|
||||
func buildEcShardInfo(volumeID uint32, collection, diskType string, diskID uint32, shardSizes map[uint32]int64) *master_pb.VolumeEcShardInformationMessage {
|
||||
shardIDs := make([]int, 0, len(shardSizes))
|
||||
for shardID := range shardSizes {
|
||||
shardIDs = append(shardIDs, int(shardID))
|
||||
}
|
||||
sort.Ints(shardIDs)
|
||||
|
||||
var bits uint32
|
||||
sizes := make([]int64, 0, len(shardIDs))
|
||||
for _, shardID := range shardIDs {
|
||||
bits |= (1 << shardID)
|
||||
sizes = append(sizes, shardSizes[uint32(shardID)])
|
||||
}
|
||||
return &master_pb.VolumeEcShardInformationMessage{
|
||||
Id: volumeID,
|
||||
Collection: collection,
|
||||
EcIndexBits: bits,
|
||||
DiskType: diskType,
|
||||
DiskId: diskID,
|
||||
ShardSizes: sizes,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package ec_repair
|
||||
|
||||
// VolumeKey identifies a set of EC shards for a volume within a collection and disk type.
|
||||
type VolumeKey struct {
|
||||
VolumeID uint32
|
||||
Collection string
|
||||
DiskType string
|
||||
}
|
||||
|
||||
// ShardLocation describes one shard copy for a volume.
|
||||
type ShardLocation struct {
|
||||
NodeID string
|
||||
NodeAddress string
|
||||
DataCenter string
|
||||
Rack string
|
||||
DiskType string
|
||||
DiskID uint32
|
||||
ShardID uint32
|
||||
Size int64
|
||||
}
|
||||
|
||||
// RepairCandidate summarizes EC shard issues for detection.
|
||||
type RepairCandidate struct {
|
||||
VolumeID uint32
|
||||
Collection string
|
||||
DiskType string
|
||||
MissingShards int
|
||||
ExtraShards int
|
||||
MismatchedShards int
|
||||
}
|
||||
|
||||
// ShardTarget describes where to place rebuilt shards.
|
||||
type ShardTarget struct {
|
||||
NodeAddress string
|
||||
DiskID uint32
|
||||
ShardIDs []uint32
|
||||
}
|
||||
|
||||
// RebuilderPlan describes the node that performs shard rebuilds.
|
||||
type RebuilderPlan struct {
|
||||
NodeAddress string
|
||||
DiskID uint32
|
||||
LocalShards []uint32
|
||||
}
|
||||
|
||||
// RepairPlan describes all actions needed to repair an EC volume.
|
||||
type RepairPlan struct {
|
||||
VolumeID uint32
|
||||
Collection string
|
||||
DiskType string
|
||||
MissingShards []uint32
|
||||
Targets []ShardTarget
|
||||
DeleteByNode map[string][]uint32
|
||||
Rebuilder RebuilderPlan
|
||||
CopySources map[uint32]string // shardID -> source node address
|
||||
}
|
||||
@@ -14,6 +14,7 @@ const (
|
||||
TaskTypeVacuum TaskType = "vacuum"
|
||||
TaskTypeErasureCoding TaskType = "erasure_coding"
|
||||
TaskTypeBalance TaskType = "balance"
|
||||
TaskTypeEcRepair TaskType = "ec_repair"
|
||||
TaskTypeReplication TaskType = "replication"
|
||||
)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user