From 69d3925167e18441827603239b4c566124860e9a Mon Sep 17 00:00:00 2001 From: Owen Date: Thu, 30 Jul 2026 16:46:30 -0400 Subject: [PATCH] move exit node ping to module --- exitnode/exitnode.go | 176 +++++++++++++++++++++++++++++++++++++++++++ newt/handlers.go | 126 +------------------------------ newt/types.go | 31 +++----- 3 files changed, 187 insertions(+), 146 deletions(-) create mode 100644 exitnode/exitnode.go diff --git a/exitnode/exitnode.go b/exitnode/exitnode.go new file mode 100644 index 0000000..1555796 --- /dev/null +++ b/exitnode/exitnode.go @@ -0,0 +1,176 @@ +// Package exitnode implements the exit-node ping dance run before +// registering with the server: request the candidate exit nodes, ping each +// one over HTTP, and report the results so the server can pick the best one. +// It is shared between newt and olm, which both register the same way. +package exitnode + +import ( + "net/http" + "strings" + "time" + + "github.com/fosrl/newt/logger" +) + +// ExitNodeData is the payload the server sends in response to a +// "*/ping/request" message. +type ExitNodeData struct { + ExitNodes []ExitNode `json:"exitNodes"` + ChainId string `json:"chainId"` +} + +// ExitNode is a candidate exit node offered by the server for ping selection. +type ExitNode struct { + ID int `json:"exitNodeId"` + Name string `json:"exitNodeName"` + Endpoint string `json:"endpoint"` + Weight float64 `json:"weight"` + WasPreviouslyConnected bool `json:"wasPreviouslyConnected"` +} + +// ExitNodePingResult is the measured latency (or error) for one exit node, +// sent back to the server in the "*/wg/register" message's pingResults field. +type ExitNodePingResult struct { + ExitNodeID int `json:"exitNodeId"` + LatencyMs int64 `json:"latencyMs"` + Weight float64 `json:"weight"` + Error string `json:"error,omitempty"` + Name string `json:"exitNodeName"` + Endpoint string `json:"endpoint"` + WasPreviouslyConnected bool `json:"wasPreviouslyConnected"` +} + +// PingExitNodes pings the given exit nodes over HTTP and returns a per-node +// ExitNodePingResult suitable for inclusion in a wg/register message's +// pingResults field, so the server can select the best exit node. +// +// If there's only one exit node, or preferEndpoint names one of them, the +// matching node is returned immediately with LatencyMs 0 and no pinging is +// done. Otherwise every node is pinged pingAttempts times over HTTP GET +// /ping and the average latency of successful attempts is used. +// +// When alreadyConnected is true, a node flagged WasPreviouslyConnected is +// excluded from the results as long as at least one other healthy node is +// available, biasing reconnects toward switching away from a possibly +// degraded node. +func PingExitNodes(exitNodes []ExitNode, preferEndpoint string, alreadyConnected bool) []ExitNodePingResult { + if len(exitNodes) == 0 { + return nil + } + + if len(exitNodes) == 1 || preferEndpoint != "" { + selected := exitNodes[0] + if preferEndpoint != "" { + for _, node := range exitNodes { + if node.Endpoint == preferEndpoint { + selected = node + break + } + } + } + + logger.Debug("Only one exit node available, using it directly: %s", selected.Endpoint) + + return []ExitNodePingResult{ + { + ExitNodeID: selected.ID, + LatencyMs: 0, + Weight: selected.Weight, + Error: "", + Name: selected.Name, + Endpoint: selected.Endpoint, + WasPreviouslyConnected: selected.WasPreviouslyConnected, + }, + } + } + + type nodeResult struct { + Node ExitNode + Latency time.Duration + Err error + } + + results := make([]nodeResult, len(exitNodes)) + const pingAttempts = 3 + for i, node := range exitNodes { + var totalLatency time.Duration + var lastErr error + successes := 0 + httpClient := &http.Client{ + Timeout: 5 * time.Second, + } + url := node.Endpoint + if !strings.HasPrefix(url, "http://") && !strings.HasPrefix(url, "https://") { + url = "http://" + url + } + if !strings.HasSuffix(url, "/ping") { + url = strings.TrimRight(url, "/") + "/ping" + } + for j := 0; j < pingAttempts; j++ { + start := time.Now() + resp, err := httpClient.Get(url) + latency := time.Since(start) + if err != nil { + lastErr = err + logger.Warn("Failed to ping exit node %d (%s) attempt %d: %v", node.ID, url, j+1, err) + continue + } + resp.Body.Close() + totalLatency += latency + successes++ + } + var avgLatency time.Duration + if successes > 0 { + avgLatency = totalLatency / time.Duration(successes) + } + if successes == 0 { + results[i] = nodeResult{Node: node, Latency: 0, Err: lastErr} + } else { + results[i] = nodeResult{Node: node, Latency: avgLatency, Err: nil} + } + } + + var pingResults []ExitNodePingResult + for _, res := range results { + errMsg := "" + if res.Err != nil { + errMsg = res.Err.Error() + } + pingResults = append(pingResults, ExitNodePingResult{ + ExitNodeID: res.Node.ID, + LatencyMs: res.Latency.Milliseconds(), + Weight: res.Node.Weight, + Error: errMsg, + Name: res.Node.Name, + Endpoint: res.Node.Endpoint, + WasPreviouslyConnected: res.Node.WasPreviouslyConnected, + }) + } + + if alreadyConnected { + var filteredPingResults []ExitNodePingResult + previouslyConnectedNodeIdx := -1 + for i, res := range pingResults { + if res.WasPreviouslyConnected { + previouslyConnectedNodeIdx = i + } + } + goodNodeCount := 0 + for i, res := range pingResults { + if i != previouslyConnectedNodeIdx && res.LatencyMs > 0 && res.Error == "" { + goodNodeCount++ + } + } + if previouslyConnectedNodeIdx != -1 && goodNodeCount > 0 { + for i, res := range pingResults { + if i != previouslyConnectedNodeIdx { + filteredPingResults = append(filteredPingResults, res) + } + } + pingResults = filteredPingResults + logger.Info("Excluding previously connected exit node from ping results due to other available nodes") + } + } + + return pingResults +} diff --git a/newt/handlers.go b/newt/handlers.go index 55c7f4e..d1530ab 100644 --- a/newt/handlers.go +++ b/newt/handlers.go @@ -10,13 +10,13 @@ import ( "net/http" "os" "os/signal" - "strings" "syscall" "time" "github.com/fosrl/newt/authdaemon" "github.com/fosrl/newt/browsergateway" "github.com/fosrl/newt/docker" + "github.com/fosrl/newt/exitnode" "github.com/fosrl/newt/healthcheck" "github.com/fosrl/newt/internal/state" "github.com/fosrl/newt/internal/telemetry" @@ -140,129 +140,7 @@ func (n *Newt) registerHandlers(ctx context.Context) { return } - if len(exitNodes) == 1 || n.config.PreferEndpoint != "" { - logger.Debug("Only one exit node available, using it directly: %s", exitNodes[0].Endpoint) - - if n.config.PreferEndpoint != "" { - for _, node := range exitNodes { - if node.Endpoint == n.config.PreferEndpoint { - exitNodes[0] = node - break - } - } - } - - pingResults := []ExitNodePingResult{ - { - ExitNodeID: exitNodes[0].ID, - LatencyMs: 0, - Weight: exitNodes[0].Weight, - Error: "", - Name: exitNodes[0].Name, - Endpoint: exitNodes[0].Endpoint, - WasPreviouslyConnected: exitNodes[0].WasPreviouslyConnected, - }, - } - - chainId := generateChainId() - n.pendingRegisterChainId = chainId - n.stopFunc = n.client.SendMessageInterval(topicWGRegister, map[string]interface{}{ - "publicKey": n.publicKey.String(), - "pingResults": pingResults, - "newtVersion": n.config.Version, - "chainId": chainId, - }, 2*time.Second) - - return - } - - type nodeResult struct { - Node ExitNode - Latency time.Duration - Err error - } - - results := make([]nodeResult, len(exitNodes)) - const pingAttempts = 3 - for i, node := range exitNodes { - var totalLatency time.Duration - var lastErr error - successes := 0 - httpClient := &http.Client{ - Timeout: 5 * time.Second, - } - url := node.Endpoint - if !strings.HasPrefix(url, "http://") && !strings.HasPrefix(url, "https://") { - url = "http://" + url - } - if !strings.HasSuffix(url, "/ping") { - url = strings.TrimRight(url, "/") + "/ping" - } - for j := 0; j < pingAttempts; j++ { - start := time.Now() - resp, err := httpClient.Get(url) - latency := time.Since(start) - if err != nil { - lastErr = err - logger.Warn("Failed to ping exit node %d (%s) attempt %d: %v", node.ID, url, j+1, err) - continue - } - resp.Body.Close() - totalLatency += latency - successes++ - } - var avgLatency time.Duration - if successes > 0 { - avgLatency = totalLatency / time.Duration(successes) - } - if successes == 0 { - results[i] = nodeResult{Node: node, Latency: 0, Err: lastErr} - } else { - results[i] = nodeResult{Node: node, Latency: avgLatency, Err: nil} - } - } - - var pingResults []ExitNodePingResult - for _, res := range results { - errMsg := "" - if res.Err != nil { - errMsg = res.Err.Error() - } - pingResults = append(pingResults, ExitNodePingResult{ - ExitNodeID: res.Node.ID, - LatencyMs: res.Latency.Milliseconds(), - Weight: res.Node.Weight, - Error: errMsg, - Name: res.Node.Name, - Endpoint: res.Node.Endpoint, - WasPreviouslyConnected: res.Node.WasPreviouslyConnected, - }) - } - - if n.connected { - var filteredPingResults []ExitNodePingResult - previouslyConnectedNodeIdx := -1 - for i, res := range pingResults { - if res.WasPreviouslyConnected { - previouslyConnectedNodeIdx = i - } - } - goodNodeCount := 0 - for i, res := range pingResults { - if i != previouslyConnectedNodeIdx && res.LatencyMs > 0 && res.Error == "" { - goodNodeCount++ - } - } - if previouslyConnectedNodeIdx != -1 && goodNodeCount > 0 { - for i, res := range pingResults { - if i != previouslyConnectedNodeIdx { - filteredPingResults = append(filteredPingResults, res) - } - } - pingResults = filteredPingResults - logger.Info("Excluding previously connected exit node from ping results due to other available nodes") - } - } + pingResults := exitnode.PingExitNodes(exitNodes, n.config.PreferEndpoint, n.connected) chainId := generateChainId() n.pendingRegisterChainId = chainId diff --git a/newt/types.go b/newt/types.go index 649111d..8b2a2e8 100644 --- a/newt/types.go +++ b/newt/types.go @@ -2,6 +2,7 @@ package newt import ( wgclients "github.com/fosrl/newt/clients" + "github.com/fosrl/newt/exitnode" "github.com/fosrl/newt/healthcheck" ) @@ -35,28 +36,14 @@ type TargetData struct { Targets []string `json:"targets"` } -type ExitNodeData struct { - ExitNodes []ExitNode `json:"exitNodes"` - ChainId string `json:"chainId"` -} - -type ExitNode struct { - ID int `json:"exitNodeId"` - Name string `json:"exitNodeName"` - Endpoint string `json:"endpoint"` - Weight float64 `json:"weight"` - WasPreviouslyConnected bool `json:"wasPreviouslyConnected"` -} - -type ExitNodePingResult struct { - ExitNodeID int `json:"exitNodeId"` - LatencyMs int64 `json:"latencyMs"` - Weight float64 `json:"weight"` - Error string `json:"error,omitempty"` - Name string `json:"exitNodeName"` - Endpoint string `json:"endpoint"` - WasPreviouslyConnected bool `json:"wasPreviouslyConnected"` -} +// ExitNodeData, ExitNode and ExitNodePingResult are aliases for the shared +// exit-node ping dance types in package exitnode, kept here so existing code +// in this package can keep referring to them unqualified. +type ( + ExitNodeData = exitnode.ExitNodeData + ExitNode = exitnode.ExitNode + ExitNodePingResult = exitnode.ExitNodePingResult +) type BlueprintResult struct { Success bool `json:"success"`