Merge branch 'main' into transwarp_2

This commit is contained in:
Michael Quigley
2021-03-09 16:00:47 -08:00
10 changed files with 124 additions and 94 deletions
+3
View File
@@ -6,6 +6,7 @@ on:
- main
- release-*
pull_request:
workflow_dispatch:
jobs:
build:
@@ -13,6 +14,8 @@ jobs:
steps:
- name: Git Checkout
uses: actions/checkout@v2
with:
persist-credentials: false
- name: Install Go
uses: actions/setup-go@v2
+63 -45
View File
@@ -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)
}
}
}
@@ -702,58 +708,70 @@ 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 s.Rerouting.CompareAndSwap(false, true) {
defer s.Rerouting.Set(false)
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 {
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 {
logrus.Errorf("error sending route to [r/%s] (%s)", cq.Path[i].Id, err)
}
}
logrus.Infof("rerouted session [s/%s]", s.Id.Token)
network.CircuitUpdated(s.Id, s.Circuit)
return nil
} else {
return err
}
} else {
logrus.Infof("not rerouting [s/%s], already in progress", s.Id.Token)
return nil
}
}
func (network *Network) smartReroute(s *Session, cq *Circuit) error {
if s.Rerouting.CompareAndSwap(false, true) {
defer s.Rerouting.Set(false)
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.Debugf("rerouted session [s/%s]", s.Id.Token)
network.CircuitUpdated(s.Id, s.Circuit)
return nil
} else {
return err
logrus.Infof("not rerouting [s/%s], already in progress", 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
+2
View File
@@ -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
}
+23 -1
View File
@@ -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
}
+3 -3
View File
@@ -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")
}
}
+10 -9
View File
@@ -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
}
+1 -1
View File
@@ -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()
}
+1 -1
View File
@@ -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)
+1 -1
View File
@@ -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)
+17 -33
View File
@@ -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