gapmind/services/pathfinder/main.go
2026-09-11 23:59:05 -04:00

196 lines
5.8 KiB
Go

package main
import (
"bytes"
"encoding/json"
"flag"
"fmt"
"io"
"math"
"net"
"net/http"
"net/url"
"os"
"sort"
"strconv"
"strings"
"time"
)
type Probe struct {
Latitude float64 `json:"latitude"`
Longitude float64 `json:"longitude"`
Milliseconds float64 `json:"milliseconds"`
}
type Estimate struct {
Latitude float64 `json:"latitude"`
Longitude float64 `json:"longitude"`
Confidence float64 `json:"confidence"`
Samples int `json:"samples"`
}
type relayCommand struct {
ID string `json:"id"`
Type string `json:"type"`
Host string `json:"host"`
Port string `json:"port"`
}
type relayEnvelope struct {
Command *relayCommand `json:"command"`
}
func main() {
listen := flag.String("listen", "127.0.0.1:8092", "HTTP listen address")
flag.Parse()
go heartbeat()
go relay()
http.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusNoContent) })
http.HandleFunc("/probe", probe)
http.HandleFunc("/triangulate", triangulate)
fmt.Println(http.ListenAndServe(*listen, nil))
}
func probe(w http.ResponseWriter, r *http.Request) {
host := r.URL.Query().Get("host")
if host == "" {
http.Error(w, "host is required", http.StatusBadRequest)
return
}
port := r.URL.Query().Get("port")
if port == "" {
port = "443"
}
milliseconds, err := measure(host, port)
if err != nil {
http.Error(w, err.Error(), http.StatusBadGateway)
return
}
_ = json.NewEncoder(w).Encode(map[string]any{"host": host, "milliseconds": milliseconds})
}
func measure(host, port string) (float64, error) {
started := time.Now()
connection, err := net.DialTimeout("tcp", net.JoinHostPort(host, port), 3*time.Second)
if err != nil {
return 0, err
}
_ = connection.Close()
return float64(time.Since(started).Microseconds()) / 1000, nil
}
func triangulate(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "POST required", http.StatusMethodNotAllowed)
return
}
var probes []Probe
if err := json.NewDecoder(r.Body).Decode(&probes); err != nil || len(probes) < 2 {
http.Error(w, "at least two valid probes are required", http.StatusBadRequest)
return
}
valid := make([]Probe, 0, len(probes))
for _, probe := range probes {
if probe.Latitude >= -90 && probe.Latitude <= 90 && probe.Longitude >= -180 && probe.Longitude <= 180 && probe.Milliseconds > 0 {
valid = append(valid, probe)
}
}
if len(valid) < 2 {
http.Error(w, "at least two valid probes are required", http.StatusBadRequest)
return
}
_ = json.NewEncoder(w).Encode(weighted(valid))
}
func weighted(probes []Probe) Estimate {
var latitude, longitude, totalWeight float64
values := make([]float64, len(probes))
for index, probe := range probes {
weight := 1 / math.Max(probe.Milliseconds, 1)
latitude += probe.Latitude * weight
longitude += probe.Longitude * weight
totalWeight += weight
values[index] = probe.Milliseconds
}
sort.Float64s(values)
median, variance := values[len(values)/2], 0.0
for _, value := range values {
variance += math.Pow(value-median, 2)
}
confidence := math.Max(0, math.Min(1, 1-math.Sqrt(variance/float64(len(values)))/math.Max(median, 1)))
return Estimate{Latitude: latitude / totalWeight, Longitude: longitude / totalWeight, Confidence: confidence, Samples: len(probes)}
}
func heartbeat() {
central := os.Getenv("CENTRAL_URL")
if central == "" {
return
}
host, _ := os.Hostname()
latitude, _ := strconv.ParseFloat(os.Getenv("NODE_LATITUDE"), 64)
longitude, _ := strconv.ParseFloat(os.Getenv("NODE_LONGITUDE"), 64)
for {
body, _ := json.Marshal(map[string]any{"nodeId": host, "kind": "pathfinder", "address": os.Getenv("PEER_PUBLIC_URL"), "version": "0.2.0", "status": "online", "metadata": map[string]any{"transport": "outbound-relay", "latitude": latitude, "longitude": longitude}})
request, _ := http.NewRequest(http.MethodPost, central+"/api/peers/heartbeat", bytes.NewReader(body))
request.Header.Set("Content-Type", "application/json")
if token := os.Getenv("PEER_API_TOKEN"); token != "" {
request.Header.Set("Authorization", "Bearer "+token)
}
response, err := http.DefaultClient.Do(request)
if err == nil {
_, _ = io.Copy(io.Discard, response.Body)
_ = response.Body.Close()
}
time.Sleep(time.Minute)
}
}
func relay() {
central := strings.TrimRight(os.Getenv("CENTRAL_URL"), "/")
if central == "" {
return
}
host, _ := os.Hostname()
client := &http.Client{Timeout: 35 * time.Second}
for {
request, err := http.NewRequest(http.MethodGet, central+"/api/peers/commands?nodeId="+url.QueryEscape(host), nil)
if err != nil {
time.Sleep(5 * time.Second)
continue
}
if token := os.Getenv("PEER_API_TOKEN"); token != "" {
request.Header.Set("Authorization", "Bearer "+token)
}
response, err := client.Do(request)
if err != nil || response.StatusCode != http.StatusOK {
if response != nil {
_ = response.Body.Close()
}
time.Sleep(5 * time.Second)
continue
}
var envelope relayEnvelope
err = json.NewDecoder(response.Body).Decode(&envelope)
_ = response.Body.Close()
if err != nil || envelope.Command == nil || envelope.Command.Type != "probe" {
continue
}
milliseconds, probeErr := measure(envelope.Command.Host, envelope.Command.Port)
result := map[string]any{"id": envelope.Command.ID}
if probeErr != nil {
result["error"] = probeErr.Error()
} else {
result["milliseconds"] = milliseconds
}
body, _ := json.Marshal(result)
request, err = http.NewRequest(http.MethodPost, central+"/api/peers/results", bytes.NewReader(body))
if err == nil {
request.Header.Set("Content-Type", "application/json")
if token := os.Getenv("PEER_API_TOKEN"); token != "" {
request.Header.Set("Authorization", "Bearer "+token)
}
if response, err := client.Do(request); err == nil {
_, _ = io.Copy(io.Discard, response.Body)
_ = response.Body.Close()
}
}
}
}