Files
ziti/router/handler_ctrl/update_ctrl_addresses.go
Paul Lorenz ae806045b5 Convert self-describing receive handlers to channel/v5. For #3983
- channel.TypedReceiveHandler -> channel.ContentTypeReceiver
- binding.AddTypedReceiveHandler(h) -> channel.AddReceiveHandlers(binding, h)

channel/v5 repurposes TypedReceiveHandler for the senders-typed handler and replaces the
self-describing pattern with ContentTypeReceiver plus the AddReceiveHandlers free function
(openziti/channel#262). Mechanical conversion; does not build on its own.
2026-06-18 12:51:02 -04:00

152 lines
4.2 KiB
Go

package handler_ctrl
import (
"sync/atomic"
"time"
"github.com/michaelquigley/pfxlog"
"github.com/openziti/channel/v5"
"github.com/openziti/ziti/v2/common/pb/ctrl_pb"
"github.com/openziti/ziti/v2/router/env"
"github.com/sirupsen/logrus"
"google.golang.org/protobuf/proto"
)
type updateCtrlAddressesHandler struct {
env env.RouterEnv
currentVersion atomic.Uint64
waitingForLeaderUpdate atomic.Bool
}
func (handler *updateCtrlAddressesHandler) NotifyIndexReset() {
handler.currentVersion.Store(0)
}
func (handler *updateCtrlAddressesHandler) ContentType() int32 {
return int32(ctrl_pb.ContentType_UpdateCtrlAddressesType)
}
func (handler *updateCtrlAddressesHandler) HandleReceive(msg *channel.Message, ch channel.Channel) {
log := pfxlog.ContextLogger(ch.Label()).Entry
upd := &ctrl_pb.UpdateCtrlAddresses{}
if err := proto.Unmarshal(msg.Body, upd); err != nil {
log.WithError(err).Error("error unmarshalling")
return
}
log = log.WithFields(logrus.Fields{
"endpoints": upd.Addresses,
"ctrlDetails": upd.Controllers,
"localVersion": handler.currentVersion.Load(),
"remoteVersion": upd.Index,
"isLeader": upd.IsLeader,
"ctrlId": ch.Id(),
})
log.Info("update ctrl endpoints message received")
if upd.IsLeader {
handler.waitingForLeaderUpdate.Store(false)
}
if handler.currentVersion.Load() == 0 || handler.currentVersion.Load() < upd.Index {
// Build address set from controllers
s := map[string]struct{}{}
for _, ctrl := range upd.Controllers {
for _, ep := range ctrl.Endpoints {
s[ep.Address] = struct{}{}
}
}
if len(s) == 0 {
log.Info("ctrl list is empty, ignoring and requesting latest set from leader")
go handler.requestCtrlListFromLeader()
return
}
ctrls := handler.env.GetNetworkControllers()
endpoints := ctrls.GetAll()
hasRemovals := false
for _, ctrl := range endpoints {
if _, ok := s[ctrl.Address()]; !ok {
hasRemovals = true
}
}
if !upd.IsLeader && hasRemovals {
log.Info("updated ctrl list is not from leader, using only additions and requesting latest set from leader")
if !handler.waitingForLeaderUpdate.Load() {
go handler.requestCtrlListFromLeader()
}
for _, ctrl := range endpoints {
s[ctrl.Address()] = struct{}{}
}
// filter controllers list to only include those whose addresses are in the final set
var filtered []*ctrl_pb.CtrlDetail
for _, ctrl := range upd.Controllers {
for _, ep := range ctrl.Endpoints {
if _, ok := s[ep.Address]; ok {
filtered = append(filtered, ctrl)
break
}
}
}
upd.Controllers = filtered
}
if len(upd.Controllers) > 0 {
log.Info("updating to newer controller endpoints")
handler.env.UpdateCtrlEndpointDetails(upd.Controllers)
handler.currentVersion.Store(upd.Index)
} else {
log.Info("ignoring empty controller endpoint list")
}
} else {
log.Info("ignoring outdated controller endpoint list")
}
if upd.IsLeader && (handler.currentVersion.Load() == 0 || handler.currentVersion.Load() <= upd.Index) {
log.Infof("updating current leader to %s", ch.Id())
handler.env.UpdateLeader(ch.Id())
}
}
func (handler *updateCtrlAddressesHandler) requestCtrlListFromLeader() {
if !handler.waitingForLeaderUpdate.CompareAndSwap(false, true) {
return
}
for handler.waitingForLeaderUpdate.Load() {
log := pfxlog.Logger().Entry
leader := handler.env.GetNetworkControllers().GetLeader()
if leader == nil {
log.Info("no leader, unable to request latest cluster members from leader")
} else {
log = log.WithField("ctrlId", leader.Channel().Id())
if !leader.IsConnected() {
log.Info("not connected to leader, unable to request latest cluster members from leader")
} else {
msg := channel.NewMessage(int32(ctrl_pb.ContentType_RequestClusterMembers), nil)
if err := leader.Channel().Send(msg); err != nil {
log.WithError(err).Error("error sending request for latest cluster members to leader")
}
}
}
time.Sleep(30 * time.Second)
}
}
func newUpdateCtrlAddressesHandler(env env.RouterEnv) channel.ContentTypeReceiver {
result := &updateCtrlAddressesHandler{
env: env,
}
env.GetIndexWatchers().AddIndexResetWatcher(result.NotifyIndexReset)
return result
}