mirror of
https://github.com/openziti/ziti.git
synced 2026-09-10 00:35:41 +00:00
ae806045b5
- 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.
296 lines
10 KiB
Go
296 lines
10 KiB
Go
package router
|
|
|
|
import (
|
|
"bufio"
|
|
"bytes"
|
|
"fmt"
|
|
"net"
|
|
"os"
|
|
"time"
|
|
|
|
"github.com/michaelquigley/pfxlog"
|
|
"github.com/openziti/channel/v5"
|
|
"github.com/openziti/ziti/v2/common/agent"
|
|
"github.com/openziti/ziti/v2/common/handler_common"
|
|
"github.com/openziti/ziti/v2/common/pb/ctrl_pb"
|
|
"github.com/openziti/ziti/v2/common/pb/mgmt_pb"
|
|
"github.com/pkg/errors"
|
|
"github.com/sirupsen/logrus"
|
|
"google.golang.org/protobuf/proto"
|
|
)
|
|
|
|
const (
|
|
AgentAppId byte = 2
|
|
DumpApiSessions byte = 128
|
|
)
|
|
|
|
func (self *Router) RegisterAgentBindHandler(bindHandler channel.BindHandler) {
|
|
self.agentBindHandlers = append(self.agentBindHandlers, bindHandler)
|
|
}
|
|
|
|
func (self *Router) RegisterDefaultAgentOps(debugEnabled bool) {
|
|
self.agentBindHandlers = append(self.agentBindHandlers, channel.BindHandlerF(func(binding channel.Binding) error {
|
|
channel.AddReceiveHandlers(binding, self.inspectHandler)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDebugDumpForwarderTablesRequestType), self.agentOpDumpForwarderTables)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDebugDumpLinksRequestType), self.agentOpsDumpLinks)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterQuiesceRequestType), self.agentOpQuiesceRouter)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDequiesceRequestType), self.agentOpDequiesceRouter)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDecommissionRequestType), self.agentOpDecommissionRouter)
|
|
|
|
if debugEnabled {
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDebugUpdateRouteRequestType), self.agentOpUpdateRoute)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDebugUnrouteRequestType), self.agentOpUnroute)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDebugForgetLinkRequestType), self.agentOpForgetLink)
|
|
binding.AddReceiveHandlerF(int32(mgmt_pb.ContentType_RouterDebugToggleCtrlChannelRequestType), self.agentOpToggleCtrlChan)
|
|
}
|
|
return nil
|
|
}))
|
|
|
|
if debugEnabled {
|
|
self.RegisterAgentOp(DumpApiSessions, self.GetStateManager().DumpApiSessions)
|
|
}
|
|
}
|
|
|
|
func (self *Router) RegisterAgentOp(opId byte, f func(c *bufio.ReadWriter) error) {
|
|
self.debugOperations[opId] = f
|
|
}
|
|
|
|
func (self *Router) bindAgentChannel(binding channel.Binding) error {
|
|
for _, bh := range self.agentBindHandlers {
|
|
if err := bh.BindChannel(binding); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (self *Router) HandleAgentAsyncOp(conn net.Conn) error {
|
|
return agent.HandleChannelConnection(conn, self.config.Id, AgentAppId, channel.BindHandlerF(self.bindAgentChannel))
|
|
}
|
|
|
|
func (self *Router) agentOpForgetLink(m *channel.Message, ch channel.Channel) {
|
|
log := pfxlog.Logger()
|
|
linkId := string(m.Body)
|
|
var found bool
|
|
if link, _ := self.xlinkRegistry.GetLinkById(linkId); link != nil {
|
|
self.xlinkRegistry.DebugForgetLink(linkId)
|
|
self.forwarder.UnregisterLink(link)
|
|
found = true
|
|
}
|
|
|
|
log.Infof("forget of link %v was requested. link found? %v", linkId, found)
|
|
result := fmt.Sprintf("link removed: %v", found)
|
|
handler_common.SendOpResult(m, ch, "link.remove", result, true)
|
|
}
|
|
|
|
func (self *Router) agentOpToggleCtrlChan(m *channel.Message, ch channel.Channel) {
|
|
ctrlId := string(m.Body)
|
|
|
|
results := &bytes.Buffer{}
|
|
toggleOn, _ := m.GetBoolHeader(int32(mgmt_pb.Header_CtrlChanToggle))
|
|
|
|
success := true
|
|
count := 0
|
|
self.ctrls.ForEach(func(controllerId string, ch channel.Channel) {
|
|
if ctrlId == "" || controllerId == ctrlId {
|
|
log := pfxlog.Logger().WithField("ctrlId", controllerId)
|
|
if toggleable, ok := ch.Underlay().(connectionToggle); ok {
|
|
if toggleOn {
|
|
if err := toggleable.Reconnect(); err != nil {
|
|
log.WithError(err).Error("control channel: failed to reconnect")
|
|
_, _ = fmt.Fprintf(results, "control channel: failed to reconnect (%v)\n", err)
|
|
success = false
|
|
} else {
|
|
log.Warn("control channel: reconnected")
|
|
_, _ = fmt.Fprint(results, "control channel: reconnected")
|
|
count++
|
|
}
|
|
} else {
|
|
if err := toggleable.Disconnect(); err != nil {
|
|
log.WithError(err).Error("control channel: failed to close")
|
|
_, _ = fmt.Fprintf(results, "control channel: failed to close (%v)\n", err)
|
|
success = false
|
|
} else {
|
|
log.Warn("control channel: closed")
|
|
_, _ = fmt.Fprint(results, "control channel: closed")
|
|
count++
|
|
}
|
|
}
|
|
} else {
|
|
log.Warn("control channel: not toggleable")
|
|
_, _ = fmt.Fprint(results, "control channel: not toggleable")
|
|
success = false
|
|
}
|
|
}
|
|
})
|
|
|
|
if count == 0 {
|
|
_, _ = fmt.Fprintf(results, "control channel: no controllers matched id [%v]", ctrlId)
|
|
success = false
|
|
}
|
|
|
|
handler_common.SendOpResult(m, ch, "ctrl.toggle", results.String(), success)
|
|
}
|
|
|
|
func (self *Router) agentOpQuiesceRouter(m *channel.Message, ch channel.Channel) {
|
|
ctrlCh := self.ctrls.AnyValidCtrlChannel()
|
|
if ctrlCh == nil {
|
|
handler_common.SendOpResult(m, ch, "quiesce", "unable to reach controller", false)
|
|
return
|
|
}
|
|
|
|
msg := channel.NewMessage(int32(ctrl_pb.ContentType_QuiesceRouterRequestType), nil)
|
|
resp, err := msg.WithTimeout(5 * time.Second).SendForReply(ctrlCh)
|
|
if err != nil {
|
|
handler_common.SendOpResult(m, ch, "quiesce", fmt.Sprintf("error in controller communications: %v", err.Error()), false)
|
|
return
|
|
}
|
|
|
|
result := channel.UnmarshalResult(resp)
|
|
handler_common.SendOpResult(m, ch, "quiesce", result.Message+"\n", result.Success)
|
|
}
|
|
|
|
func (self *Router) agentOpDequiesceRouter(m *channel.Message, ch channel.Channel) {
|
|
ctrlCh := self.ctrls.AnyValidCtrlChannel()
|
|
if ctrlCh == nil {
|
|
handler_common.SendOpResult(m, ch, "dequiesce", "unable to reach controller", false)
|
|
return
|
|
}
|
|
|
|
msg := channel.NewMessage(int32(ctrl_pb.ContentType_DequiesceRouterRequestType), nil)
|
|
resp, err := msg.WithTimeout(5 * time.Second).SendForReply(ctrlCh)
|
|
if err != nil {
|
|
handler_common.SendOpResult(m, ch, "dequiesce", fmt.Sprintf("error in controller communications: %v", err.Error()), false)
|
|
return
|
|
}
|
|
|
|
result := channel.UnmarshalResult(resp)
|
|
handler_common.SendOpResult(m, ch, "dequiesce", result.Message+"\n", result.Success)
|
|
}
|
|
|
|
func (self *Router) agentOpDecommissionRouter(m *channel.Message, ch channel.Channel) {
|
|
ctrlCh := self.ctrls.AnyValidCtrlChannel()
|
|
if ctrlCh == nil {
|
|
handler_common.SendOpResult(m, ch, "decomission", "unable to reach controller", false)
|
|
return
|
|
}
|
|
|
|
msg := channel.NewMessage(int32(ctrl_pb.ContentType_DecommissionRouterRequestType), nil)
|
|
err := msg.WithTimeout(5 * time.Second).SendAndWaitForWire(ctrlCh)
|
|
if err != nil {
|
|
handler_common.SendOpResult(m, ch, "decommission", "unexpected result, no disconnect but no error", false)
|
|
return
|
|
}
|
|
|
|
connectable, ok := ctrlCh.Underlay().(interface{ IsConnected() bool })
|
|
if !ok {
|
|
pfxlog.Logger().Warn("control channel can't be checked to see if it's connected")
|
|
}
|
|
|
|
for i := 0; i < 100; i++ {
|
|
// decommissioning router should cause control channel to close
|
|
if connectable == nil || !connectable.IsConnected() {
|
|
handler_common.SendOpResult(m, ch, "decommission", "router decommissioned\n", true)
|
|
time.Sleep(500 * time.Millisecond)
|
|
pfxlog.Logger().Info("router decommissioned, shutting down")
|
|
os.Exit(0)
|
|
}
|
|
time.Sleep(10 * time.Millisecond)
|
|
}
|
|
|
|
pfxlog.Logger().Warn("control channel still connected after decomission")
|
|
handler_common.SendOpResult(m, ch, "decommission", "controller didn't disconnect after decommission", false)
|
|
}
|
|
|
|
func (self *Router) agentOpDumpForwarderTables(m *channel.Message, ch channel.Channel) {
|
|
tables := self.forwarder.Debug()
|
|
handler_common.SendOpResult(m, ch, "dump.forwarder_tables", tables, true)
|
|
}
|
|
|
|
func (self *Router) agentOpsDumpLinks(m *channel.Message, ch channel.Channel) {
|
|
result := &bytes.Buffer{}
|
|
for link := range self.xlinkRegistry.Iter() {
|
|
line := fmt.Sprintf("id: %v dest: %v protocol: %v\n", link.Id(), link.DestinationId(), link.LinkProtocol())
|
|
_, err := result.WriteString(line)
|
|
if err != nil {
|
|
handler_common.SendOpResult(m, ch, "dump.links", err.Error(), false)
|
|
return
|
|
}
|
|
}
|
|
|
|
output := result.String()
|
|
if len(output) == 0 {
|
|
output = "no links\n"
|
|
}
|
|
handler_common.SendOpResult(m, ch, "dump.links", output, true)
|
|
}
|
|
|
|
func (self *Router) agentOpUpdateRoute(m *channel.Message, ch channel.Channel) {
|
|
logrus.Warn("received debug operation to update routes")
|
|
ctrlId, _ := m.GetStringHeader(int32(mgmt_pb.Header_ControllerId))
|
|
if ctrlId == "" {
|
|
handler_common.SendOpResult(m, ch, "update.route", "no controller id provided", false)
|
|
return
|
|
}
|
|
|
|
ctrl := self.ctrls.GetChannel(ctrlId)
|
|
if ctrl == nil {
|
|
handler_common.SendOpResult(m, ch, "update.route", fmt.Sprintf("no control channel found for [%v]", ctrlId), false)
|
|
return
|
|
}
|
|
|
|
route := &ctrl_pb.Route{}
|
|
if err := proto.Unmarshal(m.Body, route); err != nil {
|
|
handler_common.SendOpResult(m, ch, "update.route", err.Error(), false)
|
|
return
|
|
}
|
|
|
|
if err := self.forwarder.Route(ctrlId, route); err != nil {
|
|
handler_common.SendOpResult(m, ch, "update.route", errors.Wrap(err, "error adding route").Error(), false)
|
|
return
|
|
}
|
|
|
|
logrus.Warnf("route added: %+v", route)
|
|
handler_common.SendOpResult(m, ch, "update.route", "route added", true)
|
|
}
|
|
|
|
func (self *Router) agentOpUnroute(m *channel.Message, ch channel.Channel) {
|
|
logrus.Warn("received debug operation to unroute a circuit")
|
|
|
|
unroute := &ctrl_pb.Unroute{}
|
|
if err := proto.Unmarshal(m.Body, unroute); err != nil {
|
|
handler_common.SendOpResult(m, ch, "unroute", err.Error(), false)
|
|
return
|
|
}
|
|
|
|
self.forwarder.Unroute(unroute.CircuitId, unroute.Now)
|
|
logrus.Warnf("circuit unrouted: %+v", unroute)
|
|
handler_common.SendOpResult(m, ch, "unroute", fmt.Sprintf("circuit unrouted [%s]", unroute.CircuitId), true)
|
|
}
|
|
|
|
func (self *Router) HandleAgentOp(conn net.Conn) error {
|
|
bconn := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
|
|
appId, err := bconn.ReadByte()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if appId != AgentAppId {
|
|
return errors.Errorf("invalid operation for router")
|
|
}
|
|
|
|
op, err := bconn.ReadByte()
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if opF, ok := self.debugOperations[op]; ok {
|
|
if err := opF(bconn); err != nil {
|
|
return err
|
|
}
|
|
return bconn.Flush()
|
|
}
|
|
return errors.Errorf("invalid operation %v", op)
|
|
}
|