mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-29 19:25:35 +00:00
* cluster: never broadcast an empty lock ring An empty member list is never a usable ring state, but a delayed RemoveServer on a former leader can fire after the new leader already broadcast the recovered ring. That late broadcast carries a newer wall-clock version, so clients accept the empty ring and permanently reject the good one. Skip the broadcast entirely when the member list is empty, keeping the last non-empty snapshot for reconnecting clients. * cluster: periodically rebroadcast the lock ring Ring updates are purely event-driven, so one lost or poisoned update is permanent until the next membership change — with a single filer that may never come. Re-arm a per-group timer after every broadcast so the current leader keeps re-sending the ring; clients reject nothing newer than their last accepted version, so a re-sent snapshot always heals a stale view. * filer,s3api: reset the lock ring on master change Ring versions are per-master monotonic — each master stamps wall-clock nanoseconds — so a late high-version update accepted from a former leader makes the new leader's snapshot look stale forever. Detect a leader change across the reconnect gap (currentMaster is cleared between attempts, so remember the last served master) and reset the ring to bootstrap state so the new leader's view always applies. * cluster: fail lock acquisition when no lock server exists retryUntilLocked loops forever, so a filer reporting an empty lock ring wedges every append write indefinitely. Bound only the "no lock server found" case — ordinary contention is still waited out since the holder releases eventually. The constructors now return nil on failure: the filer append path and S3 object writes fail fast, while mounts degrade to their existing lockless mode. * cluster: reset only the ring version on master change Ring versions are per-master monotonic, so a version gate reset is all a leader change needs. Clearing the whole ring made every filer its own write owner until the next update and dropped the prior-owner window for keys the new leader remaps; the last ring now keeps routing until the new leader's snapshot transitions off it. * cluster: skip redundant ring installs and defer rebroadcasts An unchanged member list now only bumps the accepted version instead of installing a snapshot: periodic rebroadcasts no longer fire the topology-change callback or restart the prior-owner window. And a rebroadcast that lands inside a membership stabilization window yields to the pending timer rather than publishing an intermediate ring. * cluster,mount: bound lock unavailability, fail ops that cannot lock Only 'lock already owned' contention retries without bound now; every other failure — no lock server, or a dead ring member refusing connections — shares the same unavailability budget, so a ring naming departed filers can no longer hang a lock forever. Mount open-write, create, and rename fail with EAGAIN when the required lock cannot be acquired instead of proceeding without cross-mount serialization. * cluster: check pending stabilization inside the broadcast critical section rebroadcast released the mutex between the pending-timer check and nextBroadcastUpdate, so a membership change arriving in the gap could arm a stabilization timer while the rebroadcast emitted an intermediate ring. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: acquire path locks before mutating create/rename state Create took the DLM lock only after the filer create, so a lock failure returned EAGAIN with an eagerly persisted file left behind. Rename marked source handles renamed before acquiring locks, so a failed acquisition left them suppressing old-path flushes for a rename that never happened. Both now take the locks first; the create's lock is released again if the entry race loses to another creator and AcquireHandle takes over. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: keep the old-path lock when rename lock migration fails The migration stopped the handle's lock before acquiring the replacement, so a nil result left the handle writing with no lock at all. Acquiring the new-path lock first means failure keeps the existing lock instead of reporting success with serialization dropped. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: skip new-path rename lock when a handle already holds it A target file open for write on this mount already carries a lock on newPath; the lock manager does not grant a second lock to the same owner, so the rename would wait on itself until the handle closed. Also avoid locking twice when old and new paths coincide. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: hand the rename's target lock to the migrating handle The rename holds a lock on newPath for its duration, so the response migration's fresh acquisition waited on that same lock until the handle released — under fhLockTable, blocking the handle's own close. Adopt the rename's lock directly; nested move responses still acquire their own. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: move the replaced target's lock to the renamed handle When the target path was already locked by an open handle on this mount, the migrated source handle kept only its stale old-path lock — the target's close would then release the last lock on the new path while the renamed handle was still open. Adopt the replaced handle's lock instead. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: stop the handle lock inside the fh lock on release ReleaseHandle stopped fh.dlmLock before taking the fhLockTable slot, so a rename migration holding that slot could still observe and adopt a lock that was already stopping. Stopping under the fh lock makes the transfer serialize against the release. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: claim the replaced target's lock for the renamed handle When the target path is already locked by an open handle on this mount, adopting it at migration time keeps the renamed path protected after that handle closes, without waiting on a lock this mount already holds. If the handle was released mid-migration the claimed lock is stopped, and a fresh acquire covers the case where it was already gone. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: claim the target handle's lock before the rename runs Skipping the new-path lock when a handle already holds it let that handle's close release the lock mid-rename, leaving the path unguarded until the response migrated it. Take over the lock at check time and hold it for the rename's duration: the response adopts it for the migrating handle, or it returns to the target handle / is released on failure. The target handle lookup also falls back to the entry's stored inode for a forgotten path mapping. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> * mount: read handle locks only under the fh lock during rename The loose dlmLock reads raced ReleaseHandle, which now mutates the lock inside the handle lock; check and claim it under the same hold. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
653 lines
24 KiB
Go
653 lines
24 KiB
Go
package wdclient
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"math/rand"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"google.golang.org/grpc"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/status"
|
|
|
|
"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/stats"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
"github.com/seaweedfs/seaweedfs/weed/util/version"
|
|
)
|
|
|
|
// masterVolumeProvider implements VolumeLocationProvider by querying master
|
|
// This is rarely called since master pushes updates proactively via KeepConnected stream
|
|
type masterVolumeProvider struct {
|
|
masterClient *MasterClient
|
|
}
|
|
|
|
func isCanceledErr(err error) bool {
|
|
if err == nil {
|
|
return false
|
|
}
|
|
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
|
|
return true
|
|
}
|
|
if statusErr, ok := status.FromError(err); ok {
|
|
switch statusErr.Code() {
|
|
case codes.Canceled, codes.DeadlineExceeded:
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// LookupVolumeIds queries the master for volume locations (fallback when cache misses).
|
|
// Returns partial results with aggregated errors for volumes that failed.
|
|
// Retries on codes.Unavailable (e.g. master warming up after restart) with backoff.
|
|
func (p *masterVolumeProvider) LookupVolumeIds(ctx context.Context, volumeIds []string) (map[string][]Location, error) {
|
|
var result map[string][]Location
|
|
var lookupErrors []error
|
|
|
|
glog.V(2).Infof("Looking up %d volumes from master: %v", len(volumeIds), volumeIds)
|
|
|
|
retryErr := util.RetryWithBackoff(ctx, "lookup", 30*time.Second,
|
|
func(err error) bool {
|
|
st, ok := status.FromError(err)
|
|
return ok && st.Code() == codes.Unavailable
|
|
},
|
|
func() error {
|
|
result = make(map[string][]Location)
|
|
lookupErrors = nil
|
|
|
|
// Per-attempt timeout bounds both master resolution and the RPC
|
|
// so a single attempt cannot consume the entire retry budget.
|
|
timeoutCtx, cancel := context.WithTimeout(ctx, p.masterClient.grpcTimeout)
|
|
defer cancel()
|
|
|
|
master := p.masterClient.GetMaster(timeoutCtx)
|
|
if master == "" {
|
|
if ctx.Err() != nil {
|
|
return ctx.Err()
|
|
}
|
|
return status.Errorf(codes.Unavailable, "no master available")
|
|
}
|
|
|
|
return pb.WithMasterClient(timeoutCtx, false, master, p.masterClient.grpcDialOption, false, func(client master_pb.SeaweedClient) error {
|
|
resp, err := client.LookupVolume(timeoutCtx, &master_pb.LookupVolumeRequest{
|
|
VolumeOrFileIds: volumeIds,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, vidLoc := range resp.VolumeIdLocations {
|
|
// Preserve per-volume errors from master response
|
|
// These could indicate misconfiguration, volume deletion, etc.
|
|
if vidLoc.Error != "" {
|
|
lookupErrors = append(lookupErrors, fmt.Errorf("volume %s: %s", vidLoc.VolumeOrFileId, vidLoc.Error))
|
|
glog.V(1).Infof("volume %s lookup error from master: %s", vidLoc.VolumeOrFileId, vidLoc.Error)
|
|
continue
|
|
}
|
|
|
|
// Parse volume ID from response
|
|
parts := strings.Split(vidLoc.VolumeOrFileId, ",")
|
|
vidOnly := parts[0]
|
|
vid, err := strconv.ParseUint(vidOnly, 10, 32)
|
|
if err != nil {
|
|
lookupErrors = append(lookupErrors, fmt.Errorf("volume %s: invalid volume ID format: %w", vidLoc.VolumeOrFileId, err))
|
|
glog.Warningf("Failed to parse volume id '%s' from master response '%s': %v", vidOnly, vidLoc.VolumeOrFileId, err)
|
|
continue
|
|
}
|
|
|
|
var locations []Location
|
|
for _, masterLoc := range vidLoc.Locations {
|
|
loc := Location{
|
|
Url: masterLoc.Url,
|
|
PublicUrl: masterLoc.PublicUrl,
|
|
GrpcPort: int(masterLoc.GrpcPort),
|
|
DataCenter: masterLoc.DataCenter,
|
|
DataInRemote: masterLoc.DataInRemote,
|
|
ReadOnly: masterLoc.ReadOnly,
|
|
ReadOnlyCanDelete: masterLoc.ReadOnlyCanDelete,
|
|
}
|
|
// Update cache with the location
|
|
p.masterClient.addLocation(uint32(vid), loc)
|
|
locations = append(locations, loc)
|
|
}
|
|
|
|
if len(locations) > 0 {
|
|
result[vidOnly] = locations
|
|
}
|
|
}
|
|
return nil
|
|
})
|
|
})
|
|
if retryErr != nil {
|
|
return nil, retryErr
|
|
}
|
|
|
|
// Return partial results with detailed errors
|
|
// Callers should check both result map and error
|
|
if len(lookupErrors) > 0 {
|
|
glog.V(2).Infof("MasterClient: looked up %d volumes, found %d, %d errors", len(volumeIds), len(result), len(lookupErrors))
|
|
return result, fmt.Errorf("master volume lookup errors: %w", errors.Join(lookupErrors...))
|
|
}
|
|
|
|
glog.V(3).Infof("MasterClient: looked up %d volumes, found %d", len(volumeIds), len(result))
|
|
return result, nil
|
|
}
|
|
|
|
// MasterClient connects to master servers and maintains volume location cache
|
|
// It receives real-time updates via KeepConnected streaming and uses vidMapClient for caching
|
|
type MasterClient struct {
|
|
*vidMapClient // Embedded cache with shared logic
|
|
|
|
FilerGroup string
|
|
clientType string
|
|
clientHost pb.ServerAddress
|
|
rack string
|
|
currentMaster pb.ServerAddress
|
|
lastServedMaster pb.ServerAddress
|
|
currentMasterLock sync.RWMutex
|
|
masters pb.ServerDiscovery
|
|
grpcDialOption grpc.DialOption
|
|
grpcTimeout time.Duration // Timeout for gRPC calls to master
|
|
OnPeerUpdate func(update *master_pb.ClusterNodeUpdate, startFrom time.Time)
|
|
OnPeerUpdateLock sync.RWMutex
|
|
OnLockRingUpdate func(update *master_pb.LockRingUpdate)
|
|
OnLockRingUpdateLock sync.RWMutex
|
|
OnMasterChange func(previous, current pb.ServerAddress)
|
|
OnMasterChangeLock sync.RWMutex
|
|
}
|
|
|
|
func NewMasterClient(grpcDialOption grpc.DialOption, filerGroup string, clientType string, clientHost pb.ServerAddress, clientDataCenter string, rack string, masters pb.ServerDiscovery) *MasterClient {
|
|
mc := &MasterClient{
|
|
FilerGroup: filerGroup,
|
|
clientType: clientType,
|
|
clientHost: clientHost,
|
|
rack: rack,
|
|
masters: masters,
|
|
grpcDialOption: grpcDialOption,
|
|
grpcTimeout: 5 * time.Second, // Default: 5 seconds for gRPC calls to master
|
|
}
|
|
|
|
// Create provider that references this MasterClient
|
|
provider := &masterVolumeProvider{masterClient: mc}
|
|
|
|
// Initialize embedded vidMapClient with the provider and default cache size
|
|
mc.vidMapClient = newVidMapClient(provider, clientDataCenter, DefaultVidMapCacheSize)
|
|
|
|
return mc
|
|
}
|
|
|
|
func (mc *MasterClient) SetOnPeerUpdateFn(onPeerUpdate func(update *master_pb.ClusterNodeUpdate, startFrom time.Time)) {
|
|
mc.OnPeerUpdateLock.Lock()
|
|
mc.OnPeerUpdate = onPeerUpdate
|
|
mc.OnPeerUpdateLock.Unlock()
|
|
}
|
|
|
|
func (mc *MasterClient) SetOnLockRingUpdateFn(fn func(update *master_pb.LockRingUpdate)) {
|
|
mc.OnLockRingUpdateLock.Lock()
|
|
mc.OnLockRingUpdate = fn
|
|
mc.OnLockRingUpdateLock.Unlock()
|
|
}
|
|
|
|
func (mc *MasterClient) SetOnMasterChangeFn(fn func(previous, current pb.ServerAddress)) {
|
|
mc.OnMasterChangeLock.Lock()
|
|
mc.OnMasterChange = fn
|
|
mc.OnMasterChangeLock.Unlock()
|
|
}
|
|
|
|
func (mc *MasterClient) tryAllMasters(ctx context.Context) {
|
|
var nextHintedLeader pb.ServerAddress
|
|
failedMasters := make(map[pb.ServerAddress]struct{})
|
|
mc.masters.RefreshBySrvIfAvailable()
|
|
for _, master := range mc.masters.GetInstances() {
|
|
if _, failed := failedMasters[master]; failed {
|
|
continue
|
|
}
|
|
nextHintedLeader = mc.tryConnectToMaster(ctx, master)
|
|
for nextHintedLeader != "" {
|
|
if _, failed := failedMasters[nextHintedLeader]; failed {
|
|
break // don't follow redirect to a known-unreachable master
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
glog.V(0).Infof("Connection attempt to all masters stopped: %v", ctx.Err())
|
|
return
|
|
default:
|
|
target := nextHintedLeader
|
|
nextHintedLeader = mc.tryConnectToMaster(ctx, target)
|
|
if nextHintedLeader == "" {
|
|
// connection to target failed; remember it so we skip
|
|
// stale redirects pointing back to it this cycle
|
|
failedMasters[target] = struct{}{}
|
|
}
|
|
}
|
|
}
|
|
mc.setCurrentMaster("")
|
|
}
|
|
}
|
|
|
|
func (mc *MasterClient) tryConnectToMaster(ctx context.Context, master pb.ServerAddress) (nextHintedLeader pb.ServerAddress) {
|
|
glog.V(1).Infof("%s.%s masterClient Connecting to master %v", mc.FilerGroup, mc.clientType, master)
|
|
stats.MasterClientConnectCounter.WithLabelValues("total").Inc()
|
|
connectStartTime := time.Now()
|
|
gprcErr := pb.WithMasterClient(context.Background(), true, master, mc.grpcDialOption, false, func(client master_pb.SeaweedClient) error {
|
|
ctx, cancel := context.WithCancel(ctx)
|
|
defer cancel()
|
|
|
|
stream, err := client.KeepConnected(ctx)
|
|
if err != nil {
|
|
glog.V(1).Infof("%s.%s masterClient failed to keep connected to %s: %v", mc.FilerGroup, mc.clientType, master, err)
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.FailedToKeepConnected).Inc()
|
|
return err
|
|
}
|
|
glog.V(1).Infof("%s.%s masterClient gRPC stream established to %s in %v", mc.FilerGroup, mc.clientType, master, time.Since(connectStartTime))
|
|
|
|
if err = stream.Send(&master_pb.KeepConnectedRequest{
|
|
FilerGroup: mc.FilerGroup,
|
|
DataCenter: mc.GetDataCenter(),
|
|
Rack: mc.rack,
|
|
ClientType: mc.clientType,
|
|
ClientAddress: string(mc.clientHost),
|
|
Version: version.Version(),
|
|
}); err != nil {
|
|
glog.V(0).Infof("%s.%s masterClient failed to send to %s: %v", mc.FilerGroup, mc.clientType, master, err)
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.FailedToSend).Inc()
|
|
return err
|
|
}
|
|
glog.V(1).Infof("%s.%s masterClient Connected to %v", mc.FilerGroup, mc.clientType, master)
|
|
|
|
resp, err := stream.Recv()
|
|
if err != nil {
|
|
canceled := isCanceledErr(err) || ctx.Err() != nil
|
|
if canceled {
|
|
glog.V(1).Infof("%s.%s masterClient stream closed from %s: %v", mc.FilerGroup, mc.clientType, master, err)
|
|
} else {
|
|
glog.V(0).Infof("%s.%s masterClient failed to receive from %s: %v", mc.FilerGroup, mc.clientType, master, err)
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.FailedToReceive).Inc()
|
|
}
|
|
return err
|
|
}
|
|
|
|
// check if it is the leader to determine whether to reset the vidMap
|
|
if resp.VolumeLocation != nil {
|
|
if resp.VolumeLocation.Leader != "" && !master.Equals(pb.ServerAddress(resp.VolumeLocation.Leader)) {
|
|
glog.V(1).Infof("master %v redirected to leader %v", master, resp.VolumeLocation.Leader)
|
|
nextHintedLeader = pb.ServerAddress(resp.VolumeLocation.Leader)
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.RedirectedToLeader).Inc()
|
|
return nil
|
|
}
|
|
mc.resetVidMap()
|
|
mc.updateVidMap(resp)
|
|
} else {
|
|
// First message from master is not VolumeLocation (e.g., ClusterNodeUpdate)
|
|
// Still need to reset cache to ensure we don't use stale data from previous master
|
|
mc.resetVidMap()
|
|
}
|
|
if previous := mc.markServingMaster(master); previous != "" {
|
|
mc.OnMasterChangeLock.RLock()
|
|
if mc.OnMasterChange != nil {
|
|
glog.V(0).Infof("%s.%s masterClient master changed %s -> %s", mc.FilerGroup, mc.clientType, previous, master)
|
|
mc.OnMasterChange(previous, master)
|
|
}
|
|
mc.OnMasterChangeLock.RUnlock()
|
|
}
|
|
|
|
for {
|
|
resp, err := stream.Recv()
|
|
if err != nil {
|
|
canceled := isCanceledErr(err) || ctx.Err() != nil
|
|
if canceled {
|
|
glog.V(1).Infof("%s.%s masterClient stream closed from %s: %v", mc.FilerGroup, mc.clientType, master, err)
|
|
} else {
|
|
glog.V(0).Infof("%s.%s masterClient failed to receive from %s: %v", mc.FilerGroup, mc.clientType, master, err)
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.FailedToReceive).Inc()
|
|
}
|
|
return err
|
|
}
|
|
|
|
if resp.VolumeLocation != nil {
|
|
// Check for leader change during the stream
|
|
// If master announces a new leader, reconnect to it
|
|
if currentMaster := mc.GetMaster(ctx); resp.VolumeLocation.Leader != "" && !currentMaster.Equals(pb.ServerAddress(resp.VolumeLocation.Leader)) {
|
|
glog.V(1).Infof("currentMaster %v redirected to leader %v", currentMaster, resp.VolumeLocation.Leader)
|
|
nextHintedLeader = pb.ServerAddress(resp.VolumeLocation.Leader)
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.RedirectedToLeader).Inc()
|
|
return nil
|
|
}
|
|
mc.updateVidMap(resp)
|
|
}
|
|
if resp.ClusterNodeUpdate != nil {
|
|
update := resp.ClusterNodeUpdate
|
|
mc.OnPeerUpdateLock.RLock()
|
|
if mc.OnPeerUpdate != nil {
|
|
if update.FilerGroup == mc.FilerGroup {
|
|
if update.IsAdd {
|
|
glog.V(0).Infof("+ %s@%s noticed %s.%s %s\n", mc.clientType, mc.clientHost, update.FilerGroup, update.NodeType, update.Address)
|
|
} else {
|
|
glog.V(0).Infof("- %s@%s noticed %s.%s %s\n", mc.clientType, mc.clientHost, update.FilerGroup, update.NodeType, update.Address)
|
|
}
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.OnPeerUpdate).Inc()
|
|
mc.OnPeerUpdate(update, time.Now())
|
|
}
|
|
}
|
|
mc.OnPeerUpdateLock.RUnlock()
|
|
}
|
|
if resp.LockRingUpdate != nil {
|
|
update := resp.LockRingUpdate
|
|
mc.OnLockRingUpdateLock.RLock()
|
|
if mc.OnLockRingUpdate != nil {
|
|
if update.FilerGroup == mc.FilerGroup {
|
|
glog.V(0).Infof("LockRing: %s@%s received ring update v%d: %v", mc.clientType, mc.clientHost, update.Version, update.Servers)
|
|
mc.OnLockRingUpdate(update)
|
|
}
|
|
}
|
|
mc.OnLockRingUpdateLock.RUnlock()
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
if isCanceledErr(err) {
|
|
glog.V(1).Infof("Connection attempt to master stopped: %v", err)
|
|
} else {
|
|
glog.V(0).Infof("Connection attempt to master stopped: %v", err)
|
|
}
|
|
return err
|
|
}
|
|
}
|
|
})
|
|
if gprcErr != nil {
|
|
if isCanceledErr(gprcErr) || ctx.Err() != nil {
|
|
glog.V(1).Infof("%s.%s masterClient connection closed to %v: %v", mc.FilerGroup, mc.clientType, master, gprcErr)
|
|
return nextHintedLeader
|
|
}
|
|
stats.MasterClientConnectCounter.WithLabelValues(stats.Failed).Inc()
|
|
glog.V(1).Infof("%s.%s masterClient failed to connect with master %v: %v", mc.FilerGroup, mc.clientType, master, gprcErr)
|
|
}
|
|
return nextHintedLeader
|
|
}
|
|
|
|
// addedVids indexes added ids that are also being removed. A volume moved
|
|
// between a server's disks is reported both ways in one message, and the server
|
|
// still has it -- acting on the removal would drop a good location, whichever
|
|
// order the two lists happen to be applied in. Empty unless both lists are
|
|
// non-empty, so the full id list a client gets on connect costs nothing.
|
|
func addedVids(added, removed []uint32) map[uint32]struct{} {
|
|
if len(added) == 0 || len(removed) == 0 {
|
|
return nil
|
|
}
|
|
index := make(map[uint32]struct{}, len(added))
|
|
for _, vid := range added {
|
|
index[vid] = struct{}{}
|
|
}
|
|
return index
|
|
}
|
|
|
|
// updateVidMap applies a KeepConnectedResponse volume-location message to the
|
|
// local vidMap. NewVids adds the entries; RemoteVids names the subset of them
|
|
// backed by remote storage, which is added with DataInRemote so read paths can
|
|
// prefer the cheap local replica. DeletedVids drops the named entry unless the
|
|
// same message also added it back (volume moved between this server's disks).
|
|
// EC vid changes go through the parallel addEcLocation / deleteEcLocation pair.
|
|
func (mc *MasterClient) updateVidMap(resp *master_pb.KeepConnectedResponse) {
|
|
if resp.VolumeLocation.IsEmptyUrl() {
|
|
glog.V(0).Infof("updateVidMap ignore short heartbeat: %+v", resp)
|
|
return
|
|
}
|
|
// process new volume location
|
|
loc := Location{
|
|
Url: resp.VolumeLocation.Url,
|
|
PublicUrl: resp.VolumeLocation.PublicUrl,
|
|
DataCenter: resp.VolumeLocation.DataCenter,
|
|
GrpcPort: int(resp.VolumeLocation.GrpcPort),
|
|
}
|
|
stillOnServer := addedVids(resp.VolumeLocation.NewVids, resp.VolumeLocation.DeletedVids)
|
|
// RemoteVids repeats ids NewVids already carries, so the tier is settled
|
|
// before anything is written rather than adding each one twice.
|
|
var remoteVids map[uint32]struct{}
|
|
if len(resp.VolumeLocation.RemoteVids) > 0 {
|
|
remoteVids = make(map[uint32]struct{}, len(resp.VolumeLocation.RemoteVids))
|
|
for _, vid := range resp.VolumeLocation.RemoteVids {
|
|
remoteVids[vid] = struct{}{}
|
|
}
|
|
}
|
|
var readOnlyVids map[uint32]struct{}
|
|
if len(resp.VolumeLocation.ReadOnlyVids) > 0 {
|
|
readOnlyVids = make(map[uint32]struct{}, len(resp.VolumeLocation.ReadOnlyVids))
|
|
for _, vid := range resp.VolumeLocation.ReadOnlyVids {
|
|
readOnlyVids[vid] = struct{}{}
|
|
}
|
|
}
|
|
var readOnlyCanDeleteVids map[uint32]struct{}
|
|
if len(resp.VolumeLocation.ReadOnlyCanDeleteVids) > 0 {
|
|
readOnlyCanDeleteVids = make(map[uint32]struct{}, len(resp.VolumeLocation.ReadOnlyCanDeleteVids))
|
|
for _, vid := range resp.VolumeLocation.ReadOnlyCanDeleteVids {
|
|
readOnlyCanDeleteVids[vid] = struct{}{}
|
|
}
|
|
}
|
|
for _, newVid := range resp.VolumeLocation.NewVids {
|
|
if _, isRemote := remoteVids[newVid]; isRemote {
|
|
continue
|
|
}
|
|
newLoc := loc
|
|
if _, isReadOnly := readOnlyVids[newVid]; isReadOnly {
|
|
newLoc.ReadOnly = true
|
|
}
|
|
if _, canDelete := readOnlyCanDeleteVids[newVid]; canDelete {
|
|
newLoc.ReadOnlyCanDelete = true
|
|
}
|
|
glog.V(2).Infof("%s.%s: %s masterClient adds volume %d", mc.FilerGroup, mc.clientType, newLoc.Url, newVid)
|
|
mc.addLocation(newVid, newLoc)
|
|
}
|
|
for _, remoteVid := range resp.VolumeLocation.RemoteVids {
|
|
remoteLoc := loc
|
|
remoteLoc.DataInRemote = true
|
|
if _, isReadOnly := readOnlyVids[remoteVid]; isReadOnly {
|
|
remoteLoc.ReadOnly = true
|
|
}
|
|
if _, canDelete := readOnlyCanDeleteVids[remoteVid]; canDelete {
|
|
remoteLoc.ReadOnlyCanDelete = true
|
|
}
|
|
glog.V(2).Infof("%s.%s: %s masterClient adds remote volume %d", mc.FilerGroup, mc.clientType, remoteLoc.Url, remoteVid)
|
|
mc.addLocation(remoteVid, remoteLoc)
|
|
}
|
|
for _, deletedVid := range resp.VolumeLocation.DeletedVids {
|
|
if _, moved := stillOnServer[deletedVid]; moved {
|
|
continue
|
|
}
|
|
glog.V(2).Infof("%s.%s: %s masterClient removes volume %d", mc.FilerGroup, mc.clientType, loc.Url, deletedVid)
|
|
mc.deleteLocation(deletedVid, loc)
|
|
}
|
|
stillOnServerEc := addedVids(resp.VolumeLocation.NewEcVids, resp.VolumeLocation.DeletedEcVids)
|
|
for _, newEcVid := range resp.VolumeLocation.NewEcVids {
|
|
glog.V(2).Infof("%s.%s: %s masterClient adds ec volume %d", mc.FilerGroup, mc.clientType, loc.Url, newEcVid)
|
|
mc.addEcLocation(newEcVid, loc)
|
|
}
|
|
for _, deletedEcVid := range resp.VolumeLocation.DeletedEcVids {
|
|
if _, moved := stillOnServerEc[deletedEcVid]; moved {
|
|
continue
|
|
}
|
|
glog.V(2).Infof("%s.%s: %s masterClient removes ec volume %d", mc.FilerGroup, mc.clientType, loc.Url, deletedEcVid)
|
|
mc.deleteEcLocation(deletedEcVid, loc)
|
|
}
|
|
glog.V(1).Infof("updateVidMap(%s) %s.%s: %s volume add local: %d, remote: %d, del: %d, add ec: %d del ec: %d",
|
|
resp.VolumeLocation.DataCenter, mc.FilerGroup, mc.clientType, loc.Url,
|
|
len(resp.VolumeLocation.NewVids)-len(resp.VolumeLocation.RemoteVids),
|
|
len(resp.VolumeLocation.RemoteVids),
|
|
len(resp.VolumeLocation.DeletedVids), len(resp.VolumeLocation.NewEcVids),
|
|
len(resp.VolumeLocation.DeletedEcVids))
|
|
}
|
|
|
|
func (mc *MasterClient) WithClient(ctx context.Context, streamingMode bool, fn func(client master_pb.SeaweedClient) error) error {
|
|
getMasterF := func() pb.ServerAddress {
|
|
return mc.GetMaster(ctx)
|
|
}
|
|
return mc.WithClientCustomGetMaster(ctx, getMasterF, streamingMode, fn)
|
|
}
|
|
|
|
// WithClientCustomGetMaster bounds the wait for a master leader by ctx, so a
|
|
// caller with a deadline is not parked for the length of an election. The dial
|
|
// still gets context.Background(): fn brings its own RPC context, so nothing
|
|
// here can attribute a cancellation to the shared connection.
|
|
func (mc *MasterClient) WithClientCustomGetMaster(ctx context.Context, getMasterF func() pb.ServerAddress, streamingMode bool, fn func(client master_pb.SeaweedClient) error) error {
|
|
return util.RetryWithBackoff(ctx, "master grpc", util.RetryWaitTime, util.IsTransientError, func() error {
|
|
master := getMasterF()
|
|
if err := ctx.Err(); err != nil {
|
|
return err
|
|
}
|
|
return pb.WithMasterClient(context.Background(), streamingMode, master, mc.grpcDialOption, false, func(client master_pb.SeaweedClient) error {
|
|
return fn(client)
|
|
})
|
|
})
|
|
}
|
|
|
|
func (mc *MasterClient) getCurrentMaster() pb.ServerAddress {
|
|
mc.currentMasterLock.RLock()
|
|
defer mc.currentMasterLock.RUnlock()
|
|
return mc.currentMaster
|
|
}
|
|
|
|
func (mc *MasterClient) setCurrentMaster(master pb.ServerAddress) {
|
|
mc.currentMasterLock.Lock()
|
|
mc.currentMaster = master
|
|
mc.currentMasterLock.Unlock()
|
|
}
|
|
|
|
// markServingMaster records the master now serving this client and returns
|
|
// the previously served one when it differs. Unlike currentMaster,
|
|
// lastServedMaster survives the disconnected gap between reconnect attempts,
|
|
// so a leader change is still detected.
|
|
func (mc *MasterClient) markServingMaster(master pb.ServerAddress) (previous pb.ServerAddress) {
|
|
mc.currentMasterLock.Lock()
|
|
defer mc.currentMasterLock.Unlock()
|
|
if mc.lastServedMaster != "" && !mc.lastServedMaster.Equals(master) {
|
|
previous = mc.lastServedMaster
|
|
}
|
|
mc.currentMaster = master
|
|
mc.lastServedMaster = master
|
|
return
|
|
}
|
|
|
|
// GetMaster returns the current master address, blocking until connected.
|
|
//
|
|
// IMPORTANT: This method blocks until KeepConnectedToMaster successfully establishes
|
|
// a connection to a master server. If KeepConnectedToMaster hasn't been started in a
|
|
// background goroutine, this will block indefinitely (or until ctx is canceled).
|
|
//
|
|
// Typical initialization pattern:
|
|
//
|
|
// mc := wdclient.NewMasterClient(...)
|
|
// go mc.KeepConnectedToMaster(ctx) // Start connection management
|
|
// // ... later ...
|
|
// master := mc.GetMaster(ctx) // Will block until connected
|
|
//
|
|
// If called before KeepConnectedToMaster establishes a connection, this may cause
|
|
// unexpected timeouts in LookupVolumeIds and other operations that depend on it.
|
|
func (mc *MasterClient) GetMaster(ctx context.Context) pb.ServerAddress {
|
|
mc.WaitUntilConnected(ctx)
|
|
return mc.getCurrentMaster()
|
|
}
|
|
|
|
// GetMasters returns all configured master addresses, blocking until connected.
|
|
// See GetMaster() for important initialization contract details.
|
|
func (mc *MasterClient) GetMasters(ctx context.Context) []pb.ServerAddress {
|
|
mc.WaitUntilConnected(ctx)
|
|
return mc.masters.GetInstances()
|
|
}
|
|
|
|
// ListMasterSet returns a set of configured master addresses keyed by their
|
|
// canonical http form. Unlike GetMasters this does not wait for a connection,
|
|
// so it is safe to call from admission paths that must stay non-blocking.
|
|
func (mc *MasterClient) ListMasterSet() map[string]struct{} {
|
|
addrs := mc.masters.GetInstances()
|
|
set := make(map[string]struct{}, len(addrs))
|
|
for _, a := range addrs {
|
|
set[a.ToHttpAddress()] = struct{}{}
|
|
}
|
|
return set
|
|
}
|
|
|
|
// WaitUntilConnected blocks until a master connection is established or ctx is canceled.
|
|
// This does NOT initiate connections - it only waits for KeepConnectedToMaster to succeed.
|
|
func (mc *MasterClient) WaitUntilConnected(ctx context.Context) {
|
|
attempts := 0
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
default:
|
|
currentMaster := mc.getCurrentMaster()
|
|
if currentMaster != "" {
|
|
return
|
|
}
|
|
attempts++
|
|
if attempts%100 == 0 { // Log every 100 attempts (roughly every 20 seconds)
|
|
glog.V(0).Infof("%s.%s WaitUntilConnected still waiting for master connection (attempt %d)...", mc.FilerGroup, mc.clientType, attempts)
|
|
}
|
|
// Use select with time.After to respect context cancellation during sleep
|
|
sleepDuration := time.Duration(rand.Int31n(200)) * time.Millisecond
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-time.After(sleepDuration):
|
|
// continue to next iteration
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (mc *MasterClient) KeepConnectedToMaster(ctx context.Context) {
|
|
glog.V(1).Infof("%s.%s masterClient bootstraps with masters %v", mc.FilerGroup, mc.clientType, mc.masters)
|
|
reconnectCount := 0
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
if isCanceledErr(ctx.Err()) {
|
|
glog.V(1).Infof("Connection to masters stopped: %v", ctx.Err())
|
|
} else {
|
|
glog.V(0).Infof("Connection to masters stopped: %v", ctx.Err())
|
|
}
|
|
return
|
|
default:
|
|
reconnectStart := time.Now()
|
|
if reconnectCount > 0 {
|
|
glog.V(0).Infof("%s.%s masterClient reconnection attempt #%d", mc.FilerGroup, mc.clientType, reconnectCount)
|
|
}
|
|
mc.tryAllMasters(ctx)
|
|
reconnectCount++
|
|
glog.V(1).Infof("%s.%s masterClient connection cycle completed in %v, sleeping before retry",
|
|
mc.FilerGroup, mc.clientType, time.Since(reconnectStart))
|
|
time.Sleep(time.Second)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (mc *MasterClient) FindLeaderFromOtherPeers(myMasterAddress pb.ServerAddress) (leader string) {
|
|
for _, master := range mc.masters.GetInstances() {
|
|
if master == myMasterAddress {
|
|
continue
|
|
}
|
|
if grpcErr := pb.WithMasterClient(context.Background(), false, master, mc.grpcDialOption, false, func(client master_pb.SeaweedClient) error {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Millisecond)
|
|
defer cancel()
|
|
resp, err := client.GetMasterConfiguration(ctx, &master_pb.GetMasterConfigurationRequest{})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
leader = resp.Leader
|
|
return nil
|
|
}); grpcErr != nil {
|
|
glog.V(0).Infof("connect to %s: %v", master, grpcErr)
|
|
}
|
|
if leader != "" {
|
|
glog.V(0).Infof("existing leader is %s", leader)
|
|
return
|
|
}
|
|
}
|
|
glog.V(0).Infof("No existing leader found!")
|
|
return
|
|
}
|