@ -15,6 +15,15 @@ const (
BrokerType = "broker"
BrokerType = "broker"
)
)
type FilerGroup string
type Filers struct {
filers map [ pb . ServerAddress ] * ClusterNode
leaders * Leaders
}
type Leaders struct {
leaders [ 3 ] pb . ServerAddress
}
type ClusterNode struct {
type ClusterNode struct {
Address pb . ServerAddress
Address pb . ServerAddress
Version string
Version string
@ -22,42 +31,50 @@ type ClusterNode struct {
createdTs time . Time
createdTs time . Time
}
}
type Leaders struct {
leaders [ 3 ] pb . ServerAddress
}
type Cluster struct {
type Cluster struct {
filers map [ pb . ServerAddress ] * ClusterNode
filerGroup2filers map [ FilerGroup ] * Filers
filersLock sync . RWMutex
filersLock sync . RWMutex
filerLeaders * Leaders
brokers map [ pb . ServerAddress ] * ClusterNode
brokers map [ pb . ServerAddress ] * ClusterNode
brokersLock sync . RWMutex
brokersLock sync . RWMutex
}
}
func NewCluster ( ) * Cluster {
func NewCluster ( ) * Cluster {
return & Cluster {
return & Cluster {
filers : make ( map [ pb . ServerAddress ] * ClusterNode ) ,
filerLeaders : & Leaders { } ,
filerGroup2filers : make ( map [ FilerGroup ] * Filers ) ,
brokers : make ( map [ pb . ServerAddress ] * ClusterNode ) ,
brokers : make ( map [ pb . ServerAddress ] * ClusterNode ) ,
}
}
}
}
func ( cluster * Cluster ) AddClusterNode ( nodeType string , address pb . ServerAddress , version string ) [ ] * master_pb . KeepConnectedResponse {
switch nodeType {
case FilerType :
func ( cluster * Cluster ) getFilers ( filerGroup FilerGroup , createIfNotFound bool ) * Filers {
cluster . filersLock . Lock ( )
cluster . filersLock . Lock ( )
defer cluster . filersLock . Unlock ( )
defer cluster . filersLock . Unlock ( )
if existingNode , found := cluster . filers [ address ] ; found {
filers , found := cluster . filerGroup2filers [ filerGroup ]
if ! found && createIfNotFound {
filers = & Filers {
filers : make ( map [ pb . ServerAddress ] * ClusterNode ) ,
leaders : & Leaders { } ,
}
cluster . filerGroup2filers [ filerGroup ] = filers
}
return filers
}
func ( cluster * Cluster ) AddClusterNode ( ns , nodeType string , address pb . ServerAddress , version string ) [ ] * master_pb . KeepConnectedResponse {
filerGroup := FilerGroup ( ns )
switch nodeType {
case FilerType :
filers := cluster . getFilers ( filerGroup , true )
if existingNode , found := filers . filers [ address ] ; found {
existingNode . counter ++
existingNode . counter ++
return nil
return nil
}
}
cluster . filers [ address ] = & ClusterNode {
filers . filers [ address ] = & ClusterNode {
Address : address ,
Address : address ,
Version : version ,
Version : version ,
counter : 1 ,
counter : 1 ,
createdTs : time . Now ( ) ,
createdTs : time . Now ( ) ,
}
}
return cluster . ensureFilerLeaders ( true , nodeType , address )
return cluster . ensureFilerLeaders ( filers , true , filerGroup , nodeType , address )
case BrokerType :
case BrokerType :
cluster . brokersLock . Lock ( )
cluster . brokersLock . Lock ( )
defer cluster . brokersLock . Unlock ( )
defer cluster . brokersLock . Unlock ( )
@ -94,18 +111,21 @@ func (cluster *Cluster) AddClusterNode(nodeType string, address pb.ServerAddress
return nil
return nil
}
}
func ( cluster * Cluster ) RemoveClusterNode ( nodeType string , address pb . ServerAddress ) [ ] * master_pb . KeepConnectedResponse {
func ( cluster * Cluster ) RemoveClusterNode ( ns string , nodeType string , address pb . ServerAddress ) [ ] * master_pb . KeepConnectedResponse {
filerGroup := FilerGroup ( ns )
switch nodeType {
switch nodeType {
case FilerType :
case FilerType :
cluster . filersLock . Lock ( )
defer cluster . filersLock . Unlock ( )
if existingNode , found := cluster . filers [ address ] ; ! found {
filers := cluster . getFilers ( filerGroup , false )
if filers == nil {
return nil
}
if existingNode , found := filers . filers [ address ] ; ! found {
return nil
return nil
} else {
} else {
existingNode . counter --
existingNode . counter --
if existingNode . counter <= 0 {
if existingNode . counter <= 0 {
delete ( cluster . filers , address )
return cluster . ensureFilerLeaders ( false , nodeType , address )
delete ( filers . filers , address )
return cluster . ensureFilerLeaders ( filers , false , filerGroup , nodeType , address )
}
}
}
}
case BrokerType :
case BrokerType :
@ -142,12 +162,16 @@ func (cluster *Cluster) RemoveClusterNode(nodeType string, address pb.ServerAddr
return nil
return nil
}
}
func ( cluster * Cluster ) ListClusterNode ( nodeType string ) ( nodes [ ] * ClusterNode ) {
func ( cluster * Cluster ) ListClusterNode ( filerGroup FilerGroup , nodeType string ) ( nodes [ ] * ClusterNode ) {
switch nodeType {
switch nodeType {
case FilerType :
case FilerType :
filers := cluster . getFilers ( filerGroup , false )
if filers == nil {
return
}
cluster . filersLock . RLock ( )
cluster . filersLock . RLock ( )
defer cluster . filersLock . RUnlock ( )
defer cluster . filersLock . RUnlock ( )
for _ , node := range cluster . filers {
for _ , node := range filers . filers {
nodes = append ( nodes , node )
nodes = append ( nodes , node )
}
}
case BrokerType :
case BrokerType :
@ -161,16 +185,21 @@ func (cluster *Cluster) ListClusterNode(nodeType string) (nodes []*ClusterNode)
return
return
}
}
func ( cluster * Cluster ) IsOneLeader ( address pb . ServerAddress ) bool {
return cluster . filerLeaders . isOneLeader ( address )
func ( cluster * Cluster ) IsOneLeader ( filerGroup FilerGroup , address pb . ServerAddress ) bool {
filers := cluster . getFilers ( filerGroup , false )
if filers == nil {
return false
}
return filers . leaders . isOneLeader ( address )
}
}
func ( cluster * Cluster ) ensureFilerLeaders ( isAdd bool , nodeType string , address pb . ServerAddress ) ( result [ ] * master_pb . KeepConnectedResponse ) {
func ( cluster * Cluster ) ensureFilerLeaders ( filers * Filers , isAdd bool , filerGroup FilerGroup , nodeType string , address pb . ServerAddress ) ( result [ ] * master_pb . KeepConnectedResponse ) {
if isAdd {
if isAdd {
if cluster . filerL eaders. addLeaderIfVacant ( address ) {
if filers . l eaders. addLeaderIfVacant ( address ) {
// has added the address as one leader
// has added the address as one leader
result = append ( result , & master_pb . KeepConnectedResponse {
result = append ( result , & master_pb . KeepConnectedResponse {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
FilerGroup : string ( filerGroup ) ,
NodeType : nodeType ,
NodeType : nodeType ,
Address : string ( address ) ,
Address : string ( address ) ,
IsLeader : true ,
IsLeader : true ,
@ -180,6 +209,7 @@ func (cluster *Cluster) ensureFilerLeaders(isAdd bool, nodeType string, address
} else {
} else {
result = append ( result , & master_pb . KeepConnectedResponse {
result = append ( result , & master_pb . KeepConnectedResponse {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
FilerGroup : string ( filerGroup ) ,
NodeType : nodeType ,
NodeType : nodeType ,
Address : string ( address ) ,
Address : string ( address ) ,
IsLeader : false ,
IsLeader : false ,
@ -188,10 +218,11 @@ func (cluster *Cluster) ensureFilerLeaders(isAdd bool, nodeType string, address
} )
} )
}
}
} else {
} else {
if cluster . filerL eaders. removeLeaderIfExists ( address ) {
if filers . l eaders. removeLeaderIfExists ( address ) {
result = append ( result , & master_pb . KeepConnectedResponse {
result = append ( result , & master_pb . KeepConnectedResponse {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
FilerGroup : string ( filerGroup ) ,
NodeType : nodeType ,
NodeType : nodeType ,
Address : string ( address ) ,
Address : string ( address ) ,
IsLeader : true ,
IsLeader : true ,
@ -203,8 +234,8 @@ func (cluster *Cluster) ensureFilerLeaders(isAdd bool, nodeType string, address
var shortestDuration int64 = math . MaxInt64
var shortestDuration int64 = math . MaxInt64
now := time . Now ( )
now := time . Now ( )
var candidateAddress pb . ServerAddress
var candidateAddress pb . ServerAddress
for _ , node := range cluster . filers {
if cluster . filerL eaders. isOneLeader ( node . Address ) {
for _ , node := range filers . filers {
if filers . l eaders. isOneLeader ( node . Address ) {
continue
continue
}
}
duration := now . Sub ( node . createdTs ) . Nanoseconds ( )
duration := now . Sub ( node . createdTs ) . Nanoseconds ( )
@ -214,7 +245,7 @@ func (cluster *Cluster) ensureFilerLeaders(isAdd bool, nodeType string, address
}
}
}
}
if candidateAddress != "" {
if candidateAddress != "" {
cluster . filerL eaders. addLeaderIfVacant ( candidateAddress )
filers . l eaders. addLeaderIfVacant ( candidateAddress )
// added a new leader
// added a new leader
result = append ( result , & master_pb . KeepConnectedResponse {
result = append ( result , & master_pb . KeepConnectedResponse {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
@ -228,6 +259,7 @@ func (cluster *Cluster) ensureFilerLeaders(isAdd bool, nodeType string, address
} else {
} else {
result = append ( result , & master_pb . KeepConnectedResponse {
result = append ( result , & master_pb . KeepConnectedResponse {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
ClusterNodeUpdate : & master_pb . ClusterNodeUpdate {
FilerGroup : string ( filerGroup ) ,
NodeType : nodeType ,
NodeType : nodeType ,
Address : string ( address ) ,
Address : string ( address ) ,
IsLeader : false ,
IsLeader : false ,