diff --git a/client.py b/client.py index f7ef287..04c8553 100644 --- a/client.py +++ b/client.py @@ -8,13 +8,17 @@ import json import logging from datetime import datetime, timezone import platform -import socket # For getting local IP +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' + import psutil # For system metrics + 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,84 +29,95 @@ 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") # --- Global state --- -uptime_seconds = 0 # Will be updated by psutil.boot_time() or incremented +uptime_seconds = 0 # Will be updated by psutil.boot_time() or incremented # known_peers will store { "node_uuid_str": "ip_address_str" } -known_peers: dict[str, str] = {} +known_peers: dict[str, str] = {} # Determine local IP for self-pinging and reporting to server -LOCAL_IP = "127.0.0.1" # Default fallback +LOCAL_IP = "127.0.0.1" # Default fallback try: s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) - s.connect(("8.8.8.8", 80)) # Connect to an external host (doesn't send data) + s.connect(("8.8.8.8", 80)) # Connect to an external host (doesn't send data) LOCAL_IP = s.getsockname()[0] s.close() 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 temp_peers = {} for k, v in loaded_data.items(): - if isinstance(v, str): # Already in {uuid: ip} format + 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.") - known_peers = {} # Reset if file is corrupt + 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.""" global uptime_seconds - + # Uptime # psutil.boot_time() returns a timestamp in seconds since epoch uptime_seconds = int(time.time() - psutil.boot_time()) @@ -112,14 +127,16 @@ 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 + 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 - load_avg = [cpu_percent, cpu_percent * 0.9, cpu_percent * 0.8] # Simulate decay + 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}") # Memory Usage @@ -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.""" @@ -145,7 +163,7 @@ def perform_pings(targets: dict[str, str]) -> dict[str, float]: # pythonping returns response_time in seconds, convert to milliseconds pings_results[str(NODE_UUID)] = round(response_list.rtt_avg_ms, 2) else: - pings_results[str(NODE_UUID)] = -1.0 # Indicate failure + pings_results[str(NODE_UUID)] = -1.0 # Indicate failure logger.debug(f"Ping to self ({LOCAL_IP}): {pings_results[str(NODE_UUID)]}ms") except Exception as e: logger.warning(f"Failed to ping self ({LOCAL_IP}): {e}") @@ -154,7 +172,7 @@ def perform_pings(targets: dict[str, str]) -> dict[str, float]: # Ping other known peers for peer_uuid, peer_ip in targets.items(): if peer_uuid == str(NODE_UUID): - continue # Already pinged self + continue # Already pinged self try: # Use a longer timeout for external pings @@ -162,14 +180,37 @@ def perform_pings(targets: dict[str, str]) -> dict[str, float]: if response_list.success: 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") + pings_results[peer_uuid] = -1.0 # Indicate failure + 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 @@ -198,74 +249,124 @@ def run_client(): try: # 1. Get real system metrics status_data = get_system_metrics() - + # 2. Perform pings to known peers (and self) ping_data = perform_pings(known_peers) # 3. Construct the payload payload = { "node": str(NODE_UUID), - "timestamp": datetime.now(timezone.utc).isoformat(), + "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 - # 6. Process the response + response = requests.put( + endpoint_url, + json=payload, + headers=headers, + timeout=15, + verify=SSL_VERIFY, + ) + + # 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')}'") - - if "peers" in response_data and isinstance(response_data["peers"], dict): + 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 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 + 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.") - elif response.status_code == 422: # Pydantic validation error - logger.error(f"Server validation error (422 Unprocessable Entity): {response.json()}") + 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 + 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() \ No newline at end of file + run_client() diff --git a/monitor-agent/agent-deploy.sh b/monitor-agent/agent-deploy.sh new file mode 100644 index 0000000..66f8128 --- /dev/null +++ b/monitor-agent/agent-deploy.sh @@ -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 < /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 < /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 "------------------------------------------------" diff --git a/monitor-agent/go.mod b/monitor-agent/go.mod new file mode 100644 index 0000000..207e160 --- /dev/null +++ b/monitor-agent/go.mod @@ -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 +) diff --git a/monitor-agent/main.go b/monitor-agent/main.go new file mode 100644 index 0000000..5c2338b --- /dev/null +++ b/monitor-agent/main.go @@ -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)) + } + } +} diff --git a/test-client-multi-node-flux.py b/test-client-multi-node-flux.py deleted file mode 100644 index bb8d241..0000000 --- a/test-client-multi-node-flux.py +++ /dev/null @@ -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() diff --git a/test-client-multi-node.py b/test-client-multi-node.py deleted file mode 100644 index 61d4567..0000000 --- a/test-client-multi-node.py +++ /dev/null @@ -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() diff --git a/test-client.py b/test-client.py deleted file mode 100644 index 6e40753..0000000 --- a/test-client.py +++ /dev/null @@ -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() \ No newline at end of file