Add options for 'termination timeout' and 'route timeout' to 'network' section of controller configuration. Implement. (#151)

This commit is contained in:
Michael Quigley
2021-01-21 13:20:08 -05:00
parent 30b52af7ec
commit aceb021a77
2 changed files with 28 additions and 7 deletions
+6 -6
View File
@@ -324,7 +324,7 @@ func (network *Network) CreateSession(srcR *Router, clientId *identity.TokenId,
// 5: Route Egress
rms[len(rms)-1].Egress.PeerData = clientId.Data
peerData, err := sendRoute(circuit.Path[len(circuit.Path)-1], rms[len(rms)-1])
peerData, err := sendRoute(circuit.Path[len(circuit.Path)-1], rms[len(rms)-1], network.options.TerminationTimeout)
if err != nil {
strategy.NotifyEvent(xt.NewDialFailedEvent(terminator))
retryCount++
@@ -339,7 +339,7 @@ func (network *Network) CreateSession(srcR *Router, clientId *identity.TokenId,
// 6: Create Intermediate Routes
for i := 0; i < len(circuit.Path)-1; i++ {
_, err = sendRoute(circuit.Path[i], rms[i])
_, err = sendRoute(circuit.Path[i], rms[i], network.options.RouteTimeout)
if err != nil {
return nil, err
}
@@ -655,7 +655,7 @@ func (network *Network) rerouteSession(s *session) error {
}
for i := 0; i < len(cq.Path); i++ {
if _, err := sendRoute(cq.Path[i], rms[i]); err != nil {
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)
}
}
@@ -682,7 +682,7 @@ func (network *Network) smartReroute(s *session, cq *Circuit) error {
}
for i := 0; i < len(cq.Path); i++ {
if _, err := sendRoute(cq.Path[i], rms[i]); err != nil {
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)
}
}
@@ -721,7 +721,7 @@ func (network *Network) AcceptMetrics(metrics *metrics_pb.MetricsMessage) {
}
}
func sendRoute(r *Router, createMsg *ctrl_pb.Route) (xt.PeerData, error) {
func sendRoute(r *Router, createMsg *ctrl_pb.Route, timeout time.Duration) (xt.PeerData, error) {
pfxlog.Logger().Debugf("sending Create route message to [r/%s] for [s/%s]", r.Id, createMsg.SessionId)
body, err := proto.Marshal(createMsg)
@@ -754,7 +754,7 @@ func sendRoute(r *Router, createMsg *ctrl_pb.Route) (xt.PeerData, error) {
}
return nil, fmt.Errorf("unexpected response type %v received in reply to route request", msg.ContentType)
case <-time.After(10 * time.Second):
case <-time.After(timeout):
pfxlog.Logger().Errorf("timed out waiting for response to route message from [r/%s] for [s/%s]", r.Id, createMsg.SessionId)
return nil, errors.New("timeout")
}
+22 -1
View File
@@ -19,6 +19,7 @@ package network
import (
"errors"
"github.com/michaelquigley/pfxlog"
"time"
)
type Options struct {
@@ -27,11 +28,15 @@ type Options struct {
RerouteFraction float32
RerouteCap uint32
}
TerminationTimeout time.Duration
RouteTimeout time.Duration
}
func DefaultOptions() *Options {
options := &Options{
CycleSeconds: 15,
CycleSeconds: 15,
TerminationTimeout: 10 * time.Second,
RouteTimeout: 10 * time.Second,
}
options.Smart.RerouteFraction = 0.02
options.Smart.RerouteCap = 4
@@ -49,6 +54,22 @@ func LoadOptions(src map[interface{}]interface{}) (*Options, error) {
}
}
if value, found := src["terminationTimeoutSeconds"]; found {
if terminationTimeoutSeconds, ok := value.(int); ok {
options.TerminationTimeout = time.Duration(terminationTimeoutSeconds) * time.Second
} else {
return nil, errors.New("invalid value for 'terminationTimeoutSeconds'")
}
}
if value, found := src["routeTimeoutSeconds"]; found {
if routeTimeoutSeconds, ok := value.(int); ok {
options.RouteTimeout = time.Duration(routeTimeoutSeconds) * time.Second
} else {
return nil, errors.New("invalid value for 'routeTimeoutSeconds'")
}
}
if value, found := src["smart"]; found {
if submap, ok := value.(map[interface{}]interface{}); ok {
if value, found := submap["rerouteFraction"]; found {