Remade the client in go, now named monitor-agent. The agents now support http basic auth

This commit is contained in:
Kalzu Rekku
2026-07-25 00:01:46 +03:00
parent c98081d360
commit 48bdab08b8
7 changed files with 578 additions and 1608 deletions
+142 -41
View File
@@ -10,11 +10,15 @@ from datetime import datetime, timezone
import platform
import socket # For getting local IP
import sys
import base64
import urllib.parse
# --- Install necessary libraries if not already present ---
try:
import psutil # For system metrics
from pythonping import ping as python_ping # Renamed to avoid conflict with common 'ping'
from pythonping import (
ping as python_ping,
) # Renamed to avoid conflict with common 'ping'
except ImportError:
print("Required libraries 'psutil' and 'pythonping' not found.")
print("Please install them: pip install psutil pythonping")
@@ -25,25 +29,29 @@ except ImportError:
NODE_UUID = os.environ.get("NODE_UUID", str(uuid.uuid4()))
# The UUID of the target monitoring service (the main.py server).
# IMPORTANT: This MUST match the SERVICE_UUID of your running FastAPI server.
# You can get this from the server's initial console output or by accessing its root endpoint ('/').
TARGET_SERVICE_UUID = os.environ.get(
"TARGET_SERVICE_UUID", "REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID"
"TARGET_SERVICE_UUID", "ab73d00a-8169-46bb-997d-f13e5f760973"
)
# Optional basic auth for the server (UTF-8 compatible)
BASIC_AUTH_USERNAME = os.environ.get("BASIC_AUTH_USERNAME", "")
BASIC_AUTH_PASSWORD = os.environ.get("BASIC_AUTH_PASSWORD", "")
# The base URL of the FastAPI monitoring service
SERVER_BASE_URL = os.environ.get("SERVER_URL", "http://localhost:8000")
SERVER_BASE_URL = os.environ.get("SERVER_URL", "https://test.mystaginglab.net")
# How often to send status updates (in seconds)
UPDATE_INTERVAL_SECONDS = int(os.environ.get("UPDATE_INTERVAL_SECONDS", 5))
UPDATE_INTERVAL_SECONDS = int(os.environ.get("UPDATE_INTERVAL_SECONDS", 10))
# SSL verification (set to False if using self-signed certificates)
SSL_VERIFY = os.environ.get("SSL_VERIFY", "true").lower() == "true"
# File to store known peers' UUIDs and IPs for persistence
PEERS_FILE = os.environ.get("PEERS_FILE", f"known_peers_{NODE_UUID}.json")
# --- Logging Configuration ---
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
)
logger = logging.getLogger("NodeClient")
@@ -62,13 +70,14 @@ try:
except Exception:
logger.warning("Could not determine local IP, defaulting to 127.0.0.1 for pings.")
# --- File Operations for Peers ---
def load_peers():
"""Loads known peers (UUID: IP) from a local JSON file."""
global known_peers
if os.path.exists(PEERS_FILE):
try:
with open(PEERS_FILE, 'r') as f:
with open(PEERS_FILE, "r") as f:
loaded_data = json.load(f)
# Ensure loaded peers are in the correct {uuid: ip} format
# Handle cases where the file might contain server's full peer info
@@ -76,28 +85,34 @@ def load_peers():
for k, v in loaded_data.items():
if isinstance(v, str): # Already in {uuid: ip} format
temp_peers[k] = v
elif isinstance(v, dict) and 'ip' in v: # Server's full peer info
temp_peers[k] = v['ip']
elif isinstance(v, dict) and "ip" in v: # Server's full peer info
temp_peers[k] = v["ip"]
known_peers = temp_peers
logger.info(f"Loaded {len(known_peers)} known peers from {PEERS_FILE}")
except json.JSONDecodeError as e:
logger.error(f"Error decoding JSON from {PEERS_FILE}: {e}. Starting with no known peers.")
logger.error(
f"Error decoding JSON from {PEERS_FILE}: {e}. Starting with no known peers."
)
known_peers = {} # Reset if file is corrupt
except Exception as e:
logger.error(f"Error loading peers from {PEERS_FILE}: {e}. Starting with no known peers.")
logger.error(
f"Error loading peers from {PEERS_FILE}: {e}. Starting with no known peers."
)
known_peers = {}
else:
logger.info(f"No existing peers file found at {PEERS_FILE}.")
def save_peers():
"""Saves current known peers (UUID: IP) to a local JSON file."""
try:
with open(PEERS_FILE, 'w') as f:
with open(PEERS_FILE, "w") as f:
json.dump(known_peers, f, indent=2)
logger.debug(f"Saved {len(known_peers)} known peers to {PEERS_FILE}")
except Exception as e:
logger.error(f"Error saving peers to {PEERS_FILE}: {e}")
# --- System Metrics Collection ---
def get_system_metrics():
"""Collects actual system load and memory usage using psutil."""
@@ -112,13 +127,15 @@ def get_system_metrics():
# For cross-platform consistency, we'll use psutil.cpu_percent()
# and simulate 5/15 min averages if os.getloadavg is not available.
load_avg = [0.0, 0.0, 0.0]
if hasattr(os, 'getloadavg'):
if hasattr(os, "getloadavg"):
load_avg = list(os.getloadavg())
else: # Fallback for Windows or systems without getloadavg
# psutil.cpu_percent() gives current CPU utilization over an interval.
# It's not true load average, but a reasonable proxy for monitoring.
# We'll use a short interval to get a "current" load.
cpu_percent = psutil.cpu_percent(interval=0.5) / 100.0 # CPU usage as a fraction
cpu_percent = (
psutil.cpu_percent(interval=0.5) / 100.0
) # CPU usage as a fraction
load_avg = [cpu_percent, cpu_percent * 0.9, cpu_percent * 0.8] # Simulate decay
logger.debug(f"Using psutil.cpu_percent() for load_avg (non-Unix): {load_avg}")
@@ -129,9 +146,10 @@ def get_system_metrics():
return {
"uptime_seconds": uptime_seconds,
"load_avg": [round(l, 2) for l in load_avg],
"memory_usage_percent": round(memory_usage_percent, 2)
"memory_usage_percent": round(memory_usage_percent, 2),
}
# --- Ping Logic ---
def perform_pings(targets: dict[str, str]) -> dict[str, float]:
"""Performs actual pings to target IPs and returns latencies in ms."""
@@ -163,13 +181,36 @@ def perform_pings(targets: dict[str, str]) -> dict[str, float]:
pings_results[peer_uuid] = round(response_list.rtt_avg_ms, 2)
else:
pings_results[peer_uuid] = -1.0 # Indicate failure
logger.debug(f"Ping to {peer_uuid} ({peer_ip}): {pings_results[peer_uuid]}ms")
logger.debug(
f"Ping to {peer_uuid} ({peer_ip}): {pings_results[peer_uuid]}ms"
)
except Exception as e:
logger.warning(f"Failed to ping {peer_uuid} ({peer_ip}): {e}")
pings_results[peer_uuid] = -1.0
return pings_results
# --- Authentication Helper ---
def create_auth_headers():
"""Create authentication headers with UTF-8 support for Basic Auth."""
headers = {"Content-Type": "application/json"}
if BASIC_AUTH_USERNAME and BASIC_AUTH_PASSWORD:
try:
# Create Basic Auth header with UTF-8 encoding
credentials_str = f"{BASIC_AUTH_USERNAME}:{BASIC_AUTH_PASSWORD}"
credentials_b64 = base64.b64encode(credentials_str.encode("utf-8")).decode(
"ascii"
)
headers["Authorization"] = f"Basic {credentials_b64}"
logger.debug("Using HTTP Basic Authentication (UTF-8)")
except Exception as e:
logger.error(f"Failed to create authentication headers: {e}")
return headers
# --- Main Client Logic ---
def run_client():
global known_peers
@@ -179,15 +220,25 @@ def run_client():
logger.info(f"Target Service UUID: {TARGET_SERVICE_UUID}")
logger.info(f"Server URL: {SERVER_BASE_URL}")
logger.info(f"Update Interval: {UPDATE_INTERVAL_SECONDS} seconds")
logger.info(f"SSL Verification: {SSL_VERIFY}")
logger.info(f"Peers file: {PEERS_FILE}")
if BASIC_AUTH_USERNAME and BASIC_AUTH_PASSWORD:
logger.info(f"Basic Auth enabled for user: {BASIC_AUTH_USERNAME}")
else:
logger.info("No Basic Auth configured")
if TARGET_SERVICE_UUID == "REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID":
logger.error("-" * 50)
logger.error("ERROR: TARGET_SERVICE_UUID is not set correctly!")
logger.error("Please replace 'REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID' in the script")
logger.error(
"Please replace 'REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID' in the script"
)
logger.error("or set the environment variable TARGET_SERVICE_UUID.")
logger.error("You can find the server's UUID by running main.py and checking its console output")
logger.error("or by visiting 'http://localhost:8000/' in your browser.")
logger.error(
"You can find the server's UUID by running main.py and checking its console output"
)
logger.error("or by visiting the server's root endpoint in your browser.")
logger.error("-" * 50)
return
@@ -207,65 +258,115 @@ def run_client():
"node": str(NODE_UUID),
"timestamp": datetime.now(timezone.utc).isoformat(),
"status": status_data,
"pings": ping_data
"pings": ping_data,
}
# 4. Define the endpoint URL
endpoint_url = f"{SERVER_BASE_URL}/{TARGET_SERVICE_UUID}/{NODE_UUID}/"
# 5. Send the PUT request
# 5. Create headers with authentication
headers = create_auth_headers()
# 6. Send the PUT request
logger.info(
f"Sending update. Uptime: {status_data['uptime_seconds']}s, "
f"Load: {status_data['load_avg']}, Mem: {status_data['memory_usage_percent']}%, "
f"Pings: {len(ping_data)}"
)
response = requests.put(endpoint_url, json=payload, timeout=15) # Increased timeout
response = requests.put(
endpoint_url,
json=payload,
headers=headers,
timeout=15,
verify=SSL_VERIFY,
)
# 6. Process the response
# 7. Process the response
if response.status_code == 200:
response_data = response.json()
logger.info(f"Successfully sent update. Server message: '{response_data.get('message')}'")
logger.info(
f"Successfully sent update. Server message: '{response_data.get('message')}'"
)
if "peers" in response_data and isinstance(response_data["peers"], dict):
if "peers" in response_data and isinstance(
response_data["peers"], dict
):
# Update known_peers from server response
updated_peers = {}
# The server returns {uuid: {"last_seen": "...", "ip": "..."}}
# We only need the UUID and IP for pinging.
for peer_uuid, peer_info in response_data["peers"].items():
if 'ip' in peer_info:
updated_peers[peer_uuid] = peer_info['ip']
if "ip" in peer_info:
updated_peers[peer_uuid] = peer_info["ip"]
# Log newly discovered peers
newly_discovered = set(updated_peers.keys()) - set(known_peers.keys())
newly_discovered = set(updated_peers.keys()) - set(
known_peers.keys()
)
if newly_discovered:
logger.info(f"Discovered new peer(s): {', '.join(newly_discovered)}")
logger.info(
f"Discovered new peer(s): {', '.join(newly_discovered)}"
)
known_peers = updated_peers
save_peers() # Save updated peers to file for persistence
logger.info(f"Total known peers for pinging: {len(known_peers)}")
else:
logger.warning("Server response did not contain a valid 'peers' field or it was empty.")
logger.warning(
"Server response did not contain a valid 'peers' field or it was empty."
)
else:
logger.error(f"Failed to send update. Status code: {response.status_code}, Response: {response.text}")
if response.status_code == 404:
logger.error("Hint: The TARGET_SERVICE_UUID might be incorrect, or the server isn't running at this endpoint.")
logger.error(
f"Failed to send update. Status code: {response.status_code}, Response: {response.text}"
)
if response.status_code == 401:
logger.error("Authentication failed (401 Unauthorized).")
logger.error(
"Please check your BASIC_AUTH_USERNAME and BASIC_AUTH_PASSWORD."
)
elif response.status_code == 404:
logger.error(
"Endpoint not found (404). The TARGET_SERVICE_UUID might be incorrect, or the server isn't running at this endpoint."
)
elif response.status_code == 422: # Pydantic validation error
logger.error(f"Server validation error (422 Unprocessable Entity): {response.json()}")
try:
error_detail = response.json()
logger.error(
f"Server validation error (422 Unprocessable Entity): {error_detail}"
)
except:
logger.error(
f"Server validation error (422 Unprocessable Entity): {response.text}"
)
except requests.exceptions.SSLError as e:
logger.error(f"SSL Error: {e}")
logger.error(
"If using self-signed certificates, set SSL_VERIFY=false environment variable"
)
except requests.exceptions.Timeout:
logger.error(f"Request timed out after {15} seconds. Is the server running and responsive?")
logger.error(
f"Request timed out after 15 seconds. Is the server running and responsive?"
)
except requests.exceptions.ConnectionError as e:
logger.error(f"Connection error: {e}. Is the server running at {SERVER_BASE_URL}?")
logger.error(
f"Connection error: {e}. Is the server running at {SERVER_BASE_URL}?"
)
except requests.exceptions.RequestException as e:
logger.error(f"An unexpected request error occurred: {e}", exc_info=True)
except json.JSONDecodeError:
logger.error(f"Failed to decode JSON response: {response.text}. Is the server returning valid JSON?")
logger.error(
f"Failed to decode JSON response: {response.text}. Is the server returning valid JSON?"
)
except Exception as e:
logger.error(f"An unexpected error occurred in the client loop: {e}", exc_info=True)
logger.error(
f"An unexpected error occurred in the client loop: {e}", exc_info=True
)
# 7. Wait for the next update
# 8. Wait for the next update
time.sleep(UPDATE_INTERVAL_SECONDS)
if __name__ == "__main__":
run_client()
+103
View File
@@ -0,0 +1,103 @@
#!/bin/bash
# --- Configuration (Edit these before spreading) ---
TARGET_SERVICE_UUID=""
SERVER_URL="https://monitor.example.com"
AUTH_USER="monitor_user"
AUTH_PASS="monitor_user_pass"
UPDATE_INTERVAL=10
# --- Script Logic ---
# 1. Ensure the script is run as root
if [ "$EUID" -ne 0 ]; then
echo "Please run as root (use sudo)"
exit 1
fi
echo "Starting deployment of Monitoring Agent..."
# 2. Install dependencies (libcap2-bin is needed for setcap)
echo "Installing dependencies..."
apt-get update -qq
apt-get install -y libcap2-bin uuid-runtime -qq
# 3. Create the system user if it doesn't exist
if ! id "monitor-agent" &>/dev/null; then
echo "Creating monitor-agent system user..."
useradd -r -s /bin/false monitor-agent
fi
# 4. Create necessary directories
echo "Creating directories..."
mkdir -p /etc/monitor-agent
mkdir -p /var/lib/monitor-agent
chown -R monitor-agent:monitor-agent /var/lib/monitor-agent
# 5. Install the binary
if [ -f "./agent-go" ]; then
echo "Installing binary to /usr/local/bin..."
cp ./agent-go /usr/local/bin/monitor-agent
chmod +x /usr/local/bin/monitor-agent
# Grant ping capabilities without root
setcap cap_net_raw+ep /usr/local/bin/monitor-agent
else
echo "Error: 'monitor-agent' binary not found in current directory!"
exit 1
fi
# 6. Generate a unique NODE_UUID for this specific server
NEW_UUID=$(uuidgen)
echo "Generated unique Node UUID: $NEW_UUID"
# 7. Create the Environment File
echo "Creating configuration file..."
cat <<EOF > /etc/monitor-agent/agent.env
NODE_UUID=$NEW_UUID
TARGET_SERVICE_UUID=$TARGET_SERVICE_UUID
SERVER_URL=$SERVER_URL
BASIC_AUTH_USERNAME=$AUTH_USER
BASIC_AUTH_PASSWORD=$AUTH_PASS
UPDATE_INTERVAL_SECONDS=$UPDATE_INTERVAL
PEERS_FILE=/var/lib/monitor-agent/known_peers.json
EOF
chmod 600 /etc/monitor-agent/agent.env
chown monitor-agent:monitor-agent /etc/monitor-agent/agent.env
# 8. Create the Systemd Service File
echo "Creating systemd service..."
cat <<EOF > /etc/systemd/system/monitor-agent.service
[Unit]
Description=Node Monitoring Agent
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
User=monitor-agent
Group=monitor-agent
WorkingDirectory=/var/lib/monitor-agent
EnvironmentFile=/etc/monitor-agent/agent.env
ExecStart=/usr/local/bin/monitor-agent
CapabilityBoundingSet=CAP_NET_RAW
AmbientCapabilities=CAP_NET_RAW
Restart=always
RestartSec=5
[Install]
WantedBy=multi-user.target
EOF
# 9. Start the service
echo "Reloading systemd and starting service..."
systemctl daemon-reload
systemctl enable monitor-agent
systemctl restart monitor-agent
echo "------------------------------------------------"
echo "Deployment Complete!"
echo "Node UUID: $NEW_UUID"
echo "Check status with: systemctl status monitor-agent"
echo "View logs with: journalctl -u monitor-agent -f"
echo "------------------------------------------------"
+19
View File
@@ -0,0 +1,19 @@
module aget_go
go 1.25.0
require (
github.com/go-ole/go-ole v1.2.6 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect
github.com/power-devops/perfstat v0.0.0-20210106213030-5aafc221ea8c // indirect
github.com/prometheus-community/pro-bing v0.9.1 // indirect
github.com/shirou/gopsutil/v3 v3.24.5 // indirect
github.com/shoenig/go-m1cpu v0.1.6 // indirect
github.com/tklauser/go-sysconf v0.3.12 // indirect
github.com/tklauser/numcpus v0.6.1 // indirect
github.com/yusufpapurcu/wmi v1.2.4 // indirect
golang.org/x/net v0.56.0 // indirect
golang.org/x/sync v0.21.0 // indirect
golang.org/x/sys v0.46.0 // indirect
)
+290
View File
@@ -0,0 +1,290 @@
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))
}
}
}
-832
View File
@@ -1,832 +0,0 @@
import os
import uuid
import time
import requests
import random
import json
import logging
import threading
import socket
from datetime import datetime, timezone
from concurrent.futures import ThreadPoolExecutor
import argparse
from requests.adapters import HTTPAdapter
from urllib3.util.connection import create_connection
# --- Multi-Node Client Configuration ---
TARGET_SERVICE_UUID = os.environ.get(
"TARGET_SERVICE_UUID", "ab73d00a-8169-46bb-997d-f13e5f760973"
)
SERVER_BASE_URL = os.environ.get("SERVER_URL", "http://localhost:8000")
UPDATE_INTERVAL_SECONDS = int(os.environ.get("UPDATE_INTERVAL_SECONDS", 5))
NUM_NODES = int(os.environ.get("NUM_NODES", 3))
# Base IP for loopback binding (127.0.0.x where x starts from this base)
LOOPBACK_IP_BASE = int(os.environ.get("LOOPBACK_IP_BASE", 2)) # Start from 127.0.0.2
# Dynamic node management settings
DYNAMIC_MIN_NODES = int(os.environ.get("DYNAMIC_MIN_NODES", 3))
DYNAMIC_MAX_NODES = int(os.environ.get("DYNAMIC_MAX_NODES", 7))
NODE_CHANGE_INTERVAL = int(os.environ.get("NODE_CHANGE_INTERVAL", 30)) # seconds
NODE_LIFECYCLE_VARIANCE = int(os.environ.get("NODE_LIFECYCLE_VARIANCE", 10)) # seconds
# --- Logging Configuration ---
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s - %(name)s - %(levelname)s - [%(thread)d] - %(message)s",
)
logger = logging.getLogger("MultiNodeClient")
# --- Custom HTTP Adapter for Source IP Binding ---
class SourceIPHTTPAdapter(HTTPAdapter):
def __init__(self, source_ip, *args, **kwargs):
self.source_ip = source_ip
super().__init__(*args, **kwargs)
def init_poolmanager(self, *args, **kwargs):
# Override the socket creation to bind to specific source IP
def custom_create_connection(
address,
timeout=socket._GLOBAL_DEFAULT_TIMEOUT,
source_address=None,
socket_options=None,
):
# Force our custom source address
return create_connection(
address,
timeout,
source_address=(self.source_ip, 0), # 0 = any available port
socket_options=socket_options,
)
# Monkey patch the connection creation
original_create_connection = socket.create_connection
socket.create_connection = custom_create_connection
try:
result = super().init_poolmanager(*args, **kwargs)
finally:
# Restore original function
socket.create_connection = original_create_connection
return result
# --- Enhanced Node Class with IP Binding ---
class SimulatedNode:
def __init__(
self,
node_id: int,
total_nodes: int,
server_url: str,
service_uuid: str,
update_interval: int,
ip_base: int,
persistent_uuid: str = None, # Allow reusing UUIDs for returning nodes
):
self.node_id = node_id
self.node_uuid = persistent_uuid or str(uuid.uuid4())
self.server_url = server_url # Store server URL
self.service_uuid = service_uuid # Store service UUID
self.update_interval = update_interval # Store update interval
self.uptime_seconds = 0
self.known_peers = {}
self.total_nodes = total_nodes
self.running = False
self.is_persistent = (
persistent_uuid is not None
) # Track if this is a returning node
# Assign unique loopback IP to this node using the passed ip_base
self.source_ip = f"127.0.0.{ip_base + node_id - 1}"
# Create requests session with custom adapter for IP binding
self.session = requests.Session()
adapter = SourceIPHTTPAdapter(self.source_ip)
self.session.mount("http://", adapter)
self.session.mount("https://", adapter)
# Each node gets slightly different characteristics
if not self.is_persistent: # Only randomize for new nodes
self.base_load = random.uniform(0.2, 1.0)
self.base_memory = random.uniform(40.0, 70.0)
self.load_variance = random.uniform(0.1, 0.5)
self.memory_variance = random.uniform(5.0, 15.0)
# Some nodes might be "problematic" (higher load/memory)
if random.random() < 0.2: # 20% chance of being a "problematic" node
self.base_load *= 2.0
self.base_memory += 20.0
logger.info(
f"Node {self.node_id} ({self.node_uuid[:8]}) will simulate high resource usage (IP: {self.source_ip})"
)
else:
# Returning nodes keep some consistency but with variation
self.base_load = random.uniform(0.3, 0.8)
self.base_memory = random.uniform(45.0, 65.0)
self.load_variance = random.uniform(0.1, 0.4)
self.memory_variance = random.uniform(5.0, 12.0)
logger.info(
f"Node {self.node_id} ({'returning' if self.is_persistent else 'new'}) will bind to source IP: {self.source_ip}"
)
def test_ip_binding(self):
"""Test if we can bind to the assigned IP address."""
try:
# Try to create a socket and bind to the IP
test_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
test_socket.bind((self.source_ip, 0)) # Bind to any available port
test_socket.close()
logger.debug(
f"Node {self.node_id} successfully tested binding to {self.source_ip}"
)
return True
except OSError as e:
logger.error(f"Node {self.node_id} cannot bind to {self.source_ip}: {e}")
logger.error(
"Make sure the IP address is available on your loopback interface."
)
logger.error(
f"You might need to add it with: sudo ifconfig lo0 alias {self.source_ip} (macOS)"
)
logger.error(f"Or: sudo ip addr add {self.source_ip}/8 dev lo (Linux)")
return False
def generate_node_status_data(self):
"""Generates simulated node status metrics with per-node characteristics."""
self.uptime_seconds += self.update_interval + random.randint(-1, 2)
# Generate load with some randomness but consistent per-node baseline
load_1min = max(
0.1,
self.base_load + random.uniform(-self.load_variance, self.load_variance),
)
load_5min = max(0.1, load_1min * random.uniform(0.8, 1.0))
load_15min = max(0.1, load_5min * random.uniform(0.8, 1.0))
load_avg = [round(load_1min, 2), round(load_5min, 2), round(load_15min, 2)]
# Generate memory usage with baseline + variance
memory_usage = max(
10.0,
min(
95.0,
self.base_memory
+ random.uniform(-self.memory_variance, self.memory_variance),
),
)
return {
"uptime_seconds": self.uptime_seconds,
"load_avg": load_avg,
"memory_usage_percent": round(memory_usage, 2),
}
def generate_ping_data(self):
"""Generates simulated ping latencies to known peers."""
pings = {}
# Ping to self (loopback)
pings[self.node_uuid] = round(random.uniform(0.1, 1.5), 2)
# Ping to known peers
for peer_uuid in self.known_peers.keys():
if peer_uuid != self.node_uuid:
# Simulate network latency with some consistency per peer
base_latency = random.uniform(5.0, 100.0)
variation = random.uniform(-10.0, 10.0)
latency = max(0.1, base_latency + variation)
pings[peer_uuid] = round(latency, 2)
return pings
def send_update(self):
"""Sends a single status update to the server using bound IP."""
try:
status_data = self.generate_node_status_data()
ping_data = self.generate_ping_data()
payload = {
"node": self.node_uuid,
"timestamp": datetime.now(timezone.utc).isoformat(),
"status": status_data,
"pings": ping_data,
}
endpoint_url = f"{self.server_url}/{self.service_uuid}/{self.node_uuid}/"
logger.debug(
f"Node {self.node_id} ({self.source_ip}) sending update. "
f"Uptime: {status_data['uptime_seconds']}s, "
f"Load: {status_data['load_avg'][0]}, Memory: {status_data['memory_usage_percent']}%, "
f"Pings: {len(ping_data)}"
)
# Use the custom session with IP binding
response = self.session.put(endpoint_url, json=payload, timeout=10)
if response.status_code == 200:
response_data = response.json()
if "peers" in response_data and isinstance(
response_data["peers"], dict
):
new_peers = {k: v for k, v in response_data["peers"].items()}
# Log new peer discoveries
newly_discovered = set(new_peers.keys()) - set(
self.known_peers.keys()
)
if newly_discovered:
logger.info(
f"Node {self.node_id} ({self.source_ip}) discovered {len(newly_discovered)} new peer(s)"
)
self.known_peers = new_peers
if (
len(newly_discovered) > 0
or len(self.known_peers) != self.total_nodes - 1
):
logger.debug(
f"Node {self.node_id} ({self.source_ip}) knows {len(self.known_peers)} peers "
f"(expected {self.total_nodes - 1})"
)
return True
else:
logger.error(
f"Node {self.node_id} ({self.source_ip}) failed to send update. "
f"Status: {response.status_code}, Response: {response.text}"
)
return False
except requests.exceptions.Timeout:
logger.error(f"Node {self.node_id} ({self.source_ip}) request timed out")
return False
except requests.exceptions.ConnectionError as e:
logger.error(
f"Node {self.node_id} ({self.source_ip}) connection error: {e}"
)
return False
except Exception as e:
logger.error(
f"Node {self.node_id} ({self.source_ip}) unexpected error: {e}"
)
return False
def run(self):
"""Main loop for this simulated node."""
# Test IP binding before starting
if not self.test_ip_binding():
logger.error(f"Node {self.node_id} cannot start due to IP binding failure")
return
self.running = True
logger.info(
f"Starting Node {self.node_id} ({'returning' if self.is_persistent else 'new'}) with UUID: {self.node_uuid[:8]}... (IP: {self.source_ip})"
)
# Add some initial delay to stagger node starts
initial_delay = self.node_id * 0.5
time.sleep(initial_delay)
consecutive_failures = 0
while self.running:
try:
success = self.send_update()
if success:
consecutive_failures = 0
else:
consecutive_failures += 1
if consecutive_failures >= 3:
logger.warning(
f"Node {self.node_id} ({self.source_ip}) has failed {consecutive_failures} consecutive updates"
)
# Add some jitter to prevent thundering herd
jitter = random.uniform(-1.0, 1.0)
sleep_time = max(1.0, self.update_interval + jitter)
time.sleep(sleep_time)
except KeyboardInterrupt:
logger.info(
f"Node {self.node_id} ({self.source_ip}) received interrupt signal"
)
break
except Exception as e:
logger.error(
f"Node {self.node_id} ({self.source_ip}) unexpected error in main loop: {e}"
)
time.sleep(self.update_interval)
logger.info(f"Node {self.node_id} ({self.source_ip}) stopped")
def stop(self):
"""Stop the node."""
self.running = False
# --- Dynamic Multi-Node Manager ---
class DynamicMultiNodeManager:
def __init__(
self,
min_nodes: int,
max_nodes: int,
server_url: str,
service_uuid: str,
update_interval: int,
ip_base: int,
change_interval: int,
lifecycle_variance: int,
):
self.min_nodes = min_nodes
self.max_nodes = max_nodes
self.server_url = server_url
self.service_uuid = service_uuid
self.update_interval = update_interval
self.ip_base = ip_base
self.change_interval = change_interval
self.lifecycle_variance = lifecycle_variance
# Track active nodes and their threads
self.active_nodes = {} # node_id -> SimulatedNode
self.node_threads = {} # node_id -> Thread
self.node_counter = 1 # For assigning unique node IDs
# Track node history for potential returns
self.departed_nodes = {} # UUID -> {'characteristics', 'last_seen'}
self.running = False
self.manager_thread = None
def get_available_ip_addresses(self):
"""Get list of available IP addresses for new nodes."""
max_possible_ips = self.max_nodes * 2 # Allow some extra IPs
available_ips = []
used_ips = {node.source_ip for node in self.active_nodes.values()}
for i in range(max_possible_ips):
ip = f"127.0.0.{self.ip_base + i}"
if ip not in used_ips:
available_ips.append((i + 1, ip)) # (node_id_offset, ip)
return available_ips
def check_ip_availability(self, num_ips_needed):
"""Check if required number of IP addresses are available."""
logger.info(f"Checking availability of {num_ips_needed} IP addresses...")
available_ips = self.get_available_ip_addresses()
if len(available_ips) < num_ips_needed:
logger.error(
f"Only {len(available_ips)} IP addresses available, need {num_ips_needed}"
)
return False
# Test the first few IPs we'd actually use
test_ips = available_ips[:num_ips_needed]
all_available = True
for node_id_offset, ip in test_ips:
try:
test_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
test_socket.bind((ip, 0))
test_socket.close()
except OSError as e:
logger.error(f"Cannot bind to {ip}: {e}")
all_available = False
if not all_available:
logger.error("Some IP addresses are not available.")
self.show_setup_commands()
return False
logger.info(f"All {num_ips_needed} IP addresses are available!")
return True
def show_setup_commands(self):
"""Show commands for setting up loopback IPs."""
logger.info("=== Loopback IP Setup Commands ===")
logger.info("Run these commands to add the required loopback IP addresses:")
import platform
system = platform.system().lower()
max_ips_needed = self.max_nodes * 2
for i in range(max_ips_needed):
ip = f"127.0.0.{self.ip_base + i}"
if system == "linux":
logger.info(f"sudo ip addr add {ip}/8 dev lo")
elif system == "darwin": # macOS
logger.info(f"sudo ifconfig lo0 alias {ip}")
else:
logger.info(f"Add {ip} to loopback interface (OS: {system})")
logger.info("=" * 40)
def create_new_node(self, return_existing=False):
"""Create a new node, optionally bringing back a departed node."""
available_ips = self.get_available_ip_addresses()
if not available_ips:
logger.warning("No available IP addresses for new nodes")
return None
node_id_offset, source_ip = available_ips[0]
node_id = self.node_counter
self.node_counter += 1
# Decide if we should bring back an old node (30% chance if we have departed nodes)
persistent_uuid = None
if return_existing and self.departed_nodes and random.random() < 0.3:
# Pick a random departed node to bring back
persistent_uuid = random.choice(list(self.departed_nodes.keys()))
logger.info(
f"Bringing back departed node {persistent_uuid[:8]}... as Node {node_id}"
)
# Remove from departed list since it's returning
del self.departed_nodes[persistent_uuid]
node = SimulatedNode(
node_id=node_id,
total_nodes=len(self.active_nodes) + 1, # +1 for this new node
server_url=self.server_url,
service_uuid=self.service_uuid,
update_interval=self.update_interval,
ip_base=self.ip_base + node_id_offset - 1, # Adjust IP calculation
persistent_uuid=persistent_uuid,
)
return node
def start_node(self, node):
"""Start a single node in its own thread."""
thread = threading.Thread(target=node.run, name=f"Node-{node.node_id}")
thread.daemon = True
thread.start()
self.active_nodes[node.node_id] = node
self.node_threads[node.node_id] = thread
logger.info(
f"✅ Started Node {node.node_id} ({node.node_uuid[:8]}...) on {node.source_ip}"
)
def stop_node(self, node_id, permanently=False):
"""Stop a specific node."""
if node_id not in self.active_nodes:
return False
node = self.active_nodes[node_id]
thread = self.node_threads[node_id]
# Store node info for potential return (unless it's permanently leaving)
if not permanently and random.random() < 0.7: # 70% chance node might return
self.departed_nodes[node.node_uuid] = {
"last_seen": datetime.now(),
"characteristics": {
"base_load": node.base_load,
"base_memory": node.base_memory,
"load_variance": node.load_variance,
"memory_variance": node.memory_variance,
},
}
logger.info(
f"⏸️ Node {node_id} ({node.node_uuid[:8]}...) departing temporarily"
)
else:
logger.info(
f"❌ Node {node_id} ({node.node_uuid[:8]}...) leaving permanently"
)
node.stop()
thread.join(timeout=5.0)
del self.active_nodes[node_id]
del self.node_threads[node_id]
return True
def adjust_node_count(self):
"""Randomly adjust the number of active nodes within the specified range."""
current_count = len(self.active_nodes)
# Decide on target count
target_count = random.randint(self.min_nodes, self.max_nodes)
if target_count == current_count:
logger.debug(f"Node count staying at {current_count}")
return
logger.info(f"🔄 Adjusting node count: {current_count}{target_count}")
if target_count > current_count:
# Add nodes
nodes_to_add = target_count - current_count
for _ in range(nodes_to_add):
node = self.create_new_node(return_existing=True)
if node:
self.start_node(node)
time.sleep(random.uniform(1, 3)) # Stagger starts
elif target_count < current_count:
# Remove nodes
nodes_to_remove = current_count - target_count
active_node_ids = list(self.active_nodes.keys())
nodes_to_stop = random.sample(active_node_ids, nodes_to_remove)
for node_id in nodes_to_stop:
# 20% chance node leaves permanently
permanently = random.random() < 0.2
self.stop_node(node_id, permanently=permanently)
time.sleep(random.uniform(0.5, 2)) # Stagger stops
def manage_node_lifecycle(self):
"""Main loop for managing node lifecycle changes."""
logger.info("🚀 Starting dynamic node lifecycle management")
# Start with minimum nodes
logger.info(f"Initializing with {self.min_nodes} nodes...")
for _ in range(self.min_nodes):
node = self.create_new_node()
if node:
self.start_node(node)
time.sleep(1) # Brief delay between starts
# Main management loop
while self.running:
try:
# Wait for the change interval (with some randomness)
variance = random.randint(
-self.lifecycle_variance, self.lifecycle_variance
)
sleep_time = max(
10, self.change_interval + variance
) # Minimum 10 seconds
logger.debug(
f"Waiting {sleep_time} seconds until next potential change..."
)
time.sleep(sleep_time)
if not self.running:
break
# Randomly decide if we should make a change (70% chance)
if random.random() < 0.7:
self.adjust_node_count()
else:
logger.debug("Skipping this change cycle")
except KeyboardInterrupt:
logger.info("Node lifecycle manager received interrupt signal")
break
except Exception as e:
logger.error(f"Error in node lifecycle management: {e}", exc_info=True)
time.sleep(5) # Brief pause before continuing
def start_dynamic_management(self):
"""Start the dynamic node management system."""
if not self.check_ip_availability(self.max_nodes * 2):
logger.error(
"Cannot start dynamic management due to IP availability issues"
)
return False
self.running = True
self.manager_thread = threading.Thread(
target=self.manage_node_lifecycle, name="NodeLifecycleManager"
)
self.manager_thread.daemon = True
self.manager_thread.start()
return True
def stop_all_nodes(self):
"""Stop all nodes and the management system."""
logger.info("🛑 Stopping dynamic node management...")
self.running = False
# Stop the manager thread
if self.manager_thread and self.manager_thread.is_alive():
self.manager_thread.join(timeout=5.0)
# Stop all active nodes
for node_id in list(self.active_nodes.keys()):
self.stop_node(node_id, permanently=True)
logger.info("All nodes stopped")
def print_status(self):
"""Print current status of the dynamic system."""
active_count = len(self.active_nodes)
departed_count = len(self.departed_nodes)
logger.info(f"=== Dynamic Multi-Node Status ===")
logger.info(
f"Active nodes: {active_count} (range: {self.min_nodes}-{self.max_nodes})"
)
logger.info(f"Departed nodes (may return): {departed_count}")
for node_id, node in self.active_nodes.items():
status = "returning" if node.is_persistent else "new"
logger.info(
f" Node {node_id}: {status}, IP={node.source_ip}, UUID={node.node_uuid[:8]}..., Uptime={node.uptime_seconds}s"
)
if departed_count > 0:
logger.info("Departed nodes that might return:")
for uuid_val in list(self.departed_nodes.keys())[:5]: # Show first 5
logger.info(f" UUID={uuid_val[:8]}...")
def main():
parser = argparse.ArgumentParser(
description="Dynamic multi-node test client with fluctuating node counts"
)
parser.add_argument(
"--min-nodes",
type=int,
default=DYNAMIC_MIN_NODES,
help=f"Minimum number of nodes (default: {DYNAMIC_MIN_NODES})",
)
parser.add_argument(
"--max-nodes",
type=int,
default=DYNAMIC_MAX_NODES,
help=f"Maximum number of nodes (default: {DYNAMIC_MAX_NODES})",
)
parser.add_argument(
"--change-interval",
type=int,
default=NODE_CHANGE_INTERVAL,
help=f"Average seconds between node changes (default: {NODE_CHANGE_INTERVAL})",
)
parser.add_argument(
"--lifecycle-variance",
type=int,
default=NODE_LIFECYCLE_VARIANCE,
help=f"Random variance in change timing (default: {NODE_LIFECYCLE_VARIANCE})",
)
parser.add_argument(
"--interval",
type=int,
default=UPDATE_INTERVAL_SECONDS,
help=f"Update interval in seconds (default: {UPDATE_INTERVAL_SECONDS})",
)
parser.add_argument(
"--server",
type=str,
default=SERVER_BASE_URL,
help=f"Server URL (default: {SERVER_BASE_URL})",
)
parser.add_argument(
"--service-uuid",
type=str,
default=TARGET_SERVICE_UUID,
help="Target service UUID",
)
parser.add_argument(
"--ip-base",
type=int,
default=LOOPBACK_IP_BASE,
help=f"Starting IP for 127.0.0.X (default: {LOOPBACK_IP_BASE})",
)
parser.add_argument(
"--setup-help",
action="store_true",
help="Show commands to set up loopback IP addresses",
)
parser.add_argument(
"--verbose", "-v", action="store_true", help="Enable verbose logging"
)
# Keep the old --nodes argument for compatibility, but make it set max nodes
parser.add_argument(
"--nodes",
type=int,
help="Set max nodes (compatibility mode, same as --max-nodes)",
)
args = parser.parse_args()
# Handle compatibility
if args.nodes:
args.max_nodes = args.nodes
if (
args.min_nodes == DYNAMIC_MIN_NODES
): # Only adjust min if it wasn't explicitly set
args.min_nodes = max(
1, args.nodes - 2
) # Set min to a reasonable lower bound
min_nodes = args.min_nodes
max_nodes = args.max_nodes
change_interval = args.change_interval
lifecycle_variance = args.lifecycle_variance
update_interval = args.interval
server_url = args.server
service_uuid = args.service_uuid
ip_base = args.ip_base
if args.verbose:
logging.getLogger().setLevel(logging.DEBUG)
# Validate configuration
if min_nodes >= max_nodes:
logger.error("Minimum nodes must be less than maximum nodes")
return
if max_nodes < 1:
logger.error("Maximum nodes must be at least 1")
return
if args.setup_help:
# Create a temporary manager just to show setup commands
temp_manager = DynamicMultiNodeManager(
min_nodes,
max_nodes,
server_url,
service_uuid,
update_interval,
ip_base,
change_interval,
lifecycle_variance,
)
temp_manager.show_setup_commands()
return
if service_uuid == "REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID":
logger.error("=" * 60)
logger.error("ERROR: TARGET_SERVICE_UUID is not set correctly!")
logger.error(
"Please set it via --service-uuid argument or TARGET_SERVICE_UUID environment variable."
)
logger.error("=" * 60)
return
logger.info("=" * 60)
logger.info("Dynamic Multi-Node Test Client Configuration:")
logger.info(f" Node range: {min_nodes} - {max_nodes} nodes")
logger.info(
f" Change interval: {change_interval}s (±{lifecycle_variance}s variance)"
)
logger.info(f" Update interval: {update_interval} seconds per node")
logger.info(f" Server URL: {server_url}")
logger.info(f" Target Service UUID: {service_uuid}")
logger.info(
f" IP range: 127.0.0.{ip_base} - 127.0.0.{ip_base + max_nodes * 2 - 1}"
)
logger.info("=" * 60)
# Create and start the dynamic multi-node manager
manager = DynamicMultiNodeManager(
min_nodes,
max_nodes,
server_url,
service_uuid,
update_interval,
ip_base,
change_interval,
lifecycle_variance,
)
try:
if not manager.start_dynamic_management():
logger.error("Failed to start dynamic node management")
return
# Main monitoring loop
status_interval = 30 # Print status every 30 seconds
last_status_time = time.time()
logger.info(
"🎯 Dynamic node management started! Nodes will fluctuate between {} and {} over time.".format(
min_nodes, max_nodes
)
)
logger.info("Press Ctrl+C to stop...")
while True:
time.sleep(5) # Check every 5 seconds
current_time = time.time()
if current_time - last_status_time >= status_interval:
manager.print_status()
last_status_time = current_time
except KeyboardInterrupt:
logger.info("Received interrupt signal, shutting down...")
except Exception as e:
logger.error(f"Unexpected error in main: {e}", exc_info=True)
finally:
manager.stop_all_nodes()
if __name__ == "__main__":
main()
-548
View File
@@ -1,548 +0,0 @@
import os
import uuid
import time
import requests
import random
import json
import logging
import threading
import socket
from datetime import datetime, timezone
from concurrent.futures import ThreadPoolExecutor
import argparse
from requests.adapters import HTTPAdapter
from urllib3.util.connection import create_connection
# --- Multi-Node Client Configuration ---
TARGET_SERVICE_UUID = os.environ.get(
"TARGET_SERVICE_UUID", "ab73d00a-8169-46bb-997d-f13e5f760973"
)
SERVER_BASE_URL = os.environ.get("SERVER_URL", "http://localhost:8000")
UPDATE_INTERVAL_SECONDS = int(os.environ.get("UPDATE_INTERVAL_SECONDS", 5))
NUM_NODES = int(os.environ.get("NUM_NODES", 3))
# Base IP for loopback binding (127.0.0.x where x starts from this base)
LOOPBACK_IP_BASE = int(os.environ.get("LOOPBACK_IP_BASE", 2)) # Start from 127.0.0.2
# --- Logging Configuration ---
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s - %(name)s - %(levelname)s - [%(thread)d] - %(message)s",
)
logger = logging.getLogger("MultiNodeClient")
# --- Custom HTTP Adapter for Source IP Binding ---
class SourceIPHTTPAdapter(HTTPAdapter):
def __init__(self, source_ip, *args, **kwargs):
self.source_ip = source_ip
super().__init__(*args, **kwargs)
def init_poolmanager(self, *args, **kwargs):
# Override the socket creation to bind to specific source IP
def custom_create_connection(
address,
timeout=socket._GLOBAL_DEFAULT_TIMEOUT,
source_address=None,
socket_options=None,
):
# Force our custom source address
return create_connection(
address,
timeout,
source_address=(self.source_ip, 0), # 0 = any available port
socket_options=socket_options,
)
# Monkey patch the connection creation
original_create_connection = socket.create_connection
socket.create_connection = custom_create_connection
try:
result = super().init_poolmanager(*args, **kwargs)
finally:
# Restore original function
socket.create_connection = original_create_connection
return result
# --- Enhanced Node Class with IP Binding ---
class SimulatedNode:
def __init__(
self,
node_id: int,
total_nodes: int,
server_url: str,
service_uuid: str,
update_interval: int,
ip_base: int,
):
self.node_id = node_id
self.node_uuid = str(uuid.uuid4())
self.server_url = server_url # Store server URL
self.service_uuid = service_uuid # Store service UUID
self.update_interval = update_interval # Store update interval
self.uptime_seconds = 0
self.known_peers = {}
self.total_nodes = total_nodes
self.running = False
# Assign unique loopback IP to this node using the passed ip_base
self.source_ip = f"127.0.0.{ip_base + node_id - 1}"
# Create requests session with custom adapter for IP binding
self.session = requests.Session()
adapter = SourceIPHTTPAdapter(self.source_ip)
self.session.mount("http://", adapter)
self.session.mount("https://", adapter)
# Each node gets slightly different characteristics
self.base_load = random.uniform(0.2, 1.0)
self.base_memory = random.uniform(40.0, 70.0)
self.load_variance = random.uniform(0.1, 0.5)
self.memory_variance = random.uniform(5.0, 15.0)
# Some nodes might be "problematic" (higher load/memory)
if random.random() < 0.2: # 20% chance of being a "problematic" node
self.base_load *= 2.0
self.base_memory += 20.0
logger.info(
f"Node {self.node_id} ({self.node_uuid[:8]}) will simulate high resource usage (IP: {self.source_ip})"
)
logger.info(f"Node {self.node_id} will bind to source IP: {self.source_ip}")
def test_ip_binding(self):
"""Test if we can bind to the assigned IP address."""
try:
# Try to create a socket and bind to the IP
test_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
test_socket.bind((self.source_ip, 0)) # Bind to any available port
test_socket.close()
logger.debug(
f"Node {self.node_id} successfully tested binding to {self.source_ip}"
)
return True
except OSError as e:
logger.error(f"Node {self.node_id} cannot bind to {self.source_ip}: {e}")
logger.error(
"Make sure the IP address is available on your loopback interface."
)
logger.error(
"You might need to add it with: sudo ifconfig lo0 alias {self.source_ip} (macOS)"
)
logger.error("Or: sudo ip addr add {self.source_ip}/8 dev lo (Linux)")
return False
def generate_node_status_data(self):
"""Generates simulated node status metrics with per-node characteristics."""
self.uptime_seconds += self.update_interval + random.randint(-1, 2)
# Generate load with some randomness but consistent per-node baseline
load_1min = max(
0.1,
self.base_load + random.uniform(-self.load_variance, self.load_variance),
)
load_5min = max(0.1, load_1min * random.uniform(0.8, 1.0))
load_15min = max(0.1, load_5min * random.uniform(0.8, 1.0))
load_avg = [round(load_1min, 2), round(load_5min, 2), round(load_15min, 2)]
# Generate memory usage with baseline + variance
memory_usage = max(
10.0,
min(
95.0,
self.base_memory
+ random.uniform(-self.memory_variance, self.memory_variance),
),
)
return {
"uptime_seconds": self.uptime_seconds,
"load_avg": load_avg,
"memory_usage_percent": round(memory_usage, 2),
}
def generate_ping_data(self):
"""Generates simulated ping latencies to known peers."""
pings = {}
# Ping to self (loopback)
pings[self.node_uuid] = round(random.uniform(0.1, 1.5), 2)
# Ping to known peers
for peer_uuid in self.known_peers.keys():
if peer_uuid != self.node_uuid:
# Simulate network latency with some consistency per peer
base_latency = random.uniform(5.0, 100.0)
variation = random.uniform(-10.0, 10.0)
latency = max(0.1, base_latency + variation)
pings[peer_uuid] = round(latency, 2)
return pings
def send_update(self):
"""Sends a single status update to the server using bound IP."""
try:
status_data = self.generate_node_status_data()
ping_data = self.generate_ping_data()
payload = {
"node": self.node_uuid,
"timestamp": datetime.now(timezone.utc).isoformat(),
"status": status_data,
"pings": ping_data,
}
endpoint_url = f"{self.server_url}/{self.service_uuid}/{self.node_uuid}/"
logger.debug(
f"Node {self.node_id} ({self.source_ip}) sending update. "
f"Uptime: {status_data['uptime_seconds']}s, "
f"Load: {status_data['load_avg'][0]}, Memory: {status_data['memory_usage_percent']}%, "
f"Pings: {len(ping_data)}"
)
# Use the custom session with IP binding
response = self.session.put(endpoint_url, json=payload, timeout=10)
if response.status_code == 200:
response_data = response.json()
if "peers" in response_data and isinstance(
response_data["peers"], dict
):
new_peers = {k: v for k, v in response_data["peers"].items()}
# Log new peer discoveries
newly_discovered = set(new_peers.keys()) - set(
self.known_peers.keys()
)
if newly_discovered:
logger.info(
f"Node {self.node_id} ({self.source_ip}) discovered {len(newly_discovered)} new peer(s)"
)
self.known_peers = new_peers
if (
len(newly_discovered) > 0
or len(self.known_peers) != self.total_nodes - 1
):
logger.debug(
f"Node {self.node_id} ({self.source_ip}) knows {len(self.known_peers)} peers "
f"(expected {self.total_nodes - 1})"
)
return True
else:
logger.error(
f"Node {self.node_id} ({self.source_ip}) failed to send update. "
f"Status: {response.status_code}, Response: {response.text}"
)
return False
except requests.exceptions.Timeout:
logger.error(f"Node {self.node_id} ({self.source_ip}) request timed out")
return False
except requests.exceptions.ConnectionError as e:
logger.error(
f"Node {self.node_id} ({self.source_ip}) connection error: {e}"
)
return False
except Exception as e:
logger.error(
f"Node {self.node_id} ({self.source_ip}) unexpected error: {e}"
)
return False
def run(self):
"""Main loop for this simulated node."""
# Test IP binding before starting
if not self.test_ip_binding():
logger.error(f"Node {self.node_id} cannot start due to IP binding failure")
return
self.running = True
logger.info(
f"Starting Node {self.node_id} with UUID: {self.node_uuid} (IP: {self.source_ip})"
)
# Add some initial delay to stagger node starts
initial_delay = self.node_id * 0.5
time.sleep(initial_delay)
consecutive_failures = 0
while self.running:
try:
success = self.send_update()
if success:
consecutive_failures = 0
else:
consecutive_failures += 1
if consecutive_failures >= 3:
logger.warning(
f"Node {self.node_id} ({self.source_ip}) has failed {consecutive_failures} consecutive updates"
)
# Add some jitter to prevent thundering herd
jitter = random.uniform(-1.0, 1.0)
sleep_time = max(1.0, self.update_interval + jitter)
time.sleep(sleep_time)
except KeyboardInterrupt:
logger.info(
f"Node {self.node_id} ({self.source_ip}) received interrupt signal"
)
break
except Exception as e:
logger.error(
f"Node {self.node_id} ({self.source_ip}) unexpected error in main loop: {e}"
)
time.sleep(self.update_interval)
logger.info(f"Node {self.node_id} ({self.source_ip}) stopped")
def stop(self):
"""Stop the node."""
self.running = False
# --- Multi-Node Manager ---
class MultiNodeManager:
def __init__(
self,
num_nodes: int,
server_url: str,
service_uuid: str,
update_interval: int,
ip_base: int,
):
self.num_nodes = num_nodes
self.server_url = server_url
self.service_uuid = service_uuid
self.update_interval = update_interval
self.ip_base = ip_base
self.nodes = []
self.threads = []
self.running = False
# Create simulated nodes
for i in range(num_nodes):
node = SimulatedNode(
i + 1, num_nodes, server_url, service_uuid, update_interval, ip_base
)
self.nodes.append(node)
def check_ip_availability(self):
"""Check if all required IP addresses are available."""
logger.info("Checking IP address availability...")
all_available = True
for node in self.nodes:
if not node.test_ip_binding():
all_available = False
if not all_available:
logger.error(
"Some IP addresses are not available. See individual node errors above."
)
logger.info("To add loopback IP addresses:")
logger.info(" Linux: sudo ip addr add 127.0.0.X/8 dev lo")
logger.info(" macOS: sudo ifconfig lo0 alias 127.0.0.X")
# Use self.ip_base for the range
logger.info(
f" Where X ranges from {self.ip_base} to {self.ip_base + self.num_nodes - 1}"
)
return False
logger.info("All IP addresses are available!")
return True
def start_all_nodes(self):
"""Start all simulated nodes in separate threads."""
if not self.check_ip_availability():
return False
logger.info(
f"Starting {self.num_nodes} simulated nodes with unique IP addresses..."
)
self.running = True
for node in self.nodes:
thread = threading.Thread(target=node.run, name=f"Node-{node.node_id}")
thread.daemon = True
self.threads.append(thread)
thread.start()
logger.info(f"All {self.num_nodes} nodes started")
return True
def stop_all_nodes(self):
"""Stop all simulated nodes."""
logger.info("Stopping all nodes...")
self.running = False
for node in self.nodes:
node.stop()
# Wait for threads to finish
for thread in self.threads:
thread.join(timeout=5.0)
logger.info("All nodes stopped")
def print_status(self):
"""Print current status of all nodes."""
logger.info(f"=== Multi-Node Status ({self.num_nodes} nodes) ===")
for node in self.nodes:
logger.info(
f"Node {node.node_id}: IP={node.source_ip}, UUID={node.node_uuid[:8]}..., "
f"Uptime={node.uptime_seconds}s, Peers={len(node.known_peers)}"
)
def setup_loopback_ips(num_nodes, base_ip):
"""Helper function to show commands for setting up loopback IPs."""
logger.info("=== Loopback IP Setup Commands ===")
logger.info("Run these commands to add the required loopback IP addresses:")
logger.info("")
# Detect OS and show appropriate commands
import platform
system = platform.system().lower()
for i in range(num_nodes):
ip = f"127.0.0.{base_ip + i}"
if system == "linux":
logger.info(f"sudo ip addr add {ip}/8 dev lo")
elif system == "darwin": # macOS
logger.info(f"sudo ifconfig lo0 alias {ip}")
else:
logger.info(f"Add {ip} to loopback interface (OS: {system})")
logger.info("")
logger.info("To remove them later:")
for i in range(num_nodes):
ip = f"127.0.0.{base_ip + i}"
if system == "linux":
logger.info(f"sudo ip addr del {ip}/8 dev lo")
elif system == "darwin": # macOS
logger.info(f"sudo ifconfig lo0 -alias {ip}")
logger.info("=" * 40)
def main():
parser = argparse.ArgumentParser(
description="Multi-node test client with unique IP binding"
)
parser.add_argument(
"--nodes",
type=int,
default=NUM_NODES,
help=f"Number of simulated nodes (default: {NUM_NODES})",
)
parser.add_argument(
"--interval",
type=int,
default=UPDATE_INTERVAL_SECONDS,
help=f"Update interval in seconds (default: {UPDATE_INTERVAL_SECONDS})",
)
parser.add_argument(
"--server",
type=str,
default=SERVER_BASE_URL,
help=f"Server URL (default: {SERVER_BASE_URL})",
)
parser.add_argument(
"--service-uuid",
type=str,
default=TARGET_SERVICE_UUID,
help="Target service UUID",
)
parser.add_argument(
"--ip-base",
type=int,
default=LOOPBACK_IP_BASE,
help=f"Starting IP for 127.0.0.X (default: {LOOPBACK_IP_BASE})",
)
parser.add_argument(
"--setup-help",
action="store_true",
help="Show commands to set up loopback IP addresses",
)
parser.add_argument(
"--verbose", "-v", action="store_true", help="Enable verbose logging"
)
args = parser.parse_args()
num_nodes = args.nodes
update_interval = args.interval
server_url = args.server
service_uuid = args.service_uuid
ip_base = args.ip_base
if args.verbose:
logging.getLogger().setLevel(logging.DEBUG)
if args.setup_help:
setup_loopback_ips(num_nodes, ip_base)
return
# Validate configuration
if service_uuid == "REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID":
logger.error("=" * 60)
logger.error("ERROR: TARGET_SERVICE_UUID is not set correctly!")
logger.error(
"Please set it via --service-uuid argument or TARGET_SERVICE_UUID environment variable."
)
logger.error(
"You can find the server's UUID by running main.py and checking its console output"
)
logger.error("or by visiting the server's root endpoint in your browser.")
logger.error("=" * 60)
return
logger.info("=" * 60)
logger.info("Multi-Node Test Client Configuration:")
logger.info(f" Number of nodes: {num_nodes}")
logger.info(f" Update interval: {update_interval} seconds")
logger.info(f" Server URL: {server_url}")
logger.info(f" Target Service UUID: {service_uuid}")
logger.info(f" IP range: 127.0.0.{ip_base} - 127.0.0.{ip_base + num_nodes - 1}")
logger.info("=" * 60)
# Create and start the multi-node manager
manager = MultiNodeManager(
num_nodes, server_url, service_uuid, update_interval, ip_base
)
try:
if not manager.start_all_nodes():
logger.error("Failed to start nodes. Check IP availability.")
setup_loopback_ips(num_nodes, ip_base)
return
# Main monitoring loop
while True:
time.sleep(30) # Print status every 30 seconds
manager.print_status()
except KeyboardInterrupt:
logger.info("Received interrupt signal, shutting down...")
except Exception as e:
logger.error(f"Unexpected error in main: {e}", exc_info=True)
finally:
manager.stop_all_nodes()
if __name__ == "__main__":
main()
-163
View File
@@ -1,163 +0,0 @@
import os
import uuid
import time
import requests
import random
import json
import logging
from datetime import datetime, timezone
# --- Client Configuration ---
# The UUID of THIS client node. Generated on startup.
# Can be overridden by an environment variable for persistent client identity.
NODE_UUID = os.environ.get("NODE_UUID", str(uuid.uuid4()))
# The UUID of the target monitoring service (the main.py server).
# IMPORTANT: This MUST match the SERVICE_UUID of your running FastAPI server.
# You can get this from the server's initial console output or by accessing its root endpoint ('/').
# Replace the placeholder string below with your actual server's SERVICE_UUID.
# For example: TARGET_SERVICE_UUID = "a1b2c3d4-e5f6-7890-1234-567890abcdef"
TARGET_SERVICE_UUID = os.environ.get(
"TARGET_SERVICE_UUID", "REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID"
)
# The base URL of the FastAPI monitoring service
SERVER_BASE_URL = os.environ.get("SERVER_URL", "http://localhost:8000")
# How often to send status updates (in seconds)
UPDATE_INTERVAL_SECONDS = int(os.environ.get("UPDATE_INTERVAL_SECONDS", 5))
# --- Logging Configuration ---
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger("NodeClient")
# --- Global state for simulation ---
uptime_seconds = 0
# Dictionary to store UUIDs of other nodes received from the server
# Format: { "node_uuid_str": { "last_seen": "iso_timestamp", "ip": "..." } }
known_peers = {}
# --- Data Generation Functions ---
def generate_node_status_data():
"""Generates simulated node status metrics."""
global uptime_seconds
uptime_seconds += UPDATE_INTERVAL_SECONDS + random.randint(0, 2) # Simulate slight variation
# Simulate load average (3 values: 1-min, 5-min, 15-min)
# Load averages will fluctuate.
load_avg = [
round(random.uniform(0.1, 2.0), 2),
round(random.uniform(0.1, 1.8), 2),
round(random.uniform(0.1, 1.5), 2)
]
# Simulate memory usage percentage
memory_usage_percent = round(random.uniform(30.0, 90.0), 2)
return {
"uptime_seconds": uptime_seconds,
"load_avg": load_avg,
"memory_usage_percent": memory_usage_percent
}
def generate_ping_data():
"""Generates simulated ping latencies to known peers."""
pings = {}
# Simulate ping to self (loopback) - always very low latency
pings[str(NODE_UUID)] = round(random.uniform(0.1, 1.0), 2)
# Simulate pings to other known peers
for peer_uuid in known_peers.keys():
if peer_uuid != str(NODE_UUID): # Don't ping self twice
# Varying latency for external peers
pings[peer_uuid] = round(random.uniform(10.0, 200.0), 2)
return pings
# --- Main Client Logic ---
def run_client():
global known_peers
logger.info(f"Starting Node Client {NODE_UUID}")
logger.info(f"Target Service UUID: {TARGET_SERVICE_UUID}")
logger.info(f"Server URL: {SERVER_BASE_URL}")
logger.info(f"Update Interval: {UPDATE_INTERVAL_SECONDS} seconds")
if TARGET_SERVICE_UUID == "REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID":
logger.error("-" * 50)
logger.error("ERROR: TARGET_SERVICE_UUID is not set correctly!")
logger.error("Please replace 'REPLACE_ME_WITH_YOUR_SERVER_SERVICE_UUID' in client.py")
logger.error("or set the environment variable TARGET_SERVICE_UUID.")
logger.error("You can find the server's UUID by running main.py and checking its console output")
logger.error("or by visiting 'http://localhost:8000/' in your browser.")
logger.error("-" * 50)
return
while True:
try:
# 1. Generate status data
status_data = generate_node_status_data()
ping_data = generate_ping_data()
# 2. Construct the payload matching the StatusUpdate model
# Use datetime.now(timezone.utc) for timezone-aware UTC timestamp
payload = {
"node": str(NODE_UUID),
"timestamp": datetime.now(timezone.utc).isoformat(),
"status": status_data,
"pings": ping_data
}
# 3. Define the endpoint URL
endpoint_url = f"{SERVER_BASE_URL}/{TARGET_SERVICE_UUID}/{NODE_UUID}/"
# 4. Send the PUT request
logger.info(f"Sending update to {endpoint_url}. Uptime: {status_data['uptime_seconds']}s, Load: {status_data['load_avg']}, Pings: {len(ping_data)}")
response = requests.put(endpoint_url, json=payload, timeout=10) # 10-second timeout
# 5. Process the response
if response.status_code == 200:
response_data = response.json()
logger.info(f"Successfully sent update. Server message: '{response_data.get('message')}'")
if "peers" in response_data and isinstance(response_data["peers"], dict):
# Update known_peers, converting keys to strings from JSON
new_peers = {k: v for k, v in response_data["peers"].items()}
# Log if new peers are discovered
newly_discovered = set(new_peers.keys()) - set(known_peers.keys())
if newly_discovered:
logger.info(f"Discovered new peer(s): {', '.join(newly_discovered)}")
known_peers = new_peers
logger.info(f"Total known peers (including self if returned by server): {len(known_peers)}")
else:
logger.warning("Server response did not contain a valid 'peers' field or it was empty.")
else:
logger.error(f"Failed to send update. Status code: {response.status_code}, Response: {response.text}")
if response.status_code == 404:
logger.error("Hint: The TARGET_SERVICE_UUID might be incorrect, or the server isn't running at this endpoint.")
elif response.status_code == 422: # Pydantic validation error
logger.error(f"Server validation error (422 Unprocessable Entity): {response.json()}")
except requests.exceptions.Timeout:
logger.error(f"Request timed out after {10} seconds. Is the server running and responsive?")
except requests.exceptions.ConnectionError as e:
logger.error(f"Connection error: {e}. Is the server running at {SERVER_BASE_URL}?")
except requests.exceptions.RequestException as e:
logger.error(f"An unexpected request error occurred: {e}", exc_info=True)
except json.JSONDecodeError:
logger.error(f"Failed to decode JSON response: {response.text}. Is the server returning valid JSON?")
except Exception as e:
logger.error(f"An unexpected error occurred in the client loop: {e}", exc_info=True)
# 6. Wait for the next update
time.sleep(UPDATE_INTERVAL_SECONDS)
if __name__ == "__main__":
run_client()