mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-08 15:27:43 +02:00
Filer, S3 and broker report it on the KeepConnected registration, volume servers in their heartbeat, and masters return their own in GetMasterConfiguration. MasterClient gains SetMetricsPort so the eight callers that have no metrics listener are untouched, and the volume server takes it as a constructor argument because its heartbeat goroutine starts there. In combined "weed server" one metrics listener serves the whole shared registry, so every component advertises the same port.
197 lines
6.0 KiB
Go
197 lines
6.0 KiB
Go
package cluster
|
|
|
|
import (
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
|
)
|
|
|
|
const (
|
|
MasterType = "master"
|
|
VolumeServerType = "volumeServer"
|
|
FilerType = "filer"
|
|
BrokerType = "broker"
|
|
S3Type = "s3"
|
|
)
|
|
|
|
type FilerGroupName string
|
|
type DataCenter string
|
|
type Rack string
|
|
|
|
type ClusterNode struct {
|
|
Address pb.ServerAddress
|
|
Version string
|
|
counter int
|
|
CreatedTs time.Time
|
|
DataCenter DataCenter
|
|
Rack Rack
|
|
// MetricsPort is the node's Prometheus /metrics port, or 0 when the node
|
|
// does not run a metrics listener.
|
|
MetricsPort uint32
|
|
}
|
|
|
|
type ClusterNodeGroups struct {
|
|
groupMembers map[FilerGroupName]*GroupMembers
|
|
sync.RWMutex
|
|
}
|
|
type Cluster struct {
|
|
filerGroups *ClusterNodeGroups
|
|
brokerGroups *ClusterNodeGroups
|
|
s3Groups *ClusterNodeGroups
|
|
}
|
|
|
|
func newClusterNodeGroups() *ClusterNodeGroups {
|
|
return &ClusterNodeGroups{
|
|
groupMembers: map[FilerGroupName]*GroupMembers{},
|
|
}
|
|
}
|
|
func (g *ClusterNodeGroups) getGroupMembers(filerGroup FilerGroupName, createIfNotFound bool) *GroupMembers {
|
|
members, found := g.groupMembers[filerGroup]
|
|
if !found && createIfNotFound {
|
|
members = newGroupMembers()
|
|
g.groupMembers[filerGroup] = members
|
|
}
|
|
return members
|
|
}
|
|
|
|
func (g *ClusterNodeGroups) AddClusterNode(filerGroup FilerGroupName, nodeType string, dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string, metricsPort uint32) []*master_pb.KeepConnectedResponse {
|
|
g.Lock()
|
|
defer g.Unlock()
|
|
m := g.getGroupMembers(filerGroup, true)
|
|
if t := m.addMember(dataCenter, rack, address, version, metricsPort); t != nil {
|
|
return buildClusterNodeUpdateMessage(true, filerGroup, nodeType, address)
|
|
}
|
|
return nil
|
|
}
|
|
func (g *ClusterNodeGroups) RemoveClusterNode(filerGroup FilerGroupName, nodeType string, address pb.ServerAddress) []*master_pb.KeepConnectedResponse {
|
|
g.Lock()
|
|
defer g.Unlock()
|
|
m := g.getGroupMembers(filerGroup, false)
|
|
if m == nil {
|
|
return nil
|
|
}
|
|
if m.removeMember(address) {
|
|
return buildClusterNodeUpdateMessage(false, filerGroup, nodeType, address)
|
|
}
|
|
return nil
|
|
}
|
|
func (g *ClusterNodeGroups) ListClusterNode(filerGroup FilerGroupName) (nodes []*ClusterNode) {
|
|
g.Lock()
|
|
defer g.Unlock()
|
|
m := g.getGroupMembers(filerGroup, false)
|
|
if m == nil {
|
|
return nil
|
|
}
|
|
for _, node := range m.members {
|
|
nodes = append(nodes, node)
|
|
}
|
|
return
|
|
}
|
|
|
|
func NewCluster() *Cluster {
|
|
return &Cluster{
|
|
filerGroups: newClusterNodeGroups(),
|
|
brokerGroups: newClusterNodeGroups(),
|
|
s3Groups: newClusterNodeGroups(),
|
|
}
|
|
}
|
|
|
|
func (cluster *Cluster) AddClusterNode(ns, nodeType string, dataCenter DataCenter, rack Rack, address pb.ServerAddress, version string, metricsPort uint32) []*master_pb.KeepConnectedResponse {
|
|
filerGroup := FilerGroupName(ns)
|
|
switch nodeType {
|
|
case FilerType:
|
|
return cluster.filerGroups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version, metricsPort)
|
|
case BrokerType:
|
|
return cluster.brokerGroups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version, metricsPort)
|
|
case S3Type:
|
|
return cluster.s3Groups.AddClusterNode(filerGroup, nodeType, dataCenter, rack, address, version, metricsPort)
|
|
case MasterType:
|
|
return buildClusterNodeUpdateMessage(true, filerGroup, nodeType, address)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (cluster *Cluster) RemoveClusterNode(ns string, nodeType string, address pb.ServerAddress) []*master_pb.KeepConnectedResponse {
|
|
filerGroup := FilerGroupName(ns)
|
|
switch nodeType {
|
|
case FilerType:
|
|
return cluster.filerGroups.RemoveClusterNode(filerGroup, nodeType, address)
|
|
case BrokerType:
|
|
return cluster.brokerGroups.RemoveClusterNode(filerGroup, nodeType, address)
|
|
case S3Type:
|
|
return cluster.s3Groups.RemoveClusterNode(filerGroup, nodeType, address)
|
|
case MasterType:
|
|
return buildClusterNodeUpdateMessage(false, filerGroup, nodeType, address)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (cluster *Cluster) ListClusterNode(filerGroup FilerGroupName, nodeType string) (nodes []*ClusterNode) {
|
|
switch nodeType {
|
|
case FilerType:
|
|
return cluster.filerGroups.ListClusterNode(filerGroup)
|
|
case BrokerType:
|
|
return cluster.brokerGroups.ListClusterNode(filerGroup)
|
|
case S3Type:
|
|
return cluster.s3Groups.ListClusterNode(filerGroup)
|
|
case MasterType:
|
|
}
|
|
return
|
|
}
|
|
|
|
// ListClusterNodeUpdates reports the current members as add updates, so a
|
|
// client that just connected can rebuild the membership it missed while it was
|
|
// away.
|
|
func (cluster *Cluster) ListClusterNodeUpdates(filerGroup FilerGroupName, nodeType string) (updates []*master_pb.KeepConnectedResponse) {
|
|
for _, node := range cluster.ListClusterNode(filerGroup, nodeType) {
|
|
updates = append(updates, buildClusterNodeUpdateMessage(true, filerGroup, nodeType, node.Address)...)
|
|
}
|
|
return
|
|
}
|
|
|
|
// IsKnownNode reports whether address is currently registered under nodeType
|
|
// in any filer group. The lookup is intentionally group-agnostic because callers
|
|
// (e.g. Ping admission) only know the target address, not the group it joined.
|
|
func (cluster *Cluster) IsKnownNode(nodeType string, address pb.ServerAddress) bool {
|
|
var groups *ClusterNodeGroups
|
|
switch nodeType {
|
|
case FilerType:
|
|
groups = cluster.filerGroups
|
|
case BrokerType:
|
|
groups = cluster.brokerGroups
|
|
case S3Type:
|
|
groups = cluster.s3Groups
|
|
default:
|
|
return false
|
|
}
|
|
groups.RLock()
|
|
defer groups.RUnlock()
|
|
for _, members := range groups.groupMembers {
|
|
if _, found := members.members[address]; found {
|
|
return true
|
|
}
|
|
// fall back to a port-tolerant comparison so callers that omit the
|
|
// grpc-port suffix still match a registered peer
|
|
for stored := range members.members {
|
|
if stored.Equals(address) {
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func buildClusterNodeUpdateMessage(isAdd bool, filerGroup FilerGroupName, nodeType string, address pb.ServerAddress) (result []*master_pb.KeepConnectedResponse) {
|
|
result = append(result, &master_pb.KeepConnectedResponse{
|
|
ClusterNodeUpdate: &master_pb.ClusterNodeUpdate{
|
|
FilerGroup: string(filerGroup),
|
|
NodeType: nodeType,
|
|
Address: string(address),
|
|
IsAdd: isAdd,
|
|
},
|
|
})
|
|
return
|
|
}
|