/* Copyright NetFoundry Inc. Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at https://www.apache.org/licenses/LICENSE-2.0 Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License. */ package handler_ctrl import ( "github.com/google/uuid" "github.com/michaelquigley/pfxlog" "github.com/openziti/channel/v5" "github.com/openziti/ziti/v2/common/capabilities" "github.com/openziti/ziti/v2/common/ctrlchan" "github.com/openziti/ziti/v2/common/pb/ctrl_pb" "github.com/openziti/ziti/v2/controller/change" "github.com/openziti/ziti/v2/controller/network" "github.com/openziti/ziti/v2/controller/xctrl" "github.com/pkg/errors" "google.golang.org/protobuf/proto" ) type CtrlAccepter struct { network *network.Network xctrls []xctrl.Xctrl options *channel.Options heartbeatOptions *channel.HeartbeatOptions traceHandler *channel.TraceHandler } func NewCtrlAccepter(network *network.Network, xctrls []xctrl.Xctrl, options *channel.Options, heartbeatOptions *channel.HeartbeatOptions, traceHandler *channel.TraceHandler) *CtrlAccepter { return &CtrlAccepter{ network: network, xctrls: xctrls, options: options, heartbeatOptions: heartbeatOptions, traceHandler: traceHandler, } } // NewMultiListener returns a HelloAcceptor that handles both grouped (multi-underlay) and // ungrouped (single underlay) connections from routers. As a HelloAcceptor it defers the // hello acknowledgement until the group is registered, closing the race where a second // underlay for the same group could arrive before the group was known. func (self *CtrlAccepter) NewMultiListener() channel.HelloAcceptor { return channel.NewMultiListener(self.HandleGroupedUnderlay, self.AcceptUnderlay) } // HandleGroupedUnderlay handles incoming grouped connections from routers that support // multi-underlay control channels. It creates a MultiChannel with ListenerCtrlChannel. func (self *CtrlAccepter) HandleGroupedUnderlay(underlay channel.Underlay, closeCallback func()) (channel.Channel, error) { if _, hasSecret := underlay.Headers()[channel.GroupSecretHeader]; !hasSecret { underlay.Headers()[channel.GroupSecretHeader] = []byte(uuid.NewString()) } listenerCtrlChan := ctrlchan.NewListenerCtrlChannel() multiConfig := channel.Config{ LogicalName: "ctrl/" + underlay.Id(), Options: self.options, Underlay: underlay, Binder: channel.MakeBinder(channel.BindHandlerF(func(binding channel.Binding) error { binding.AddCloseHandler(channel.CloseHandlerF(func(ch channel.Channel) { closeCallback() })) return self.Bind(binding) })), Senders: listenerCtrlChan, MessageSourceProvider: listenerCtrlChan, UnderlayEventListeners: []channel.UnderlayEventListener{listenerCtrlChan}, // Multi-underlay-capable so the high/low-priority underlays are accepted; // MinTotalUnderlays closes the channel only when its last underlay is lost. Constraints: listenerCtrlChan.GetConstraints(), MinTotalUnderlays: 1, } mc, err := channel.NewChannel(&multiConfig) if err != nil { pfxlog.Logger().WithError(err).Errorf("failure accepting ctrl channel %v with multi-underlay", underlay.Label()) return nil, err } return mc, nil } // AcceptUnderlay handles incoming ungrouped connections from routers that don't support // multi-underlay control channels (backward compatibility). func (self *CtrlAccepter) AcceptUnderlay(underlay channel.Underlay) error { _, err := self.HandleGroupedUnderlay(underlay, func() {}) return err } func (self *CtrlAccepter) Bind(binding channel.Binding) error { binding.GetChannel().SetLogicalName(binding.GetChannel().Id()) ch := binding.GetChannel() log := pfxlog.Logger().WithField("routerId", ch.Id()) // A fresh instance per connection, carrying this channel: that is what lets connect and disconnect tell // two connections for one router apart, and keeps them from writing over each other's state. r, err := self.network.NewCtrlChanRouter(ch) if err != nil { return err } var ctrlChanListeners map[string][]string if ch.Underlay().Headers() != nil { if versionValue, found := ch.Underlay().Headers()[channel.HelloVersionHeader]; found { if versionInfo, err := self.network.VersionProvider.EncoderDecoder().Decode(versionValue); err == nil { r.VersionInfo = versionInfo log = log.WithField("version", r.VersionInfo.Version). WithField("revision", r.VersionInfo.Revision). WithField("buildDate", r.VersionInfo.BuildDate). WithField("os", r.VersionInfo.OS). WithField("arch", r.VersionInfo.Arch) } else { return errors.Wrap(err, "could not parse version info from router hello, not accepting router connection") } } else { return errors.New("no version info header, not accepting router connection") } r.SetLinkListeners(nil) headers := ch.Underlay().Headers() // Determine header locations based on router capabilities. 2.0+ routers // send a CapabilitiesHeader with RouterMultiChannel set and use header IDs // in the 1000+ range. Pre-2.0 routers use legacy IDs (10-12) and don't // send a CapabilitiesHeader. r.Capabilities = capabilities.GetCapabilities[capabilities.RouterCapability](headers) useNewHeaders := r.Capabilities.IsSet(capabilities.RouterMultiChannel) listenersHeaderId := ctrl_pb.LegacyListenersHeader if useNewHeaders { listenersHeaderId = int32(ctrl_pb.ControlHeaders_ListenersHeader) } if val, found := headers[listenersHeaderId]; found { listeners := &ctrl_pb.Listeners{} if err = proto.Unmarshal(val, listeners); err != nil { log.WithError(err).Error("unable to unmarshall listeners value") } else { r.SetLinkListeners(listeners.Listeners) for _, listener := range listeners.Listeners { log.WithField("address", listener.GetAddress()). WithField("protocol", listener.GetProtocol()). Debug("router listener") } } } else { log.Debug("no advertised listeners") } if val, found := ch.Underlay().Headers()[int32(ctrl_pb.ControlHeaders_CtrlChanListenersHeader)]; found { ctrlListeners := &ctrl_pb.CtrlChanListeners{} if err = proto.Unmarshal(val, ctrlListeners); err != nil { log.WithError(err).Error("unable to unmarshal ctrl chan listeners value") } else { ctrlChanListeners = make(map[string][]string, len(ctrlListeners.Listeners)) for _, listener := range ctrlListeners.Listeners { ctrlChanListeners[listener.Address] = listener.Groups } } } changeCtx := change.NewControlChannelChange(r.Id, r.Name, "router.connect", ch) self.network.Router.UpdateCtrlChanListeners(r, ctrlChanListeners, changeCtx) } else { return errors.New("channel provided no headers, not accepting router connection as version info not provided") } if err = newBindHandler(self.heartbeatOptions, r, self.network, self.xctrls).BindChannel(binding); err != nil { return errors.Wrap(err, "error binding router") } if self.traceHandler != nil { binding.AddPeekHandler(self.traceHandler) } if err = self.network.ConnectRouter(r); err != nil { if network.IsConnectRejected(err) { log.Info("router connect rejected; another connection is already current, router will redial") } // Returning the error fails the bind, so NewChannel closes this channel's underlay without // starting rx or registering it. That preserves the rx-gate for a rejected connection. return err } log.Info("accepted new router connection") return nil }