From 777993dfbbcbc539cbb94278b3c00a73fa9321b4 Mon Sep 17 00:00:00 2001 From: Diego Freitas Date: Mon, 15 Dec 2025 10:42:45 -0300 Subject: [PATCH] ajustado asincronia ao atualizar saude do IPB --- .../workers/health_worker/modulos/ponte_ip.py | 624 ++++++++++++++++++ 1 file changed, 624 insertions(+) create mode 100644 AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/modulos/ponte_ip.py diff --git a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/modulos/ponte_ip.py b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/modulos/ponte_ip.py new file mode 100644 index 000000000..4e3b7380e --- /dev/null +++ b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/modulos/ponte_ip.py @@ -0,0 +1,624 @@ +import time +import threading +import statistics +import socket +import psutil +import subprocess +from collections import deque +import paho.mqtt.client as mqtt +from shared.contexto_global_redis import ContextoGlobalRedis +from shared.enums import StatusModulo, T_Code +from health_worker.modulos.base import ModuloDiagnosticoBase + +class ModuloIPBribge(ModuloDiagnosticoBase): + def __init__(self, window_size: int = 5): + self.t_code = T_Code.Ipb + self.nome = "IP_Brigde" + self.timeout = 5 + + self._thread_saude = None + self._running = False + + self.nic_name = None + self.window_size = window_size + + self.rtts = deque(maxlen=window_size) # ping válidos (ms) + self.timeouts = deque(maxlen=window_size) # True/False por tentativa + self.loses = deque(maxlen=window_size) # True/False por tentativa + + self.last_nic_counters = None + self.nic_error_rates = deque(maxlen=window_size) + + self.last_bw_counters = None + self.last_bw_ts = None + self.bw_rx_mbps = deque(maxlen=window_size) + self.bw_tx_mbps = deque(maxlen=window_size) + + self.last_heartbeat_ts = time.time() # atualize isso de fora + + self.rover_id = None + self._mqtt_conectado = False + self.sub = False + self._start_mqtt_heartbeat() + + + def atualizar_saude(self): + if self._running: + return + self._running = True + self._thread_saude = threading.Thread(target=self._job_saude, daemon=True) + self._thread_saude.start() + + def _job_saude(self): + try: + self.atualizar_saude_interno() + except Exception as e: + print(f"Erro ao atualizar saude: {e}") + finally: + self._running = False + + def _start_mqtt_heartbeat(self): + try: + if self._mqtt_conectado: + return + base_ip = self.get_base_ip() + if base_ip is None: + self._mqtt_conectado = False + return + self._mqtt = mqtt.Client() + self._mqtt.on_connect = self._on_connect + self._mqtt.on_message = self._on_message + self._mqtt.on_disconnect = self._on_disconnect + self._mqtt.connect(base_ip, 1883, 60) + self._mqtt.loop_start() + except: + pass + #print("Erro ao se conectar no broker mqtt") + + def _on_connect(self, client, userdata, flags, rc): + if rc == 0: + print("[HEARTBEAT] MQTT conectado") + self._mqtt_conectado = True + + self._reset_janelas() + + self.rover_id = ContextoGlobalRedis.get_equipamento().get("serial_number") + + if self.sub == False: + topic = f"agrobot/v1/rover/{self.rover_id}/heartbeat" + client.subscribe(topic) + print("[HEARTBEAT] Subscribado em:", topic) + self.sub = True + else: + print("Erro ao conectar MQTT:", rc) + + def _on_disconnect(self, client, userdata, rc): + print("[HEARTBEAT] MQTT desconectado!", rc) + self._mqtt_conectado = False + self.sub = False + + # se rc != 0 significa desconexão inesperada + if rc != 0: + print("[HEARTBEAT] Desconexão inesperada — pode ter perdido o link 900 MHz") + + # Marca como "sem heartbeat" instantâneo (opcional) + # Aqui podemos forçar atraso grande + self.last_heartbeat_ts = 0 + + def _on_message(self, client, userdata, msg): + topic = msg.topic + if topic == f"agrobot/v1/rover/{self.rover_id}/heartbeat": + self.last_heartbeat_ts = time.time() + #print("[HEARTBEAT] recebido") + + def _reset_janelas(self): + self.rtts.clear() + self.timeouts.clear() + self.loses.clear() + self.nic_error_rates.clear() + self.bw_rx_mbps.clear() + self.bw_tx_mbps.clear() + self.last_nic_counters = None + self.last_bw_counters = None + self.last_bw_ts = None + + def _status_por_score(self, score: float) -> StatusModulo: + if score >= 80: + return StatusModulo.OPERANTE + elif score >= 50: + return StatusModulo.ALERTA + else: + return StatusModulo.FALHA + + def get_base_ip(self): + return ContextoGlobalRedis.get_equipamento().get("base_ip") + + def descobrir_nic_para_base(self) -> str | None: + try: + base_ip = self.get_base_ip() + if base_ip is None: + return None + + # Abre socket UDP só para descobrir o IP local usado na rota + s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + s.connect((base_ip, 80)) # porta não importa + local_ip = s.getsockname()[0] + s.close() + except Exception as e: + print("[Ponte IP] Erro ao descobrir IP local da rota:", e) + return None + + # Agora, com o local_ip, encontrar qual NIC possui esse IP + addrs = psutil.net_if_addrs() + + for nic_name, addr_list in addrs.items(): + for addr in addr_list: + if addr.family == socket.AF_INET and addr.address == local_ip: + return nic_name + + return None + + def _update_nic_bandwidth(self): + if self.nic_name is None: + self.nic_name = self.descobrir_nic_para_base() + if self.nic_name is None: + return + + counters = psutil.net_io_counters(pernic=True).get(self.nic_name) + if not counters: + return + + now = time.time() + if self.last_bw_counters is not None and self.last_bw_ts is not None: + dt = now - self.last_bw_ts + if dt > 0: + delta_rx = counters.bytes_recv - self.last_bw_counters.bytes_recv + delta_tx = counters.bytes_sent - self.last_bw_counters.bytes_sent + + # bytes → bits → Mbit/s + rx_mbps = (delta_rx * 8.0) / (dt * 1_000_000.0) + tx_mbps = (delta_tx * 8.0) / (dt * 1_000_000.0) + + self.bw_rx_mbps.append(rx_mbps) + self.bw_tx_mbps.append(tx_mbps) + + self.last_bw_counters = counters + self.last_bw_ts = now + + def _ping_once(self): + """ + Retorna RTT em ms se ok, ou None se timeout. + Aqui exemplo pra Windows usando 'ping -n 1'. + Ajuste pro teu ambiente se precisar. + """ + try: + base_ip = self.get_base_ip() + if base_ip is None: + return None, None + + # timeout de 1000 ms + proc = subprocess.run( + ["ping", "-n", "5", "-w", "500", base_ip], + capture_output=True, text=True + ) + if proc.returncode != 0: + #print("timeout") + return None, 100 + + #print(proc.stdout) + + perda = None + media_tempo = None + for line in proc.stdout.splitlines(): + line = line.lower() + if ("%" in line and perda is None): + perda = line.split('(')[1].split('%')[0] + perda = float(perda) if perda.isdigit() else None + elif (" = " in line and "ms" in line and media_tempo is None): + media_tempo = line.split(' = ')[3].split('ms')[0] + media_tempo = float(media_tempo) if media_tempo.isdigit() else None + if perda is not None and media_tempo is not None: + break + return media_tempo, perda + except Exception: + print("erro timeout") + return None, 100 + + def _update_nic_errors(self): + if self.nic_name is None: + self.nic_name = self.descobrir_nic_para_base() + if self.nic_name is None: + return + + counters = psutil.net_io_counters(pernic=True).get(self.nic_name) + if not counters: + return + + if self.last_nic_counters is not None: + delta_err = (counters.errin - self.last_nic_counters.errin) + \ + (counters.errout - self.last_nic_counters.errout) + delta_pkts = (counters.packets_sent - self.last_nic_counters.packets_sent) + \ + (counters.packets_recv - self.last_nic_counters.packets_recv) + + if delta_pkts > 0: + # erros por mil pacotes + rate = (delta_err / delta_pkts) * 1000.0 + self.nic_error_rates.append(rate) + + self.last_nic_counters = counters + + def _metric_scores(self): + # valores brutos (default) + avg_rtt = None + loss_pct = 100.0 + jitter = 0.0 + err_rate = 0.0 + atraso = time.time() - self.last_heartbeat_ts + + # LATÊNCIA + if self.rtts: + avg_rtt = statistics.mean(self.rtts) + if avg_rtt <= 30: + latency_score = 100 + elif avg_rtt <= 50: + latency_score = 80 + elif avg_rtt <= 100: + latency_score = 50 + elif avg_rtt <= 200: + latency_score = 20 + else: + latency_score = 0 + else: + latency_score = 0 + + # TIMEOUT + if self.timeouts: + timeout_pct = (sum(self.timeouts) / len(self.timeouts)) * 100.0 + else: + timeout_pct = 100.0 + + if timeout_pct <= 5: + timeout_score = 100 + elif timeout_pct <= 10: + timeout_score = 80 + elif timeout_pct <= 15: + timeout_score = 50 + elif timeout_pct <= 30: + timeout_score = 20 + else: + timeout_score = 0 + + # LOSS + if self.loses: + loss_pct = (sum(self.loses) / len(self.loses)) + else: + loss_pct = 100.0 + + if loss_pct <= 2: + loss_score = 100 + elif loss_pct <= 5: + loss_score = 80 + elif loss_pct <= 8: + loss_score = 50 + elif loss_pct <= 15: + loss_score = 20 + else: + loss_score = 0 + + # JITTER + if len(self.rtts) > 1: + jitter = statistics.pstdev(self.rtts) + if jitter <= 10: + jitter_score = 100 + elif jitter <= 20: + jitter_score = 70 + elif jitter <= 40: + jitter_score = 30 + else: + jitter_score = 0 + else: + jitter = 0.0 + # Sem dados suficientes, melhor considerar "desconhecido" => score neutro/baixo + jitter_score = 0 + + # NIC ERRORS + if self.nic_error_rates: + err_rate = statistics.mean(self.nic_error_rates) + else: + err_rate = 0.0 + + if err_rate <= 0.1: + nic_score = 100 + elif err_rate <= 1: + nic_score = 70 + elif err_rate <= 5: + nic_score = 40 + else: + nic_score = 0 + + # HEARTBEAT (score usando "atraso") + if atraso <= 1.5: + hb_score = 100 + elif atraso <= 2: + hb_score = 80 + elif atraso <= 5: + hb_score = 30 + else: + hb_score = 0 + + # retorna scores + brutos + return ( + latency_score, timeout_score, loss_score, jitter_score, nic_score, hb_score, + avg_rtt, timeout_pct, loss_pct, jitter, err_rate, atraso + ) + + def atualizar_saude_interno(self): + try: + self._start_mqtt_heartbeat() + + SAUDE_MIN_ALERTA = 80 + + motivos = [] + condicoes = [] + saude_individual = [] + conectado = self.get_base_ip() is not None and self.nic_name is not None and self._mqtt_conectado + + (rtt, loss) = self._ping_once() + #print(f"rtt: {rtt}, loss: {loss}") + self.timeouts.append(loss == 100) + if rtt is not None: + self.rtts.append(rtt) + if loss is not None: + self.loses.append(loss) + + self._update_nic_errors() + self._update_nic_bandwidth() + + (ls, tmout, los, js, ns, hbs, avg_rtt, tmout_pct, loss_pct, jitter, err_rate, atraso) = self._metric_scores() + + health = ( + 0.25 * ls + + 0.10 * tmout + + 0.30 * los + + 0.15 * js + + 0.10 * ns + + 0.10 * hbs + ) + saude = round(health, 1) + + # -------- SAÚDE INDIVIDUAL / CONDIÇÕES / MOTIVOS -------- + + # 1) Latência + cond_lat = [] + if ls < SAUDE_MIN_ALERTA: + severidade = int(max(0, 100 - ls)) + c = { + "label": "Latência ICMP", + "valor": avg_rtt if avg_rtt is not None else -1, + "severidade": severidade, + "descricao": "Latência alta na comunicação com a base", + "acoes": [ + "Reduzir taxa de envio de telemetria", + "Verificar alinhamento das antenas 900 MHz", + "Verificar interferência ou obstáculos entre base e rover" + ] + } + condicoes.append(c) + cond_lat.append(c) + motivos.append( + f"Latência ICMP elevada " + f"({avg_rtt:.0f} ms, score {ls:.0f})." + if avg_rtt is not None else + f"Latência ICMP comprometida (score {ls:.0f})." + ) + + saude_individual.append({ + "id": "ip_latency", + "label": "Latência ICMP", + "status": self._status_por_score(ls).value, + "saude": round(ls, 1), + "motivos": [m for m in motivos if "Latência" in m], + "condicoes_operacionais": cond_lat, + "em_uso": True, + }) + + # 2) Timeouts + cond_timeout = [] + if tmout < SAUDE_MIN_ALERTA: + severidade = int(max(0, 100 - tmout)) + c = { + "label": "Timeout", + "valor": round(tmout_pct, 2), + "severidade": severidade, + "descricao": "Timeout na comunicação ICMP", + "acoes": [ + "Checar conectores e cabo da ponte 900 MHz", + "Verificar nível de ruído / interferência no enlace", + "Reduzir banda utilizada (streaming de vídeo, logs, etc.)" + ] + } + condicoes.append(c) + cond_timeout.append(c) + motivos.append( + f"Timeout elevado ({tmout_pct:.1f}%, score {tmout:.0f})." + ) + + saude_individual.append({ + "id": "ip_timeout", + "label": "Timeout", + "status": self._status_por_score(tmout).value, + "saude": round(tmout, 1), + "motivos": [m for m in motivos if "Timeout" in m], + "condicoes_operacionais": cond_timeout, + "em_uso": True, + }) + + # 2) Perda de pacotes + cond_loss = [] + if los < SAUDE_MIN_ALERTA: + severidade = int(max(0, 100 - los)) + c = { + "label": "Perda de Pacotes", + "valor": round(loss_pct, 2), + "severidade": severidade, + "descricao": "Perda de pacotes na comunicação ICMP", + "acoes": [ + "Checar conectores e cabo da ponte 900 MHz", + "Verificar nível de ruído / interferência no enlace", + "Reduzir banda utilizada (streaming de vídeo, logs, etc.)" + ] + } + condicoes.append(c) + cond_loss.append(c) + motivos.append( + f"Perda de pacotes elevada ({loss_pct:.1f}%, score {los:.0f})." + ) + + saude_individual.append({ + "id": "ip_loss", + "label": "Perda de pacotes", + "status": self._status_por_score(los).value, + "saude": round(los, 1), + "motivos": [m for m in motivos if "Perda de pacotes" in m], + "condicoes_operacionais": cond_loss, + "em_uso": True, + }) + + # 3) Jitter + cond_jit = [] + if js < SAUDE_MIN_ALERTA: + severidade = int(max(0, 100 - js)) + c = { + "label": "Jitter (variação de latência)", + "valor": round(jitter, 2), + "severidade": severidade, + "descricao": "Variação de latência acima do ideal", + "acoes": [ + "Evitar tráfego pesado na mesma rede da ponte", + "Reduzir taxa de envio de mensagens de controle", + "Verificar qualidade do enlace 900 MHz" + ] + } + condicoes.append(c) + cond_jit.append(c) + motivos.append( + f"Jitter elevado ({jitter:.1f} ms, score {js:.0f})." + ) + + saude_individual.append({ + "id": "ip_jitter", + "label": "Jitter da conexão", + "status": self._status_por_score(js).value, + "saude": round(js, 1), + "motivos": [m for m in motivos if "Jitter" in m], + "condicoes_operacionais": cond_jit, + "em_uso": True, + }) + + # 4) Erros na NIC + cond_nic = [] + if ns < SAUDE_MIN_ALERTA: + severidade = int(max(0, 100 - ns)) + c = { + "label": "Erros na interface de rede", + "valor": round(err_rate, 3), + "severidade": severidade, + "descricao": "Erros na placa de rede ligada à ponte 900 MHz", + "acoes": [ + "Verificar cabo de rede e conectores", + "Checar se há colisões ou problemas físicos no link", + "Substituir cabo ou porta de switch se necessário" + ] + } + condicoes.append(c) + cond_nic.append(c) + motivos.append( + f"Erros na NIC acima do normal ({err_rate:.3f} erros/mil pacotes, score {ns:.0f})." + ) + + saude_individual.append({ + "id": "ip_nic_errors", + "label": "Erros da interface de rede", + "status": self._status_por_score(ns).value, + "saude": round(ns, 1), + "motivos": [m for m in motivos if "NIC" in m or "interface de rede" in m], + "condicoes_operacionais": cond_nic, + "em_uso": True, + }) + + # 5) Heartbeat + cond_hb = [] + if hbs < SAUDE_MIN_ALERTA: + severidade = int(max(0, 100 - hbs)) + c = { + "label": "Atraso de heartbeat", + "valor": round(atraso, 2), + "severidade": severidade, + "descricao": "Atraso na recepção de heartbeat do rover", + "acoes": [ + "Verificar estado do serviço de telemetria no rover", + "Checar fila de mensagens MQTT/Redis", + "Garantir prioridade para mensagens de controle" + ] + } + condicoes.append(c) + cond_hb.append(c) + motivos.append( + f"Heartbeat atrasado ({atraso:.1f} s sem atualização, score {hbs:.0f})." + ) + + saude_individual.append({ + "id": "ip_heartbeat", + "label": "Heartbeat do rover", + "status": self._status_por_score(hbs).value, + "saude": round(hbs, 1), + "motivos": [m for m in motivos if "Heartbeat" in m], + "condicoes_operacionais": cond_hb, + "em_uso": True, + }) + + rx_mbps = statistics.mean(self.bw_rx_mbps) if self.bw_rx_mbps else 0.0 + tx_mbps = statistics.mean(self.bw_tx_mbps) if self.bw_tx_mbps else 0.0 + bw_total = rx_mbps + tx_mbps + + # -------- STATUS GERAL DO MÓDULO -------- + status = StatusModulo.OPERANTE + if not conectado: + status = StatusModulo.DESCONECTADO + elif saude <= 0: + status = StatusModulo.FALHA + elif saude < SAUDE_MIN_ALERTA: + status = StatusModulo.ALERTA + + payload = { + "conectado": conectado, + "status": status.value, + "saude": saude, + "motivos": motivos, + "saude_individual": saude_individual, + "condicoes_operacionais": condicoes, + "detalhes": { + "ls": ls, + "tmout": tmout, + "los": los, + "js": js, + "ns": ns, + "hbs": hbs, + "avg_rtt": avg_rtt, + "tmout_pct": tmout_pct, + "loss_pct": loss_pct, + "jitter": jitter, + "err_rate": err_rate, + "atraso_hb": atraso, + "rx_mbps": rx_mbps, + "tx_mbps": tx_mbps, + "bw_total_mbps": bw_total + } + } + + ContextoGlobalRedis.atualizar_ctx_dict( + ContextoGlobalRedis.ModKey(self.t_code), + saude=payload + ) + except Exception as e: + print(f"Erro ao atualizar saude do modulo {self.t_code.name}: {e}")