package main import ( "bytes" "encoding/base64" "encoding/json" "fmt" "io" "log" "net" "net/http" "os" "runtime" "strconv" "sync" "time" "github.com/google/uuid" probing "github.com/prometheus-community/pro-bing" "github.com/shirou/gopsutil/v3/cpu" "github.com/shirou/gopsutil/v3/host" "github.com/shirou/gopsutil/v3/load" "github.com/shirou/gopsutil/v3/mem" ) // --- Configuration --- var ( NodeUUID = getEnv("NODE_UUID", uuid.New().String()) TargetServiceUUID = getEnv("TARGET_SERVICE_UUID", "ab73d00a-8169-46bb-997d-f13e5f760973") BasicAuthUser = getEnv("BASIC_AUTH_USERNAME", "") BasicAuthPass = getEnv("BASIC_AUTH_PASSWORD", "") ServerBaseURL = getEnv("SERVER_URL", "https://test.mystaginglab.net") UpdateInterval = getEnvInt("UPDATE_INTERVAL_SECONDS", 10) PeersFile = getEnv("PEERS_FILE", fmt.Sprintf("known_peers_%s.json", NodeUUID)) LocalIP = getLocalIP() KnownPeers = make(map[string]string) KnownPeersMu sync.RWMutex ) // --- Structs for JSON --- type StatusData struct { UptimeSeconds uint64 `json:"uptime_seconds"` LoadAvg []float64 `json:"load_avg"` MemoryUsagePercent float64 `json:"memory_usage_percent"` } type Payload struct { Node string `json:"node"` Timestamp string `json:"timestamp"` Status StatusData `json:"status"` Pings map[string]float64 `json:"pings"` } type PeerInfo struct { IP string `json:"ip"` LastSeen string `json:"last_seen"` } type ServerResponse struct { Message string `json:"message"` Peers map[string]PeerInfo `json:"peers"` } // --- Helper Functions --- func getEnv(key, fallback string) string { if value, ok := os.LookupEnv(key); ok { return value } return fallback } func getEnvInt(key string, fallback int) int { if value, ok := os.LookupEnv(key); ok { if i, err := strconv.Atoi(value); err == nil { return i } } return fallback } func getLocalIP() string { conn, err := net.Dial("udp", "8.8.8.8:80") if err != nil { return "127.0.0.1" } defer conn.Close() localAddr := conn.LocalAddr().(*net.UDPAddr) return localAddr.IP.String() } // --- Peer Persistence --- func loadPeers() { KnownPeersMu.Lock() defer KnownPeersMu.Unlock() data, err := os.ReadFile(PeersFile) if err != nil { log.Printf("No existing peers file found or error reading: %v", err) return } // The Python script handles both {uuid: ip} and {uuid: {ip: ip}} // We'll decode into a generic map first to handle flexibility var raw map[string]interface{} if err := json.Unmarshal(data, &raw); err != nil { log.Printf("Error decoding peers JSON: %v", err) return } for k, v := range raw { if ipStr, ok := v.(string); ok { KnownPeers[k] = ipStr } else if ipMap, ok := v.(map[string]interface{}); ok { if ip, ok := ipMap["ip"].(string); ok { KnownPeers[k] = ip } } } log.Printf("Loaded %d known peers from %s", len(KnownPeers), PeersFile) } func savePeers() { KnownPeersMu.RLock() defer KnownPeersMu.RUnlock() data, err := json.MarshalIndent(KnownPeers, "", " ") if err != nil { log.Printf("Error marshaling peers: %v", err) return } if err := os.WriteFile(PeersFile, data, 0644); err != nil { log.Printf("Error saving peers to file: %v", err) } } // --- Metrics Collection --- func getSystemMetrics() StatusData { uptime, _ := host.Uptime() var loadValues []float64 avg, err := load.Avg() if err != nil { // Fallback for Windows/Non-Unix like the Python script c, _ := cpu.Percent(500*time.Millisecond, false) cpuUsage := 0.0 if len(c) > 0 { cpuUsage = c[0] / 100.0 } loadValues = []float64{cpuUsage, cpuUsage * 0.9, cpuUsage * 0.8} } else { loadValues = []float64{avg.Load1, avg.Load5, avg.Load15} } v, _ := mem.VirtualMemory() return StatusData{ UptimeSeconds: uptime, LoadAvg: loadValues, MemoryUsagePercent: v.UsedPercent, } } // --- Ping Logic --- func performPings(targets map[string]string) map[string]float64 { results := make(map[string]float64) // Ping self results[NodeUUID] = runPing(LocalIP) // Ping peers for uuid, ip := range targets { if uuid == NodeUUID { continue } results[uuid] = runPing(ip) } return results } func runPing(ip string) float64 { pinger, err := probing.NewPinger(ip) if err != nil { return -1.0 } // On Linux, this requires sudo or 'setcap cap_net_raw+ep' // If not root, you can try: pinger.SetPrivileged(false) which uses UDP pings if runtime.GOOS == "windows" { pinger.SetPrivileged(true) } else { pinger.SetPrivileged(false) // Try unprivileged first } pinger.Count = 1 pinger.Timeout = 2 * time.Second err = pinger.Run() if err != nil { return -1.0 } stats := pinger.Statistics() if stats.PacketsRecv > 0 { return float64(stats.AvgRtt.Microseconds()) / 1000.0 } return -1.0 } // --- Main Loop --- func main() { log.Printf("Starting Go Node Client %s", NodeUUID) log.Printf("Local IP: %s | Server: %s", LocalIP, ServerBaseURL) loadPeers() ticker := time.NewTicker(time.Duration(UpdateInterval) * time.Second) client := &http.Client{Timeout: 15 * time.Second} for ; ; <-ticker.C { metrics := getSystemMetrics() KnownPeersMu.RLock() targets := make(map[string]string) for k, v := range KnownPeers { targets[k] = v } KnownPeersMu.RUnlock() pingResults := performPings(targets) payload := Payload{ Node: NodeUUID, Timestamp: time.Now().UTC().Format(time.RFC3339), Status: metrics, Pings: pingResults, } jsonPayload, _ := json.Marshal(payload) url := fmt.Sprintf("%s/%s/%s/", ServerBaseURL, TargetServiceUUID, NodeUUID) req, err := http.NewRequest("PUT", url, bytes.NewBuffer(jsonPayload)) if err != nil { log.Printf("Error creating request: %v", err) continue } req.Header.Set("Content-Type", "application/json") if BasicAuthUser != "" { auth := BasicAuthUser + ":" + BasicAuthPass req.Header.Set("Authorization", "Basic "+base64.StdEncoding.EncodeToString([]byte(auth))) } resp, err := client.Do(req) if err != nil { log.Printf("Request failed: %v", err) continue } body, _ := io.ReadAll(resp.Body) resp.Body.Close() if resp.StatusCode == http.StatusOK { var serverResp ServerResponse if err := json.Unmarshal(body, &serverResp); err == nil { log.Printf("Update sent. Server: %s", serverResp.Message) // Update peers newPeers := make(map[string]string) for id, info := range serverResp.Peers { newPeers[id] = info.IP } KnownPeersMu.Lock() KnownPeers = newPeers KnownPeersMu.Unlock() savePeers() } } else { log.Printf("Server returned error %d: %s", resp.StatusCode, string(body)) } } }