From a3f31f6e20e65c0caac6eed0a55758717a6795da Mon Sep 17 00:00:00 2001 From: Paul Lorenz Date: Mon, 8 Mar 2021 11:40:47 -0500 Subject: [PATCH 1/8] Try and fix build after push from update-dependencies run --- .github/workflows/main.yml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index e13380fd6..101d20a23 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -13,6 +13,8 @@ jobs: steps: - name: Git Checkout uses: actions/checkout@v2 + with: + persist-credentials: false - name: Install Go uses: actions/setup-go@v2 From 931d1e4794550fefaabc182f53cbfa8fdd8c9932 Mon Sep 17 00:00:00 2001 From: Paul Lorenz Date: Mon, 8 Mar 2021 12:37:37 -0500 Subject: [PATCH 2/8] all manual trigger of build --- .github/workflows/main.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 101d20a23..5cab1d2ba 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -6,6 +6,7 @@ on: - main - release-* pull_request: + workflow_dispatch: jobs: build: From 5c22c9e531d88bc90ad4484ec13a1a222d9a8e7a Mon Sep 17 00:00:00 2001 From: Michael Quigley Date: Tue, 9 Mar 2021 11:35:51 -0500 Subject: [PATCH 3/8] Initial parallel network re-routing, and session-level concurrency control. (#206) --- controller/network/network.go | 106 ++++++++++++++++++++-------------- controller/network/session.go | 2 + 2 files changed, 66 insertions(+), 42 deletions(-) diff --git a/controller/network/network.go b/controller/network/network.go index a8bf29594..59edcf9d1 100644 --- a/controller/network/network.go +++ b/controller/network/network.go @@ -635,25 +635,21 @@ func (network *Network) AddRouterPresenceHandler(h RouterPresenceHandler) { } func (network *Network) Run() { - log := pfxlog.Logger() - defer log.Error("exited") - log.Info("started") + defer logrus.Error("exited") + logrus.Info("started") for { select { case r := <-network.routerChanged: - log.Infof("changed router [r/%s]", r.Id) + logrus.Infof("changed router [r/%s]", r.Id) network.assemble() network.clean() case l := <-network.linkChanged: - log.Infof("changed link [l/%s]", l.Id.Token) - if err := network.rerouteLink(l); err != nil { - log.Errorf("unexpected error rerouting link (%s)", err) - } + go network.handleLinkChanged(l) case ffr := <-network.forwardingFaults: - network.fault(ffr) + go network.handleForwardingFaults(ffr) network.clean() case <-time.After(time.Duration(network.options.CycleSeconds) * time.Second): @@ -669,6 +665,17 @@ func (network *Network) Run() { } } +func (network *Network) handleLinkChanged(l *Link) { + logrus.Infof("changed link [l/%s]", l.Id.Token) + if err := network.rerouteLink(l); err != nil { + logrus.Errorf("unexpected error rerouting link (%s)", err) + } +} + +func (network *Network) handleForwardingFaults(ffr *ForwardingFaultReport) { + network.fault(ffr) +} + func (network *Network) AddCapability(capability string) { network.lock.Lock() defer network.lock.Unlock() @@ -682,17 +689,16 @@ func (network *Network) GetCapabilities() []string { } func (network *Network) rerouteLink(l *Link) error { - log := pfxlog.Logger() - log.Infof("link [l/%s] changed", l.Id.Token) + logrus.Infof("link [l/%s] changed", l.Id.Token) sessions := network.sessionController.all() for _, s := range sessions { if s.Circuit.usesLink(l) { - log.Infof("session [s/%s] uses link [l/%s]", s.Id.Token, l.Id.Token) + logrus.Infof("session [s/%s] uses link [l/%s]", s.Id.Token, l.Id.Token) if err := network.rerouteSession(s); err != nil { - log.Errorf("error rerouting session [s/%s], removing", s.Id.Token) + logrus.Errorf("error rerouting session [s/%s], removing", s.Id.Token) if err := network.RemoveSession(s.Id, true); err != nil { - log.Errorf("error removing session [s/%s] (%s)", s.Id.Token, err) + logrus.Errorf("error removing session [s/%s] (%s)", s.Id.Token, err) } } } @@ -703,9 +709,47 @@ func (network *Network) rerouteLink(l *Link) error { func (network *Network) rerouteSession(s *Session) error { log := pfxlog.Logger() - log.Warnf("rerouting [s/%s]", s.Id.Token) - if cq, err := network.UpdateCircuit(s.Circuit); err == nil { + if s.Rerouting.CompareAndSwap(false, true) { + defer s.Rerouting.Set(false) + + log.Warnf("rerouting [s/%s]", s.Id.Token) + + if cq, err := network.UpdateCircuit(s.Circuit); err == nil { + s.Circuit = cq + + rms, err := cq.CreateRouteMessages(SmartRerouteAttempt, s.Id, s.Terminator.GetAddress()) + if err != nil { + log.Errorf("error creating route messages (%s)", err) + return err + } + + for i := 0; i < len(cq.Path); i++ { + if _, err := sendRoute(cq.Path[i], rms[i], network.options.RouteTimeout); err != nil { + log.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err) + } + } + + log.Infof("rerouted session [s/%s]", s.Id.Token) + + network.CircuitUpdated(s.Id, s.Circuit) + + return nil + } else { + return err + } + } else { + log.Warnf("already re-routing [s/%s], ignoring", s.Id.Token) + return nil + } +} + +func (network *Network) smartReroute(s *Session, cq *Circuit) error { + log := pfxlog.Logger() + + if s.Rerouting.CompareAndSwap(false, true) { + defer s.Rerouting.Set(false) + s.Circuit = cq rms, err := cq.CreateRouteMessages(SmartRerouteAttempt, s.Id, s.Terminator.GetAddress()) @@ -720,40 +764,18 @@ func (network *Network) rerouteSession(s *Session) error { } } - log.Infof("rerouted session [s/%s]", s.Id.Token) + log.Debugf("rerouted session [s/%s]", s.Id.Token) network.CircuitUpdated(s.Id, s.Circuit) return nil + } else { - return err + log.Warnf("elready re-routing [s/%s], ignoring", s.Id.Token) + return nil } } -func (network *Network) smartReroute(s *Session, cq *Circuit) error { - log := pfxlog.Logger() - - s.Circuit = cq - - rms, err := cq.CreateRouteMessages(SmartRerouteAttempt, s.Id, s.Terminator.GetAddress()) - if err != nil { - log.Errorf("error creating route messages (%s)", err) - return err - } - - for i := 0; i < len(cq.Path); i++ { - if _, err := sendRoute(cq.Path[i], rms[i], network.options.RouteTimeout); err != nil { - log.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err) - } - } - - log.Debugf("rerouted session [s/%s]", s.Id.Token) - - network.CircuitUpdated(s.Id, s.Circuit) - - return nil -} - func (network *Network) AcceptMetrics(metrics *metrics_pb.MetricsMessage) { if metrics.SourceId == network.nodeId.Token { return // ignore metrics coming from the controller itself diff --git a/controller/network/session.go b/controller/network/session.go index 97f1db8dc..38467fd3e 100644 --- a/controller/network/session.go +++ b/controller/network/session.go @@ -19,6 +19,7 @@ package network import ( "github.com/openziti/fabric/controller/xt" "github.com/openziti/foundation/identity/identity" + "github.com/openziti/foundation/util/concurrenz" "github.com/orcaman/concurrent-map" ) @@ -28,6 +29,7 @@ type Session struct { Service *Service Terminator xt.Terminator Circuit *Circuit + Rerouting concurrenz.AtomicBoolean PeerData xt.PeerData } From 8d3ea033063d653148472f00486248eac1f76afe Mon Sep 17 00:00:00 2001 From: Paul Lorenz Date: Tue, 9 Mar 2021 12:10:04 -0500 Subject: [PATCH 4/8] Use bitset for xgress flags --- router/xgress/link_send_buffer.go | 2 +- router/xgress/retransmitter.go | 2 +- router/xgress/xgress.go | 50 +++++++++++-------------------- 3 files changed, 19 insertions(+), 35 deletions(-) diff --git a/router/xgress/link_send_buffer.go b/router/xgress/link_send_buffer.go index 0e96ac69c..ee7197bb1 100644 --- a/router/xgress/link_send_buffer.go +++ b/router/xgress/link_send_buffer.go @@ -207,7 +207,7 @@ func (buffer *LinkSendBuffer) run() { case ack := <-buffer.newlyReceivedAcks: buffer.receiveAcknowledgement(ack) buffer.retransmit() - if buffer.closeWhenEmpty.Get() && len(buffer.buffer) == 0 && !buffer.x.closed.Get() && buffer.x.IsEndOfSessionSent() { + if buffer.closeWhenEmpty.Get() && len(buffer.buffer) == 0 && !buffer.x.Closed() && buffer.x.IsEndOfSessionSent() { go buffer.x.Close() } diff --git a/router/xgress/retransmitter.go b/router/xgress/retransmitter.go index 5cafa8052..2d1a0e6a3 100644 --- a/router/xgress/retransmitter.go +++ b/router/xgress/retransmitter.go @@ -147,7 +147,7 @@ func (retransmitter *Retransmitter) retransmitSender() { if !retransmit.isAcked() { if err := retransmitter.forwarder.ForwardPayload(retransmit.x.address, retransmit.payload); err != nil { // if xgress is closed, don't log the error. We still want to try retransmitting in case we're re-sending end of session - if !retransmit.x.closed.Get() { + if !retransmit.x.Closed() { logger.WithError(err).Errorf("unexpected error while retransmitting payload from [@/%v]", retransmit.x.address) retransmissionFailures.Mark(1) retransmitter.faultReporter.ReportForwardingFault(retransmit.payload.SessionId) diff --git a/router/xgress/xgress.go b/router/xgress/xgress.go index f60541eb9..dac879069 100644 --- a/router/xgress/xgress.go +++ b/router/xgress/xgress.go @@ -34,6 +34,11 @@ import ( const ( HeaderKeyUUID = 0 + + closedFlag = 0 + rxerStartedFlag = 1 + endOfSessionRecvdFlag = 2 + endOfSessionSentFlag = 3 ) type Address string @@ -107,11 +112,7 @@ type Xgress struct { linkRxBuffer *LinkReceiveBuffer closeHandler CloseHandler peekHandlers []PeekHandler - closed concurrenz.AtomicBoolean - rxerStarted bool - endOfSessionRecvd bool - endOfSessionSent bool - flagsMutex sync.Mutex + flags concurrenz.AtomicBitSet timeOfLastRxFromLink int64 } @@ -165,31 +166,19 @@ func (self *Xgress) AddPeekHandler(peekHandler PeekHandler) { } func (self *Xgress) IsEndOfSessionReceived() bool { - self.flagsMutex.Lock() - defer self.flagsMutex.Unlock() - return self.endOfSessionRecvd + return self.flags.IsSet(endOfSessionRecvdFlag) } func (self *Xgress) markSessionEndReceived() { - self.flagsMutex.Lock() - defer self.flagsMutex.Unlock() - self.endOfSessionRecvd = true + self.flags.Set(endOfSessionRecvdFlag, true) } func (self *Xgress) IsSessionStarted() bool { - self.flagsMutex.Lock() - defer self.flagsMutex.Unlock() - return !self.IsTerminator() || self.rxerStarted + return !self.IsTerminator() || self.flags.IsSet(rxerStartedFlag) } func (self *Xgress) firstSessionStartReceived() bool { - self.flagsMutex.Lock() - defer self.flagsMutex.Unlock() - if self.rxerStarted { - return false - } - self.rxerStarted = true - return true + return self.flags.CompareAndSet(rxerStartedFlag, false, true) } func (self *Xgress) Start() { @@ -243,20 +232,15 @@ func (self *Xgress) GetEndSession() *Payload { } func (self *Xgress) ForwardEndOfSession(sendF func(payload *Payload) bool) { - self.flagsMutex.Lock() - defer self.flagsMutex.Unlock() - // for now always send end of session. too many is better than not enough - if !self.endOfSessionRecvd { + if !self.IsEndOfSessionSent() { sendF(self.GetEndSession()) - self.endOfSessionSent = true + self.flags.Set(endOfSessionSentFlag, true) } } func (self *Xgress) IsEndOfSessionSent() bool { - self.flagsMutex.Lock() - defer self.flagsMutex.Unlock() - return self.endOfSessionSent + return self.flags.IsSet(endOfSessionSentFlag) } func (self *Xgress) CloseTimeout(duration time.Duration) { @@ -277,7 +261,7 @@ Things which can trigger close func (self *Xgress) Close() { log := pfxlog.ContextLogger(self.Label()) - if self.closed.CompareAndSwap(false, true) { + if self.flags.CompareAndSet(closedFlag, false, true) { log.Debug("closing xgress peer") if err := self.peer.Close(); err != nil { log.WithError(err).Warn("error while closing xgress peer") @@ -301,11 +285,11 @@ func (self *Xgress) Close() { } func (self *Xgress) Closed() bool { - return self.closed.Get() + return self.flags.IsSet(closedFlag) } func (self *Xgress) SendPayload(payload *Payload) error { - if self.closed.Get() { + if self.Closed() { return nil } @@ -479,7 +463,7 @@ func (self *Xgress) rx() { return } - if self.closed.Get() { + if self.Closed() { return } start := 0 From f4d9be42715075896137ae25cd487c2d7c9bc165 Mon Sep 17 00:00:00 2001 From: Michael Quigley Date: Tue, 9 Mar 2021 12:42:26 -0500 Subject: [PATCH 5/8] Log lint. (#206) --- controller/network/network.go | 22 +++++++++------------- 1 file changed, 9 insertions(+), 13 deletions(-) diff --git a/controller/network/network.go b/controller/network/network.go index 59edcf9d1..38cd08a59 100644 --- a/controller/network/network.go +++ b/controller/network/network.go @@ -708,29 +708,27 @@ func (network *Network) rerouteLink(l *Link) error { } func (network *Network) rerouteSession(s *Session) error { - log := pfxlog.Logger() - if s.Rerouting.CompareAndSwap(false, true) { defer s.Rerouting.Set(false) - log.Warnf("rerouting [s/%s]", s.Id.Token) + logrus.Warnf("rerouting [s/%s]", s.Id.Token) if cq, err := network.UpdateCircuit(s.Circuit); err == nil { s.Circuit = cq rms, err := cq.CreateRouteMessages(SmartRerouteAttempt, s.Id, s.Terminator.GetAddress()) if err != nil { - log.Errorf("error creating route messages (%s)", err) + logrus.Errorf("error creating route messages (%s)", err) return err } for i := 0; i < len(cq.Path); i++ { if _, err := sendRoute(cq.Path[i], rms[i], network.options.RouteTimeout); err != nil { - log.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err) + logrus.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err) } } - log.Infof("rerouted session [s/%s]", s.Id.Token) + logrus.Infof("rerouted session [s/%s]", s.Id.Token) network.CircuitUpdated(s.Id, s.Circuit) @@ -739,14 +737,12 @@ func (network *Network) rerouteSession(s *Session) error { return err } } else { - log.Warnf("already re-routing [s/%s], ignoring", s.Id.Token) + logrus.Warnf("not rerouting [s/%s], already in progress", s.Id.Token) return nil } } func (network *Network) smartReroute(s *Session, cq *Circuit) error { - log := pfxlog.Logger() - if s.Rerouting.CompareAndSwap(false, true) { defer s.Rerouting.Set(false) @@ -754,24 +750,24 @@ func (network *Network) smartReroute(s *Session, cq *Circuit) error { rms, err := cq.CreateRouteMessages(SmartRerouteAttempt, s.Id, s.Terminator.GetAddress()) if err != nil { - log.Errorf("error creating route messages (%s)", err) + logrus.Errorf("error creating route messages (%s)", err) return err } for i := 0; i < len(cq.Path); i++ { if _, err := sendRoute(cq.Path[i], rms[i], network.options.RouteTimeout); err != nil { - log.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err) + logrus.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err) } } - log.Debugf("rerouted session [s/%s]", s.Id.Token) + logrus.Debugf("rerouted session [s/%s]", s.Id.Token) network.CircuitUpdated(s.Id, s.Circuit) return nil } else { - log.Warnf("elready re-routing [s/%s], ignoring", s.Id.Token) + logrus.Warnf("not rerouting [s/%s], already in progress", s.Id.Token) return nil } } From 3ec45972b90d417d9c4e03a1f9ed8f073f555eeb Mon Sep 17 00:00:00 2001 From: Michael Quigley Date: Tue, 9 Mar 2021 12:52:41 -0500 Subject: [PATCH 6/8] Warn->Info (#206) --- controller/network/network.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/controller/network/network.go b/controller/network/network.go index 38cd08a59..ffa00d291 100644 --- a/controller/network/network.go +++ b/controller/network/network.go @@ -737,7 +737,7 @@ func (network *Network) rerouteSession(s *Session) error { return err } } else { - logrus.Warnf("not rerouting [s/%s], already in progress", s.Id.Token) + logrus.Infof("not rerouting [s/%s], already in progress", s.Id.Token) return nil } } @@ -767,7 +767,7 @@ func (network *Network) smartReroute(s *Session, cq *Circuit) error { return nil } else { - logrus.Warnf("not rerouting [s/%s], already in progress", s.Id.Token) + logrus.Infof("not rerouting [s/%s], already in progress", s.Id.Token) return nil } } From ece7b4546eada832ec28466f7416f127e4e25db2 Mon Sep 17 00:00:00 2001 From: Paul Lorenz Date: Tue, 9 Mar 2021 13:59:34 -0500 Subject: [PATCH 7/8] Allow configuring metrics report interval and queue depth --- router/config.go | 24 +++++++++++++++++++++++- router/router.go | 19 ++++++++++--------- router/xgress/ordering_test.go | 2 +- 3 files changed, 34 insertions(+), 11 deletions(-) diff --git a/router/config.go b/router/config.go index 3aa6adc34..702296a08 100644 --- a/router/config.go +++ b/router/config.go @@ -85,7 +85,11 @@ type Config struct { Dialers map[string]xgress.OptionsData Listeners []listenerBinding Transport map[interface{}]interface{} - src map[interface{}]interface{} + Metrics struct { + ReportInterval time.Duration + MessageQueueSize int + } + src map[interface{}]interface{} } func (config *Config) Configure(sub config.Subconfig) error { @@ -309,6 +313,24 @@ func LoadConfig(path string) (*Config, error) { } } + cfg.Metrics.ReportInterval = 15 * time.Second + cfg.Metrics.MessageQueueSize = 10 + if value, found := cfgmap["metrics"]; found { + if submap, ok := value.(map[interface{}]interface{}); ok { + if value, found := submap["reportInterval"]; found { + if cfg.Metrics.ReportInterval, err = time.ParseDuration(value.(string)); err != nil { + return nil, errors.Wrap(err, "invalid value for metrics.reportInterval") + } + } + if value, found := submap["messageQueueSize"]; found { + if intVal, ok := value.(int); ok { + cfg.Metrics.MessageQueueSize = intVal + } else { + return nil, errors.Wrap(err, "invalid value for metrics.messageQueueSize") + } + } + } + } return cfg, nil } diff --git a/router/router.go b/router/router.go index 1a2b85808..7eb93aa72 100644 --- a/router/router.go +++ b/router/router.go @@ -89,12 +89,7 @@ func Create(config *Config, versionProvider common.VersionProvider) *Router { closeNotify := make(chan struct{}) eventDispatcher := event.NewDispatcher(closeNotify) - metricsConfig := &metrics.Config{ - Source: config.Id.Token, - ReportInterval: time.Second * 15, - EventSink: metrics.NilHandler{}, - } - metricsRegistry := metrics.NewUsageRegistryFromConfig(metricsConfig, closeNotify) + metricsRegistry := metrics.NewUsageRegistry(config.Id.Token, map[string]string{}, closeNotify) xgress.InitMetrics(metricsRegistry) faulter := forwarder.NewFaulter(config.Forwarder.FaultTxInterval, closeNotify) @@ -190,8 +185,14 @@ func (self *Router) Run() error { } func (self *Router) showOptions() { - if ctrl, err := json.Marshal(self.config.Ctrl.Options); err == nil { - pfxlog.Logger().Infof("ctrl = %s", string(ctrl)) + if output, err := json.Marshal(self.config.Ctrl.Options); err == nil { + pfxlog.Logger().Infof("ctrl = %s", string(output)) + } else { + logrus.Fatalf("unable to display options (%v)", err) + } + + if output, err := json.Marshal(self.config.Metrics); err == nil { + pfxlog.Logger().Infof("metrics = %s", string(output)) } else { logrus.Fatalf("unable to display options (%v)", err) } @@ -338,7 +339,7 @@ func (self *Router) startControlPlane() error { } self.metricsReporter = metrics.NewChannelReporter(self.ctrl) - self.metricsRegistry.SetEventSink(self.metricsReporter) + self.metricsRegistry.StartReporting(self.metricsReporter, self.config.Metrics.ReportInterval, self.config.Metrics.MessageQueueSize) return nil } diff --git a/router/xgress/ordering_test.go b/router/xgress/ordering_test.go index 1f02df7ac..253b81aee 100644 --- a/router/xgress/ordering_test.go +++ b/router/xgress/ordering_test.go @@ -56,8 +56,8 @@ func (n noopReceiveHandler) HandleXgressReceive(payload *Payload, x *Xgress) { } func Test_Ordering(t *testing.T) { - metricsRegistry := metrics.NewUsageRegistry("test", map[string]string{}, time.Minute, noopMetricsHandler{}, nil) closeNotify := make(chan struct{}) + metricsRegistry := metrics.NewUsageRegistry("test", map[string]string{}, closeNotify) InitPayloadIngester(closeNotify) InitMetrics(metricsRegistry) InitAcker(&noopForwarder{}, metricsRegistry, closeNotify) From 8c14bfd607ffac1790a7f45c7ee771312df4f06b Mon Sep 17 00:00:00 2001 From: Michael Quigley Date: Tue, 9 Mar 2021 16:56:18 -0500 Subject: [PATCH 8/8] Fix for xgressDialQueueLength option in forwarder.options. --- router/forwarder/options.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/router/forwarder/options.go b/router/forwarder/options.go index 68f7eb42e..3c4868735 100644 --- a/router/forwarder/options.go +++ b/router/forwarder/options.go @@ -87,14 +87,14 @@ func LoadOptions(src map[interface{}]interface{}) (*Options, error) { } } - if value, found := src["xgressDialQueueLength "]; found { + if value, found := src["xgressDialQueueLength"]; found { if length, ok := value.(int); ok { if length <= 0 || length > 10000 { - return nil, errors.New("invalid value for 'xgressDialQueueLength ', expected integer between 1 and 1000") + return nil, errors.New("invalid value for 'xgressDialQueueLength', expected integer between 1 and 1000") } options.XgressDial.QueueLength = uint16(length) } else { - return nil, errors.New("invalid value for 'xgressDialQueueLength ', expected integer between 1 and 1000") + return nil, errors.New("invalid value for 'xgressDialQueueLength', expected integer between 1 and 1000") } }