Files
seaweedfs/weed/wdclient/masterclient.go
T
Peter Dodd 52c6df3bec fix(filer): leave the lock ring before stopping gRPC on shutdown (#11604)
A draining filer stayed in the master's lock ring until its process
exited, while gRPC GracefulStop was already refusing new connections. S3
gateways and peer filers kept routing object-write locks and owner-routed
writes to it for the whole graceful-stop window.

On shutdown the filer now sends leave_lock_ring on its open KeepConnected
stream. The master removes it from the lock ring only; it stays a cluster
member so peers keep following its metadata log through the drain. Once
the ring update without it arrives, the filer has already transferred its
locks to the new owners, and it keeps serving through the prior-owner
window (plus a second for peers that apply the update later) before gRPC
and HTTP begin draining. The whole leave is bounded at 10s so slow lock
transfers or a stuck stream cannot hold up the drain. A lone filer, or a
master that ignores the message, falls back to the previous behavior.
2026-10-06 09:53:40 +08:00

686 lines
25 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
// streamLock publishes the stream only after its registration send, so a
// leave sent afterwards never races that send.
streamLock sync.Mutex
stream master_pb.Seaweed_KeepConnectedClient
leavingLockRing bool
}
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))
mc.streamLock.Lock()
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(),
LeaveLockRing: mc.leavingLockRing,
})
if err == nil {
mc.stream = stream
}
mc.streamLock.Unlock()
if 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
}
defer func() {
mc.streamLock.Lock()
if mc.stream == stream {
mc.stream = nil
}
mc.streamLock.Unlock()
}()
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
}
// LeaveLockRing asks the master to drop this client from the lock ring while it
// stays connected, and keeps it out of the ring on any later reconnect.
func (mc *MasterClient) LeaveLockRing() error {
mc.streamLock.Lock()
mc.leavingLockRing = true
stream := mc.stream
mc.streamLock.Unlock()
if stream == nil {
return nil
}
return stream.Send(&master_pb.KeepConnectedRequest{LeaveLockRing: true})
}
// 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
}