mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-20 13:30:46 +02:00
* master: let VolumeList ask for the volumes it wants The request carried nothing, so every caller was answered with the whole cluster. A dashboard opening one volume's page, or a capacity probe adding up one bucket, was served all 800k of them and threw away the rest -- and the master built every one of those messages first. The topology, its disks and their counters are still reported in full: a caller reading free space or replica placement needs the cluster whichever volumes it asked about. Only what is listed under a disk is selected, ec shards included. An empty collection and a zero volume id take everything, the way volume.list already reads its own -collectionPattern and -volumeId, so a caller that forgets to narrow is answered too much rather than answered wrongly. That leaves the default collection unnameable, since it is the one the empty string names, so it gets a field of its own. An older client sends none of it and is answered exactly as before. * admin: ask the master for the volume the page is showing A volume's detail page was pulling every volume in the cluster to find one and its replicas, and discarding the rest. * admin: ask the master for the ec volume the page is showing Same as the volume detail page: one volume's shards were found by pulling every ec shard in the cluster. * s3: ask the master for the bucket's own collection The SOSAPI capacity probe summed one collection's volumes out of a listing of every volume in the cluster. Cluster capacity still comes out the same: it is read from the disk counters, which a filtered listing reports in full. * topology: read the disk usage counters atomically They are written with atomic.AddInt64 from heartbeats but were read plainly by the two listings and by FreeSpace, and the map they sit in was iterated without the lock its neighbour takes. Under -race a listing concurrent with a heartbeat trips on both.
454 lines
15 KiB
Go
454 lines
15 KiB
Go
package weed_server
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"math/rand/v2"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/stats"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/topology"
|
|
|
|
"github.com/seaweedfs/raft"
|
|
|
|
"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/security"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/status"
|
|
)
|
|
|
|
const (
|
|
volumeGrowStepCount = 2
|
|
)
|
|
|
|
func (ms *MasterServer) DoAutomaticVolumeGrow(req *topology.VolumeGrowRequest) {
|
|
if ms.option.VolumeGrowthDisabled {
|
|
glog.V(1).Infof("automatic volume grow disabled")
|
|
return
|
|
}
|
|
glog.V(1).Infoln("starting automatic volume grow")
|
|
start := time.Now()
|
|
newVidLocations, err := ms.vg.AutomaticGrowByType(req.Option, ms.grpcDialOption, ms.Topo, req.Count)
|
|
glog.V(1).Infoln("finished automatic volume grow, cost ", time.Now().Sub(start))
|
|
if err != nil {
|
|
glog.V(1).Infof("automatic volume grow failed: %+v", err)
|
|
return
|
|
}
|
|
ms.broadcastVolumeLocationsToClients(newVidLocations)
|
|
}
|
|
|
|
func (ms *MasterServer) ProcessGrowRequest() {
|
|
go func() {
|
|
ctx := context.Background()
|
|
firstRun := true
|
|
for {
|
|
if firstRun {
|
|
firstRun = false
|
|
} else {
|
|
time.Sleep(5*time.Minute + time.Duration(30*rand.Float32())*time.Second)
|
|
}
|
|
if !ms.Topo.IsLeader() {
|
|
continue
|
|
}
|
|
dcs := ms.Topo.ListDCAndRacks()
|
|
for _, vlc := range ms.Topo.ListVolumeLayoutCollections() {
|
|
var err error
|
|
vl := vlc.VolumeLayout
|
|
lastGrowCount := vl.GetLastGrowCount()
|
|
if vl.HasGrowRequest() {
|
|
continue
|
|
}
|
|
writable, crowded := vl.GetWritableVolumeCount()
|
|
mustGrow := int(lastGrowCount) - writable
|
|
vgr := vlc.ToVolumeGrowRequest()
|
|
underReplicated := vl.CountUnderReplicatedVolumes()
|
|
stats.MasterVolumeLayoutWritable.WithLabelValues(vlc.Collection, vgr.DiskType, vgr.Replication, vgr.Ttl).Set(float64(writable))
|
|
stats.MasterVolumeLayoutCrowded.WithLabelValues(vlc.Collection, vgr.DiskType, vgr.Replication, vgr.Ttl).Set(float64(crowded))
|
|
stats.MasterUnderReplicatedVolumes.WithLabelValues(vlc.Collection, vgr.DiskType, vgr.Replication, vgr.Ttl).Set(float64(underReplicated))
|
|
|
|
switch {
|
|
case mustGrow > 0:
|
|
if rp, rpErr := super_block.NewReplicaPlacementFromString(vgr.Replication); rpErr != nil {
|
|
glog.V(0).Infof("failed to parse replica placement %s: %v", vgr.Replication, rpErr)
|
|
} else {
|
|
vgr.WritableVolumeCount = uint32(mustGrow)
|
|
if ms.Topo.AvailableSpaceFor(&topology.VolumeGrowOption{DiskType: types.ToDiskType(vgr.DiskType)}) >= int64(vgr.WritableVolumeCount*uint32(rp.GetCopyCount())) {
|
|
_, err = ms.VolumeGrow(ctx, vgr)
|
|
}
|
|
}
|
|
case lastGrowCount > 0 && writable < int(lastGrowCount*2) && float64(crowded+volumeGrowStepCount) > float64(writable)*topology.VolumeGrowStrategy.Threshold:
|
|
vgr.WritableVolumeCount = volumeGrowStepCount
|
|
_, err = ms.VolumeGrow(ctx, vgr)
|
|
}
|
|
if err != nil {
|
|
glog.V(0).Infof("volume grow request failed: %+v", err)
|
|
}
|
|
for _, plan := range vl.PlanRackAwareGrowth(dcs, lastGrowCount, volumeGrowStepCount) {
|
|
vgr.DataCenter = plan.DataCenter
|
|
vgr.Rack = plan.Rack
|
|
vgr.WritableVolumeCount = plan.WritableVolumeCount
|
|
if _, err = ms.VolumeGrow(ctx, vgr); err != nil {
|
|
glog.V(0).Infof("volume grow request for dc:%s rack:%s failed: %+v", plan.DataCenter, plan.Rack, err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
go func() {
|
|
filter := sync.Map{}
|
|
for {
|
|
req, ok := <-ms.volumeGrowthRequestChan
|
|
if !ok {
|
|
break
|
|
}
|
|
|
|
option := req.Option
|
|
vl := ms.Topo.GetVolumeLayout(option.Collection, option.ReplicaPlacement, option.Ttl, option.DiskType)
|
|
|
|
if !ms.Topo.IsLeader() {
|
|
//discard buffered requests
|
|
time.Sleep(time.Second * 1)
|
|
vl.DoneGrowRequest()
|
|
continue
|
|
}
|
|
|
|
// filter out identical requests being processed
|
|
found := false
|
|
filter.Range(func(k, v interface{}) bool {
|
|
existingReq := k.(*topology.VolumeGrowRequest)
|
|
if existingReq.Equals(req) {
|
|
found = true
|
|
}
|
|
return !found
|
|
})
|
|
|
|
// not atomic but it's okay
|
|
if found || (!req.Force && !vl.ShouldGrowVolumes()) {
|
|
glog.V(4).Infoln("discard volume grow request")
|
|
time.Sleep(time.Millisecond * 211)
|
|
vl.DoneGrowRequest()
|
|
continue
|
|
}
|
|
|
|
filter.Store(req, nil)
|
|
// we have lock called inside vg
|
|
glog.V(0).Infof("volume grow %+v", req)
|
|
go func(req *topology.VolumeGrowRequest, vl *topology.VolumeLayout) {
|
|
// defer so a panic can't strand growRequest.
|
|
defer filter.Delete(req)
|
|
defer vl.DoneGrowRequest()
|
|
ms.DoAutomaticVolumeGrow(req)
|
|
}(req, vl)
|
|
}
|
|
}()
|
|
}
|
|
|
|
func (ms *MasterServer) LookupVolume(ctx context.Context, req *master_pb.LookupVolumeRequest) (*master_pb.LookupVolumeResponse, error) {
|
|
|
|
resp := &master_pb.LookupVolumeResponse{}
|
|
volumeLocations := ms.lookupVolumeId(req.VolumeOrFileIds, req.Collection)
|
|
|
|
notFoundCount := 0
|
|
for _, volumeOrFileId := range req.VolumeOrFileIds {
|
|
vid := volumeOrFileId
|
|
commaSep := strings.Index(vid, ",")
|
|
if commaSep > 0 {
|
|
vid = vid[0:commaSep]
|
|
}
|
|
if result, found := volumeLocations[vid]; found {
|
|
var locations []*master_pb.Location
|
|
for _, loc := range result.Locations {
|
|
locations = append(locations, &master_pb.Location{
|
|
Url: loc.Url,
|
|
PublicUrl: loc.PublicUrl,
|
|
DataCenter: loc.DataCenter,
|
|
GrpcPort: uint32(loc.GrpcPort),
|
|
})
|
|
}
|
|
var auth string
|
|
if commaSep > 0 { // this is a file id
|
|
auth = string(security.GenJwtForVolumeServer(ms.guard.SigningKey(), ms.guard.ExpiresAfterSec(), result.VolumeOrFileId))
|
|
}
|
|
if result.NotFound {
|
|
notFoundCount++
|
|
}
|
|
resp.VolumeIdLocations = append(resp.VolumeIdLocations, &master_pb.LookupVolumeResponse_VolumeIdLocation{
|
|
VolumeOrFileId: result.VolumeOrFileId,
|
|
Locations: locations,
|
|
Error: result.Error,
|
|
Auth: auth,
|
|
})
|
|
}
|
|
}
|
|
|
|
// Only return Unavailable during warmup when every requested ID was a transient not-found
|
|
if len(req.VolumeOrFileIds) > 0 && notFoundCount == len(req.VolumeOrFileIds) && ms.Topo.IsLeader() && ms.Topo.IsWarmingUp() {
|
|
glog.V(0).Infof("lookup volume warming up: topology is still loading (%d not found)", notFoundCount)
|
|
return nil, status.Errorf(codes.Unavailable, "master is warming up, topology is still loading")
|
|
}
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
// CollectionStatistics summarises every collection, so callers tracking usage
|
|
// do not pull the whole volume list to add it up themselves.
|
|
func (ms *MasterServer) CollectionStatistics(ctx context.Context, req *master_pb.CollectionStatisticsRequest) (*master_pb.CollectionStatisticsResponse, error) {
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
|
|
stats := ms.Topo.CollectionStatistics()
|
|
resp := &master_pb.CollectionStatisticsResponse{
|
|
Collections: make([]*master_pb.CollectionStatistics, 0, len(stats)),
|
|
}
|
|
for _, s := range stats {
|
|
resp.Collections = append(resp.Collections, &master_pb.CollectionStatistics{
|
|
Collection: s.Collection,
|
|
FileCount: s.FileCount,
|
|
DeleteCount: s.DeleteCount,
|
|
DeletedByteCount: s.DeletedByteCount,
|
|
Size: s.Size,
|
|
PhysicalSize: s.PhysicalSize,
|
|
VolumeCount: s.VolumeCount,
|
|
})
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) Statistics(ctx context.Context, req *master_pb.StatisticsRequest) (*master_pb.StatisticsResponse, error) {
|
|
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
|
|
// an empty collection means all collections, and a named collection covers
|
|
// all its layouts and EC volumes
|
|
stats := ms.Topo.CollectionVolumeStats(req.Collection)
|
|
totalSize := uint64(ms.Topo.GetDiskUsages().GetMaxVolumeCount() * int64(ms.option.VolumeSizeLimitMB) * 1024 * 1024)
|
|
// capacity is cluster-wide, so what is left over is what every collection
|
|
// has not taken, not just the one asked about
|
|
clusterUsedSize := stats.UsedSize
|
|
if req.Collection != "" {
|
|
clusterUsedSize = ms.Topo.CollectionVolumeStats("").UsedSize
|
|
}
|
|
// and the free space holds that many copies fewer of whatever the caller writes
|
|
var freeSize uint64
|
|
if totalSize > clusterUsedSize {
|
|
freeSize = (totalSize - clusterUsedSize) / uint64(ms.replicaCopyCount(req.Replication))
|
|
}
|
|
resp := &master_pb.StatisticsResponse{
|
|
TotalSize: totalSize,
|
|
UsedSize: stats.UsedSize,
|
|
FileCount: stats.FileCount,
|
|
LogicalTotalSize: stats.LogicalUsedSize + freeSize,
|
|
LogicalUsedSize: stats.LogicalUsedSize,
|
|
}
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
// replicaCopyCount returns how many copies the given replication makes, falling
|
|
// back to the master default and then to a single copy. Unparsable input is
|
|
// reported rather than failing the call: statistics are informational.
|
|
func (ms *MasterServer) replicaCopyCount(replication string) int {
|
|
for _, s := range []string{replication, ms.option.DefaultReplicaPlacement} {
|
|
if s == "" {
|
|
continue
|
|
}
|
|
rp, err := super_block.NewReplicaPlacementFromString(s)
|
|
if err != nil {
|
|
glog.V(1).Infof("statistics replication %q: %v", s, err)
|
|
continue
|
|
}
|
|
return rp.GetCopyCount()
|
|
}
|
|
return 1
|
|
}
|
|
|
|
func (ms *MasterServer) VolumeList(ctx context.Context, req *master_pb.VolumeListRequest) (*master_pb.VolumeListResponse, error) {
|
|
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
|
|
resp := &master_pb.VolumeListResponse{
|
|
TopologyInfo: ms.Topo.ToTopologyInfo(topology.NewVolumeFilter(req)),
|
|
VolumeSizeLimitMb: uint64(ms.option.VolumeSizeLimitMB),
|
|
}
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) LookupEcVolume(ctx context.Context, req *master_pb.LookupEcVolumeRequest) (*master_pb.LookupEcVolumeResponse, error) {
|
|
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
|
|
resp := &master_pb.LookupEcVolumeResponse{}
|
|
|
|
ecLocations, found := ms.Topo.LookupEcShards(needle.VolumeId(req.VolumeId))
|
|
|
|
if !found {
|
|
return resp, fmt.Errorf("ec volume %d not found", req.VolumeId)
|
|
}
|
|
|
|
resp.VolumeId = req.VolumeId
|
|
|
|
for shardId, shardLocations := range ecLocations.Locations {
|
|
var locations []*master_pb.Location
|
|
for _, dn := range shardLocations {
|
|
locations = append(locations, &master_pb.Location{
|
|
Url: dn.Url(),
|
|
PublicUrl: dn.PublicUrl,
|
|
DataCenter: dn.GetDataCenterId(),
|
|
// without this, clients derive grpc as httpPort+10000
|
|
GrpcPort: uint32(dn.GrpcPort),
|
|
})
|
|
}
|
|
resp.ShardIdLocations = append(resp.ShardIdLocations, &master_pb.LookupEcVolumeResponse_EcShardIdLocation{
|
|
ShardId: uint32(shardId),
|
|
Locations: locations,
|
|
})
|
|
}
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) VacuumVolume(ctx context.Context, req *master_pb.VacuumVolumeRequest) (*master_pb.VacuumVolumeResponse, error) {
|
|
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
|
|
resp := &master_pb.VacuumVolumeResponse{}
|
|
|
|
ms.Topo.Vacuum(ms.grpcDialOption, float64(req.GarbageThreshold), ms.option.MaxParallelVacuumPerServer, req.VolumeId, req.Collection, ms.preallocateSize, false)
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) DisableVacuum(ctx context.Context, req *master_pb.DisableVacuumRequest) (*master_pb.DisableVacuumResponse, error) {
|
|
// The caller explicitly indicates whether this disable request comes
|
|
// from the vacuum plugin monitor. Track ownership so the safety net
|
|
// in the vacuum loop won't override an operator's intentional disable.
|
|
if req.GetByPlugin() {
|
|
ms.Topo.DisableVacuumByPlugin()
|
|
} else {
|
|
ms.Topo.DisableVacuum()
|
|
}
|
|
resp := &master_pb.DisableVacuumResponse{}
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) EnableVacuum(ctx context.Context, req *master_pb.EnableVacuumRequest) (*master_pb.EnableVacuumResponse, error) {
|
|
if req.GetByPlugin() {
|
|
ms.Topo.EnableVacuumByPlugin()
|
|
} else {
|
|
ms.Topo.EnableVacuum()
|
|
}
|
|
resp := &master_pb.EnableVacuumResponse{}
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) VolumeMarkReadonly(ctx context.Context, req *master_pb.VolumeMarkReadonlyRequest) (*master_pb.VolumeMarkReadonlyResponse, error) {
|
|
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
|
|
resp := &master_pb.VolumeMarkReadonlyResponse{}
|
|
|
|
replicaPlacement, _ := super_block.NewReplicaPlacementFromByte(byte(req.ReplicaPlacement))
|
|
vl := ms.Topo.GetVolumeLayout(req.Collection, replicaPlacement, needle.LoadTTLFromUint32(req.Ttl), types.ToDiskType(req.DiskType))
|
|
vid := needle.VolumeId(req.VolumeId)
|
|
dataNodes := ms.Topo.Lookup(req.Collection, vid)
|
|
|
|
found := false
|
|
for _, dn := range dataNodes {
|
|
if dn.Ip == req.Ip && dn.Port == int(req.Port) {
|
|
found = true
|
|
if req.IsReadonly {
|
|
vl.SetVolumeReadOnly(dn, vid)
|
|
if pending := vl.GetPendingSize(vid); pending > 0 {
|
|
glog.V(0).Infof("volume %d marked readonly with %d pending bytes", vid, pending)
|
|
}
|
|
} else {
|
|
vl.SetVolumeWritable(dn, vid)
|
|
}
|
|
}
|
|
}
|
|
|
|
// A vacuum worker marks the volume writable after a successful commit. If the
|
|
// index lost it meanwhile (e.g. a disconnect race during the vacuum), the loop
|
|
// found nothing to update; re-register it from the node that still holds it so
|
|
// LookupVolume stops returning "not found".
|
|
if !found && !req.IsReadonly {
|
|
if dn := ms.Topo.LookupDataNodeByAddress(pb.NewServerAddress(req.Ip, int(req.Port), 0)); dn != nil {
|
|
if vi, err := dn.GetVolumesById(vid); err == nil {
|
|
ms.Topo.RegisterVolumeLayout(vi, dn)
|
|
glog.V(0).Infof("volume %d re-registered from %s:%d, was missing from the lookup index at mark-writable", vid, req.Ip, req.Port)
|
|
}
|
|
}
|
|
}
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
func (ms *MasterServer) VolumeGrow(ctx context.Context, req *master_pb.VolumeGrowRequest) (*master_pb.VolumeGrowResponse, error) {
|
|
if !ms.Topo.IsLeader() {
|
|
return nil, raft.NotLeaderError
|
|
}
|
|
if req.Replication == "" {
|
|
req.Replication = ms.option.DefaultReplicaPlacement
|
|
}
|
|
replicaPlacement, err := super_block.NewReplicaPlacementFromString(req.Replication)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
ttl, err := needle.ReadTTL(req.Ttl)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if req.DataCenter != "" && !ms.Topo.DataCenterExists(req.DataCenter) {
|
|
return nil, fmt.Errorf("data center not exists")
|
|
}
|
|
|
|
ver := needle.GetCurrentVersion()
|
|
volumeGrowOption := topology.VolumeGrowOption{
|
|
Collection: req.Collection,
|
|
ReplicaPlacement: replicaPlacement,
|
|
Ttl: ttl,
|
|
DiskType: types.ToDiskType(req.DiskType),
|
|
Preallocate: ms.preallocateSize,
|
|
DataCenter: req.DataCenter,
|
|
Rack: req.Rack,
|
|
DataNode: req.DataNode,
|
|
MemoryMapMaxSizeMb: req.MemoryMapMaxSizeMb,
|
|
Version: uint32(ver),
|
|
}
|
|
volumeGrowRequest := topology.VolumeGrowRequest{
|
|
Option: &volumeGrowOption,
|
|
Count: req.WritableVolumeCount,
|
|
Force: true,
|
|
Reason: "grpc volume grow",
|
|
}
|
|
replicaCount := int64(req.WritableVolumeCount * uint32(replicaPlacement.GetCopyCount()))
|
|
|
|
if ms.Topo.AvailableSpaceFor(&volumeGrowOption) < replicaCount {
|
|
return nil, fmt.Errorf("only %d volumes left, not enough for %d", ms.Topo.AvailableSpaceFor(&volumeGrowOption), replicaCount)
|
|
}
|
|
|
|
ms.DoAutomaticVolumeGrow(&volumeGrowRequest)
|
|
|
|
return &master_pb.VolumeGrowResponse{}, nil
|
|
}
|