From 5fb8caa88cfd1ef859095058df480f439000804a Mon Sep 17 00:00:00 2001 From: Diego Freitas Date: Sat, 19 Sep 2026 00:41:03 -0300 Subject: [PATCH] centralizacao no corredor e resync de mapa --- .gitignore | 1 + .../Operadores/HealthWorkerService.cs | 125 +++-- .../Scripts/workers/manager_worker/main.py | 249 +++++++--- .../workers/manager_worker/modulos/mpc.py | 443 ++++++++++++++++-- .../workers/shared/contexto_global_redis.py | 425 ++++++++++++----- 5 files changed, 974 insertions(+), 269 deletions(-) diff --git a/.gitignore b/.gitignore index 8d72c4f87..7ccb4536d 100644 --- a/.gitignore +++ b/.gitignore @@ -98,3 +98,4 @@ Python/OAK/datasets/oak-fcc-3/audit/ /Python/OAK/datasets/oak-d/benchmark_corridor /AgroBase/livox_visual_debugger/x64/Debug /AgroBase/AgroBase/.vs/ +/.vs diff --git a/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs b/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs index 73047fdd4..58e289225 100644 --- a/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs +++ b/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs @@ -63,7 +63,40 @@ namespace AgroBase.Services.Operadores case HealthWorkerCommandType.ScriptCarregado: DadosLeitura.Pronto = true; DadosLeitura.ProntoEm = DateTime.Now; - MostrarLog($"Comando recebido: cmd = {cmd.ToString()}"); + + _ = Task.Run(async () => + { + try + { + var op = Variaveis.OperacaoEmAndamento; + + if ( + op != null && + (op.Parametros?.ControleAutomatico ?? false) && + (op.Trajetoria?._TrajetoriaFixa?.Any() ?? false) + ) + { + Variaveis.MostrarLog( + "[MAPSYNC/RECONNECT] Python reconectado; " + + "ressincronizando mapa/MPC atual." + ); + + bool ok = await SincronizarMapaCriticoAsync(op) + .ConfigureAwait(false); + + Variaveis.MostrarLog( + $"[MAPSYNC/RECONNECT] Resync concluído | sucesso={ok}" + ); + } + } + catch (Exception ex) + { + Variaveis.MostrarLog( + $"[MAPSYNC/RECONNECT] Erro: {ex.Message}" + ); + } + }); + break; default: @@ -285,7 +318,7 @@ namespace AgroBase.Services.Operadores } } - public static void AtualizarDadosOperacao(bool forcar = false, bool sincronizarMapa = true, bool configurado = true) + public static void AtualizarDadosOperacao(bool forcar = false, bool sincronizarMapa = false, bool configurado = true) { var op = Variaveis.OperacaoEmAndamento; @@ -391,6 +424,9 @@ namespace AgroBase.Services.Operadores horizonte = pControle.MpcHorizonte, passos_atraso = 1, ku_direcional = VariaveisEquipamento.KuDirecional, + tracking_row_hybrid_diagonal_enabled = true, + entry_capture_enabled = false, + tracking_axis_debug = false, } } ), @@ -649,12 +685,8 @@ namespace AgroBase.Services.Operadores ("mapa_quantidade_pontos_esperada", quantidadeEsperada), ("mapa_primeiro_idx_esperado", primeiroIdxEsperado), ("mapa_ultimo_idx_esperado", ultimoIdxEsperado), - ("mapa_revisao_processada", 0L), - ("mapa_quantidade_pontos_processada", 0), - ("mapa_revisao_mpc_aplicada", 0L), - ("mapa_quantidade_pontos_mpc_aplicada", 0), - ("mapa_command_id_mpc_aplicado", ""), - ("mapa_status_mpc", "pendente"), + // ACKs do Health e do Manager vivem em chaves próprias. + // DadosOperacao guarda apenas a intenção/snapshot desejado. ("mapa_confirmado", false), ("mapa_status", "pendente") ); @@ -744,7 +776,7 @@ namespace AgroBase.Services.Operadores ); if (tentativa < 3) - await Task.Delay(250).ConfigureAwait(false); + await Task.Delay(500).ConfigureAwait(false); } RedisService.AtualizarCampos( @@ -812,7 +844,7 @@ namespace AgroBase.Services.Operadores bool sucesso = await AguardarMapaAplicadoAsync( commandId, mapRevision, - TimeSpan.FromSeconds(4) + TimeSpan.FromSeconds(6) ).ConfigureAwait(false); if (!sucesso) @@ -840,53 +872,64 @@ namespace AgroBase.Services.Operadores return false; } - // Segunda confirmacao, diretamente no estado transacional usado - // pelo Manager. O ACK so vale se a MESMA revisao e a MESMA - // quantidade de pontos realmente chegaram ao MPC. - long revisaoProcessada = RedisService.GetField( - CtxKey.DadosOperacao, - "mapa_revisao_processada", + // Segunda confirmação por ownership: + // - HealthWorker confirma conversão/transporte em DadosHealthWorker; + // - Manager confirma o MPC em DadosManagerWorker; + // - DadosOperacao NÃO recebe ACK de outros processos. + string healthCmd = RedisService.GetField( + CtxKey.DadosHealthWorker, + "ultimo_comando_mapa", + "" + ); + long healthRev = RedisService.GetField( + CtxKey.DadosHealthWorker, + "mapa_revisao_aplicada", 0L ); - - int qtdProcessada = RedisService.GetField( - CtxKey.DadosOperacao, - "mapa_quantidade_pontos_processada", - 0 + string healthStatus = RedisService.GetField( + CtxKey.DadosHealthWorker, + "mapa_status", + "" ); - long revisaoMpc = RedisService.GetField( - CtxKey.DadosOperacao, - "mapa_revisao_mpc_aplicada", + string managerCmd = RedisService.GetField( + CtxKey.DadosManagerWorker, + "mapa_command_id_aplicado", + "" + ); + long managerRev = RedisService.GetField( + CtxKey.DadosManagerWorker, + "mapa_revisao_aplicada", 0L ); - - int qtdMpc = RedisService.GetField( - CtxKey.DadosOperacao, - "mapa_quantidade_pontos_mpc_aplicada", + int managerQtd = RedisService.GetField( + CtxKey.DadosManagerWorker, + "mapa_quantidade_pontos_aplicada", 0 ); - - string statusMpc = RedisService.GetField( - CtxKey.DadosOperacao, - "mapa_status_mpc", + string managerStatus = RedisService.GetField( + CtxKey.DadosManagerWorker, + "mapa_status", "" ); bool confirmacaoExata = - revisaoProcessada == mapRevision && - qtdProcessada == quantidadeEsperada && - revisaoMpc == mapRevision && - qtdMpc == quantidadeEsperada && - statusMpc == "aplicado"; + healthCmd == commandId && + healthRev == mapRevision && + healthStatus == "aplicado" && + managerCmd == commandId && + managerRev == mapRevision && + managerQtd == quantidadeEsperada && + managerStatus == "aplicado"; if (!confirmacaoExata) { Variaveis.MostrarLog( - $"[MAPSYNC] ACK rejeitado por divergencia de snapshot | " + - $"revEsperada={mapRevision} | revProcessada={revisaoProcessada} | revMpc={revisaoMpc} | " + - $"qtdEsperada={quantidadeEsperada} | qtdProcessada={qtdProcessada} | qtdMpc={qtdMpc} | " + - $"statusMpc={statusMpc}" + $"[MAPSYNC] ACK rejeitado por divergencia de ownership | " + + $"cmdEsperado={commandId} | healthCmd={healthCmd} | managerCmd={managerCmd} | " + + $"revEsperada={mapRevision} | healthRev={healthRev} | managerRev={managerRev} | " + + $"qtdEsperada={quantidadeEsperada} | managerQtd={managerQtd} | " + + $"statusHealth={healthStatus} | statusManager={managerStatus}" ); return false; } diff --git a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/main.py b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/main.py index c02a6c35b..75605b9a4 100644 --- a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/main.py +++ b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/main.py @@ -1,6 +1,7 @@ import sys import os import time +import threading def main(): @@ -25,6 +26,15 @@ def main(): loop_fixo = True frequencia_loop = 15 # hz + # Serializa somente as aplicações MAPSYNC. Redis Pub/Sub pode entregar + # replays do mesmo command_id enquanto a primeira aplicação ainda termina. + _mapsync_apply_lock = threading.Lock() + + # Enquanto uma aplicação de mapa está em andamento, suspendemos apenas o + # heartbeat do DadosManagerWorker. Assim o ACK transacional não compete + # com escritas de telemetria de 15 Hz no MESMO JSON Redis. + _mapsync_em_andamento = threading.Event() + def loop_ativo(): t0 = time.time() while True: @@ -43,12 +53,16 @@ def main(): mostrar_log(f"❌ Erro no loop do Manager: {e}") latencia = time.time() - t0 - ContextoGlobalRedis.atualizar_ctx_dict( - CtxKey.DadosManagerWorker, - momento=agora, - frequencia=frequencia, - latencia=latencia - ) + + # MAPSYNC tem prioridade sobre heartbeat/telemetria nesta chave. + # O heartbeat volta automaticamente assim que o ACK for gravado. + if not _mapsync_em_andamento.is_set(): + ContextoGlobalRedis.atualizar_ctx_dict( + CtxKey.DadosManagerWorker, + momento=agora, + frequencia=frequencia, + latencia=latencia + ) delay_corrigido = max(0, (1.0 / frequencia_loop) - latencia) if loop_fixo else 0.1 time.sleep(delay_corrigido) @@ -73,10 +87,15 @@ def main(): _mpc_iniciado = get_iniciado() if _mpc_iniciado and _p is None: - ContextoGlobalRedis()._atualizar_mpc( - do_manager=True, - iniciar_operacao=True + # O mapa já foi carregado e confirmado pelo MAPSYNC. O + # antigo fluxo chamava _atualizar_mpc(..., True), mas esse + # quarto argumento chega em mpc.inicializar() como + # forcar=True e recriava a instância sem necessidade. + mostrar_log( + "[MAPSYNC/MANAGER] início confirmado; " + "instância MPC preservada" ) + return elif _p is not None: parametros = _p.get("parametros") @@ -87,17 +106,9 @@ def main(): ) p_ref = _p.get("p_ref", []) - command_id = _p.get( - "command_id" - ) - - map_revision = _p.get( - "map_revision" - ) - - expected_point_count = _p.get( - "expected_point_count" - ) + command_id = _p.get("command_id") + map_revision = _p.get("map_revision") + expected_point_count = _p.get("expected_point_count") if expected_point_count is not None: expected_point_count = int(expected_point_count) @@ -108,45 +119,144 @@ def main(): f"esperado={expected_point_count}, recebido={quantidade_recebida}" ) - mostrar_log( - f"[MAPSYNC/MANAGER] aplicando mapa | " - f"cmd={command_id} | " - f"rev={map_revision} | " - f"pontos={len(mapa or [])}" - ) - - iniciar_mpc( - parametros, - mapa, - p_ref, - iniciar_operacao - ) - - # --------------------------------------------- - # ACK somente depois de inicializar() - # retornar sem exception. - # --------------------------------------------- - - if ( + # MAPSYNC transacional deve sempre trazer identidade. + # Comandos legados sem identidade continuam possíveis, mas + # nunca contam como ACK da barreira crítica. + eh_mapsync = ( command_id is not None and map_revision is not None - ): - ContextoGlobalRedis.atualizar_ctx_dict( - CtxKey.DadosOperacao, - mapa_command_id_mpc_aplicado=str(command_id), - mapa_revisao_mpc_aplicada=map_revision, - mapa_quantidade_pontos_mpc_aplicada=len(mapa or []), - mapa_status_mpc="aplicado", - mapa_erro_mpc="", - mapa_mpc_aplicado_em=time.time(), - ) + and expected_point_count is not None + ) - mostrar_log( - f"[MAPSYNC/MANAGER] mapa aplicado | " - f"cmd={command_id} | " - f"rev={map_revision} | " - f"pontos={len(mapa or [])}" - ) + lock_ctx = _mapsync_apply_lock if eh_mapsync else threading.Lock() + + # Pare imediatamente o heartbeat concorrente antes de + # tocar no estado transacional do Manager. + if eh_mapsync: + _mapsync_em_andamento.set() + + try: + with lock_ctx: + if eh_mapsync: + estado_mgr = ContextoGlobalRedis.get( + CtxKey.DadosManagerWorker, {} + ) + if not isinstance(estado_mgr, dict): + estado_mgr = {} + + ack_cmd = str( + estado_mgr.get("mapa_command_id_aplicado", "") or "" + ) + ack_rev = int( + estado_mgr.get("mapa_revisao_aplicada", 0) or 0 + ) + ack_qtd = int( + estado_mgr.get("mapa_quantidade_pontos_aplicada", 0) or 0 + ) + ack_status = str( + estado_mgr.get("mapa_status", "") or "" + ) + + incoming_rev = int(map_revision) + incoming_qtd = int(expected_point_count) + + # Uma revisão mais nova já aplicada nunca pode ser + # rebaixada por uma mensagem atrasada do Pub/Sub. + if ack_status == "aplicado" and ack_rev > incoming_rev: + mostrar_log( + f"[MAPSYNC/MANAGER] replay obsoleto ignorado | " + f"cmd={command_id} | rev={incoming_rev} | " + f"rev_atual={ack_rev}" + ) + return + + # Idempotência por revisão: command_id é identidade + # de transporte; map_revision identifica o snapshot. + # Se esta MESMA revisão/quantidade já está aplicada, + # apenas reafirmamos o ACK para o command_id recebido. + if ( + ack_status == "aplicado" + and ack_rev == incoming_rev + and ack_qtd == incoming_qtd + ): + ContextoGlobalRedis.atualizar_ctx_dict( + CtxKey.DadosManagerWorker, + mapa_command_id_aplicado=str(command_id), + mapa_revisao_aplicada=incoming_rev, + mapa_quantidade_pontos_aplicada=incoming_qtd, + mapa_status="aplicado", + mapa_erro="", + mapa_aplicado_em=time.time(), + ) + + tipo_replay = ( + "duplicado" if ack_cmd == str(command_id) + else "rebind" + ) + mostrar_log( + f"[MAPSYNC/MANAGER] ACK reafirmado ({tipo_replay}) | " + f"cmd={command_id} | rev={incoming_rev} | " + f"pontos={incoming_qtd}" + ) + return + + # Sinaliza recepção/posse da transação ANTES de + # construir o MPC. O Health vê "aplicando", para de + # retransmitir e concede uma janela de conclusão. + if eh_mapsync: + ContextoGlobalRedis.atualizar_ctx_dict( + CtxKey.DadosManagerWorker, + mapa_command_id_aplicado=str(command_id), + mapa_revisao_aplicada=int(map_revision), + mapa_quantidade_pontos_aplicada=int(expected_point_count), + mapa_status="aplicando", + mapa_erro="", + mapa_aplicado_em=0.0, + ) + + mostrar_log( + f"[MAPSYNC/MANAGER] aplicando mapa | " + f"cmd={command_id} | " + f"rev={map_revision} | " + f"pontos={len(mapa or [])}" + ) + + resultado_mpc = iniciar_mpc( + parametros, + mapa, + p_ref, + iniciar_operacao + ) + + if resultado_mpc is False: + raise RuntimeError( + "MPC recusou a inicialização do snapshot recebido" + ) + + if ( + command_id is not None + and map_revision is not None + ): + ContextoGlobalRedis.atualizar_ctx_dict( + CtxKey.DadosManagerWorker, + mapa_command_id_aplicado=str(command_id), + mapa_revisao_aplicada=map_revision, + mapa_quantidade_pontos_aplicada=len(mapa or []), + mapa_status="aplicado", + mapa_erro="", + mapa_aplicado_em=time.time(), + ) + + mostrar_log( + f"[MAPSYNC/MANAGER] mapa aplicado | " + f"cmd={command_id} | " + f"rev={map_revision} | " + f"pontos={len(mapa or [])}" + ) + + finally: + if eh_mapsync: + _mapsync_em_andamento.clear() else: ContextoGlobalRedis()._atualizar_mpc( @@ -161,12 +271,13 @@ def main(): ): try: ContextoGlobalRedis.atualizar_ctx_dict( - CtxKey.DadosOperacao, - mapa_command_id_mpc_aplicado=str(command_id), - mapa_revisao_mpc_aplicada=map_revision, - mapa_quantidade_pontos_mpc_aplicada=0, - mapa_status_mpc="erro", - mapa_erro_mpc=str(e), + CtxKey.DadosManagerWorker, + mapa_command_id_aplicado=str(command_id), + mapa_revisao_aplicada=map_revision, + mapa_quantidade_pontos_aplicada=0, + mapa_status="erro", + mapa_erro=str(e), + mapa_aplicado_em=time.time(), ) except Exception: pass @@ -194,6 +305,20 @@ def main(): def inicializar(): mostrar_log("🚀 Iniciando Manager Worker...") + + # ACK pertence à memória viva deste processo. Se o Manager reiniciou, + # uma confirmação antiga no Redis não pode liberar a operação como se + # o MPC ainda estivesse carregado nesta nova instância. + ContextoGlobalRedis.atualizar_ctx_dict( + CtxKey.DadosManagerWorker, + mapa_command_id_aplicado="", + mapa_revisao_aplicada=0, + mapa_quantidade_pontos_aplicada=0, + mapa_status="", + mapa_erro="worker reiniciado; aguardando MAPSYNC", + mapa_aplicado_em=0.0, + ) + ContextoGlobalRedis.ouvir_comandos(CmdKey.ManagerWorkerRx, redis_callback) ContextoGlobalRedis.publicar_comando(CmdKey.ManagerWorkerTx, { "cmd": ManagerWorkerCommandType.ScriptCarregado.value }) loop_ativo() diff --git a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/modulos/mpc.py b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/modulos/mpc.py index 699d5ba54..2234ffd3b 100644 --- a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/modulos/mpc.py +++ b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/manager_worker/modulos/mpc.py @@ -547,17 +547,18 @@ class ControladorMPC: # ------------------------------------------------------------------ # O campo mostrou que, ao terminar a ferradura, o rover pode entrar na # nova passada com 20-70 cm de cross-track e ainda 5-12 graus de yaw. - # A política antiga priorizava yaw e só depois liberava o crab, fazendo - # o rover consumir alguns metros dentro da rua antes de centralizar. + # V21.1 - ENTRY CAPTURE legado fica DESLIGADO por padrão. # - # Agora o fim do Pure Pursuit arma uma janela curta "centro primeiro": - # - se o erro lateral ainda é relevante e |heading| <= 15 graus, - # MovimentoDiagonal recebe prioridade mesmo sem heading perfeito; - # - heading muito grande continua com Arco/Dianteira normalmente; - # - assim que |e_lat| <= 10 cm, a política normal volta a mandar; - # - existe limite de distância para o latch nunca permanecer indefinido. + # Nos testes de entrada de corredor, o crab pós-Dubins podia assumir o + # volante ainda com heading moderado (até 15 graus) e preservar uma pose + # ruim logo após a ferradura. A V20/V21 já consegue capturar a centerline + # com Dianteira/Traseira de forma contínua, portanto a transição normal + # passa a ser a política preferida. + # + # O código permanece disponível apenas para A/B/rollback explícito: + # entry_capture_enabled = True self._entry_capture_enabled = bool( - parametros_mpc.get("entry_capture_enabled", True) + parametros_mpc.get("entry_capture_enabled", False) ) self._entry_capture_heading_max_graus = _safe_float( parametros_mpc.get("entry_capture_heading_max_graus", 15.0), @@ -731,6 +732,18 @@ class ControladorMPC: ) self._tracking_eixo_yaw_ultimo_log = None + # V21.1 - diagnóstico opcional do seletor Dianteira x Traseira. + # Não altera a escolha; apenas mostra periodicamente os dois custos para + # comprovar que RodasTraseiras continua participando da competição. + self._tracking_axis_debug = bool( + parametros_mpc.get("tracking_axis_debug", False) + ) + self._tracking_axis_debug_period_s = _safe_float( + parametros_mpc.get("tracking_axis_debug_period_s", 2.0), + 2.0, min_value=0.25, max_value=30.0, + ) + self._tracking_axis_debug_last_t = 0.0 + # ------------------------------------------------------------------ # V16 - CINEMATICA 4WS UNIFICADA # ------------------------------------------------------------------ @@ -791,15 +804,106 @@ class ControladorMPC: 9.0, min_value=0.0, max_value=self.angulo_max_graus, ) self._tracking_curve_diag_emergency_m = _safe_float( - parametros_mpc.get("tracking_curve_diag_emergency_m", 0.24), - 0.24, min_value=self._tracking_center_recenter_m, max_value=0.80, + parametros_mpc.get("tracking_curve_diag_emergency_m", 0.12), + 0.12, min_value=self._tracking_center_recenter_m, max_value=0.80, ) self._tracking_curve_diag_release_m = _safe_float( - parametros_mpc.get("tracking_curve_diag_release_m", 0.14), - 0.14, min_value=self._tracking_center_hold_m, + parametros_mpc.get("tracking_curve_diag_release_m", 0.06), + 0.06, min_value=self._tracking_center_hold_m, max_value=self._tracking_curve_diag_emergency_m, ) + # ------------------------------------------------------------------ + # V20 - TRACKING CONTINUO DE CORREDOR + # + # Em CaminhandoRua não alternamos mais entre "corrigir lateral com + # Diagonal" e "corrigir yaw com Dianteira/Traseira". Esse chaveamento + # desacopla duas grandezas que, geometricamente, precisam ser resolvidas + # juntas. O controlador nominal da passada passa a ser contínuo: + # + # delta = feed-forward(curvatura) + # + feedback(heading) + # + feedback(cross-track) + # + # Dianteira/Traseira continuam competindo no seletor axle-aware. + # MovimentoDiagonal fica preservado para Entry Capture e fluxos legados, + # mas não é o HOLD nominal de CaminhandoRua quando esta opção está ativa. + self._tracking_row_continuous_enabled = bool( + parametros_mpc.get("tracking_row_continuous_enabled", True) + ) + self._tracking_row_lateral_gain = _safe_float( + parametros_mpc.get("tracking_row_lateral_gain", 0.85), + 0.85, min_value=0.0, max_value=3.0, + ) + self._tracking_row_lateral_max_graus = _safe_float( + parametros_mpc.get("tracking_row_lateral_max_graus", 12.0), + 12.0, min_value=1.0, max_value=self.angulo_max_graus, + ) + + # O PP da cabeceira responde imediatamente porque MovimentoArco não + # passa pelo filtro de tracking. Em passada mantemos amortecimento, mas + # com autoridade suficiente para não gastar metros "acordando" a roda. + self._tracking_row_filter_tau_s = _safe_float( + parametros_mpc.get("tracking_row_filter_tau_s", 0.12), + 0.12, min_value=0.0, max_value=1.0, + ) + self._tracking_row_rate_graus_s = _safe_float( + parametros_mpc.get("tracking_row_rate_graus_s", 35.0), + 35.0, min_value=5.0, max_value=120.0, + ) + + # ------------------------------------------------------------------ + # V21 - RECOVERY HIBRIDO POR MOVIMENTO DIAGONAL + # + # A V20 continua sendo a lei NOMINAL da passada. O crab/Diagonal só + # recebe o volante quando existe assinatura de derrapagem lateral: + # cross-track grande, corpo ainda alinhado e centerline reta ou apenas + # suavemente curva. Assim preservamos o tracking continuo da V20 e + # recuperamos a principal vantagem pratica do 4WS em terreno inclinado. + # + # Flag desligada => comportamento V20 puro, bit-a-bit na politica. + self._tracking_row_hybrid_diagonal_enabled = bool( + parametros_mpc.get("tracking_row_hybrid_diagonal_enabled", False) + ) + self._tracking_row_hybrid_enter_m = _safe_float( + parametros_mpc.get("tracking_row_hybrid_enter_m", 0.15), + 0.15, min_value=0.08, max_value=0.80, + ) + self._tracking_row_hybrid_exit_m = _safe_float( + parametros_mpc.get("tracking_row_hybrid_exit_m", 0.07), + 0.07, min_value=self._tracking_center_hold_m, + max_value=self._tracking_row_hybrid_enter_m, + ) + self._tracking_row_hybrid_heading_enter_graus = _safe_float( + parametros_mpc.get("tracking_row_hybrid_heading_enter_graus", 3.0), + 3.0, min_value=0.5, max_value=10.0, + ) + self._tracking_row_hybrid_heading_abort_graus = _safe_float( + parametros_mpc.get("tracking_row_hybrid_heading_abort_graus", 6.0), + 6.0, min_value=self._tracking_row_hybrid_heading_enter_graus, + max_value=15.0, + ) + self._tracking_row_hybrid_curve_enter_max_graus = _safe_float( + parametros_mpc.get("tracking_row_hybrid_curve_enter_max_graus", 1.5), + 1.5, min_value=0.0, max_value=8.0, + ) + self._tracking_row_hybrid_curve_abort_graus = _safe_float( + parametros_mpc.get("tracking_row_hybrid_curve_abort_graus", 2.5), + 2.5, min_value=self._tracking_row_hybrid_curve_enter_max_graus, + max_value=12.0, + ) + self._tracking_row_hybrid_max_dist_m = _safe_float( + parametros_mpc.get("tracking_row_hybrid_max_dist_m", 2.50), + 2.50, min_value=0.50, max_value=10.0, + ) + self._tracking_row_hybrid_cooldown_m = _safe_float( + parametros_mpc.get("tracking_row_hybrid_cooldown_m", 0.75), + 0.75, min_value=0.0, max_value=5.0, + ) + self._tracking_row_hybrid_active = False + self._tracking_row_hybrid_dist_m = 0.0 + self._tracking_row_hybrid_cooldown_restante_m = 0.0 + # Estado interno apenas para histerese do recovery lateral em curva. # Não sai do MPC e não altera o contrato C#/Redis. self._tracking_curve_diag_recovery_active = False @@ -1892,6 +1996,33 @@ class ControladorMPC: ) return float(np.clip(delta, -lim, lim)) + def _angulo_lateral_row_tracking(self, e_lat, velocidade): + """V20: parcela lateral contínua da passada para Dianteira/Traseira. + + É uma forma Stanley-like: quanto maior o cross-track, maior o heading + corretivo pedido; a velocidade reduz a agressividade naturalmente. + Diferente do crab, esta parcela NÃO desliga a autoridade de yaw. + """ + e_lat = float(e_lat) + v = max(0.0, float(velocidade)) + + e_eff_abs = max(0.0, abs(e_lat) - self._tracking_center_hold_m) + if e_eff_abs <= 0.0 or self._tracking_row_lateral_gain <= 0.0: + return 0.0 + + e_eff = math.copysign(e_eff_abs, e_lat) + delta = math.atan2( + self._tracking_row_lateral_gain * e_eff, + v + self._velocidade_offset_lateral, + ) + lim = math.radians( + min( + self._tracking_row_lateral_max_graus, + self.angulo_max_graus, + ) + ) + return float(np.clip(delta, -lim, lim)) + def _angulo_diagonal_tracking(self, e_lat, velocidade): """Comando crab para centralizar sem alterar o heading do corpo.""" e_lat = float(e_lat) @@ -1912,6 +2043,114 @@ class ControladorMPC: ) return float(np.clip(delta, -lim, lim)) + def _atualizar_tracking_row_hybrid_real( + self, x, y, theta, velocidade, idx_base, status_carro, contexto + ): + """V21: atualiza o latch do recovery Diagonal usando apenas a pose REAL. + + A decisao nao e feita dentro das simulacoes do beam search. Isso evita + que candidatos hipoteticos armem/desarmem o estado do controlador. + O Diagonal entra somente com assinatura de deriva lateral: erro lateral + grande, heading pequeno e baixa curvatura da centerline. + + Histerese: + enter_m > exit_m; + heading_enter < heading_abort; + curve_enter < curve_abort. + + Existe ainda um limite de distancia e cooldown para impedir que o crab + fique preso indefinidamente caso a derrapagem supere sua autoridade. + """ + enabled = bool(self._tracking_row_hybrid_diagonal_enabled) + em_passada = _status_in(status_carro, [StatusCarroMapa.CaminhandoRua]) + entry_ativo = bool(getattr(self, "_entry_capture_active", False)) + + dt = _safe_float( + getattr(self, "tempo_execucao_local", 0.5), + 0.5, min_value=0.05, max_value=1.5, + ) + ds = abs(float(velocidade)) * dt + + # Cooldown e medido em metros, portanto permanece coerente em baixa e + # alta velocidade. + self._tracking_row_hybrid_cooldown_restante_m = max( + 0.0, + float(getattr(self, "_tracking_row_hybrid_cooldown_restante_m", 0.0)) - ds, + ) + + if not enabled or not em_passada or entry_ativo: + self._tracking_row_hybrid_active = False + self._tracking_row_hybrid_dist_m = 0.0 + return + + try: + ref = self._referencia_caminho_lookahead( + float(x), float(y), int(idx_base), float(velocidade), status_carro + ) + e_lat_abs = abs(float(ref["e_lat"])) + theta_ref = float(ref["theta_ref"]) + e_head_deg = abs(math.degrees( + self._erro_heading_caminho(float(theta), theta_ref) + )) + _, delta_ff = self._feedforward_curva_tracking( + ref.get("s_proj"), int(idx_base), contexto + ) + ff_deg_abs = abs(math.degrees(float(delta_ff))) + except Exception as e: + # Qualquer duvida geometrica devolve o volante a V20. + if bool(getattr(self, "_tracking_row_hybrid_active", False)): + _log(f"[MPC/V21 HYBRID] OFF por referencia invalida: {e}") + self._tracking_row_hybrid_active = False + self._tracking_row_hybrid_dist_m = 0.0 + return + + ativo = bool(getattr(self, "_tracking_row_hybrid_active", False)) + + if ativo: + self._tracking_row_hybrid_dist_m += ds + motivo_off = None + + if e_lat_abs <= float(self._tracking_row_hybrid_exit_m): + motivo_off = "centralizado" + elif e_head_deg > float(self._tracking_row_hybrid_heading_abort_graus): + motivo_off = "heading" + elif ff_deg_abs > float(self._tracking_row_hybrid_curve_abort_graus): + motivo_off = "curvatura" + elif self._tracking_row_hybrid_dist_m >= float(self._tracking_row_hybrid_max_dist_m): + motivo_off = "distancia_max" + self._tracking_row_hybrid_cooldown_restante_m = float( + self._tracking_row_hybrid_cooldown_m + ) + + if motivo_off is not None: + _log( + "[MPC/V21 HYBRID] Diagonal OFF | " + f"motivo={motivo_off} | e_lat={e_lat_abs:.3f}m | " + f"e_head={e_head_deg:.2f}deg | ff={ff_deg_abs:.2f}deg | " + f"dist={self._tracking_row_hybrid_dist_m:.2f}m" + ) + self._tracking_row_hybrid_active = False + self._tracking_row_hybrid_dist_m = 0.0 + return + + if self._tracking_row_hybrid_cooldown_restante_m > 1e-6: + return + + pode_entrar = bool( + e_lat_abs >= float(self._tracking_row_hybrid_enter_m) + and e_head_deg <= float(self._tracking_row_hybrid_heading_enter_graus) + and ff_deg_abs <= float(self._tracking_row_hybrid_curve_enter_max_graus) + ) + + if pode_entrar: + self._tracking_row_hybrid_active = True + self._tracking_row_hybrid_dist_m = 0.0 + _log( + "[MPC/V21 HYBRID] Diagonal ON | " + f"e_lat={e_lat_abs:.3f}m | e_head={e_head_deg:.2f}deg | " + f"ff={ff_deg_abs:.2f}deg" + ) + def _entre_eixos_tracking(self, contexto): equipamento = _as_dict(_as_dict(contexto or {}).get("Equipamento", {})) fallback = float(getattr(self, "_cinematica_entre_eixos_m", 0.94)) @@ -2233,28 +2472,38 @@ class ControladorMPC: } def _log_escolha_eixo_v15(self, escolhido, frente_dbg, traseira_dbg): - """Loga somente quando o eixo escolhido muda, evitando spam.""" + """V21.1: diagnóstico opcional do seletor de eixo, sem alterar controle.""" try: nome = _enum_name( escolhido, TipoMovimentoDirecional, default=str(escolhido) ) - if self._tracking_eixo_yaw_ultimo_log == nome: - return + mudou = self._tracking_eixo_yaw_ultimo_log != nome self._tracking_eixo_yaw_ultimo_log = nome - #_log( - # "[MPC/V15 AXIS] " - # f"eixo={nome} | " - # f"Jf={frente_dbg['custo']:.2f} " - # f"Jt={traseira_dbg['custo']:.2f} | " - # f"F0={frente_dbg['e_front_0']:+.3f} " - # f"T0={frente_dbg['e_rear_0']:+.3f} | " - # f"Ff={frente_dbg['e_front']:+.3f} " - # f"Tf={frente_dbg['e_rear']:+.3f} | " - # f"Ft={traseira_dbg['e_front']:+.3f} " - # f"Tt={traseira_dbg['e_rear']:+.3f} | " - # f"preserva={frente_dbg['melhor_eixo_0']}" - #) + if not bool(getattr(self, "_tracking_axis_debug", False)): + return + + agora = time.monotonic() + periodo = float(getattr(self, "_tracking_axis_debug_period_s", 2.0)) + ultimo = float(getattr(self, "_tracking_axis_debug_last_t", 0.0)) + if not mudou and (agora - ultimo) < periodo: + return + self._tracking_axis_debug_last_t = agora + + _log( + "[MPC/V21 AXIS] " + f"eixo={nome} | " + f"Jf={frente_dbg['custo']:.3f} " + f"Jt={traseira_dbg['custo']:.3f} | " + f"F0={frente_dbg['e_front_0']:+.3f} " + f"T0={frente_dbg['e_rear_0']:+.3f} | " + f"Ff={frente_dbg['e_front']:+.3f} " + f"Tf={frente_dbg['e_rear']:+.3f} | " + f"Ft={traseira_dbg['e_front']:+.3f} " + f"Tt={traseira_dbg['e_rear']:+.3f} | " + f"preservaF={frente_dbg['melhor_eixo_0']} " + f"preservaT={traseira_dbg['melhor_eixo_0']}" + ) except Exception: pass @@ -2339,21 +2588,16 @@ class ControladorMPC: s_proj=None, contexto=None, ): - """V15: supervisor de tracking da passada com CURVE_TRACK axle-aware. + """V20: supervisor contínuo de tracking da passada, axle-aware. - RETA: - - Diagonal é HOLD nominal e corrige cross-track sem criar yaw; - - Dianteira/Traseira entram quando heading realmente sai da banda; - - Arco continua sendo recuperação grosseira. + Com tracking_row_continuous_enabled=True: + - Arco fica reservado para recuperação grosseira de heading; + - Dianteira/Traseira controlam continuamente a centerline; + - comando = feed-forward(curvatura) + heading + cross-track; + - não existe chaveamento RETA/CURVA para decidir entre yaw e crab; + - Diagonal deixa de ser o HOLD nominal de CaminhandoRua. - CURVA: - - Dianteira/Traseira acompanham continuamente a curvatura da rua; - - comando = feed-forward(curvatura) + heading + lateral; - - não esperamos acumular heading/cross-track para começar a virar; - - Diagonal só interrompe CURVE_TRACK em erro lateral de emergência. - - Isso evita o deadlock da V13 em que a diagonal tentava centralizar - enquanto a própria centerline continuava girando. + O bloco V15 antigo permanece abaixo para rollback/A-B. """ e_lat = float(e_lat) e_lat_abs = abs(e_lat) @@ -2391,6 +2635,77 @@ class ControladorMPC: kappa, delta_ff = self._feedforward_curva_tracking( s_proj, idx_base, contexto ) + + # -------------------------------------------------------------- + # V21) Recovery hibrido: crab curto para derrapagem lateral pura. + # + # O latch e atualizado somente no ciclo REAL. Aqui apenas respeitamos + # sua decisao. Se por algum motivo a geometria ja mudou entre os dois + # pontos, as guardas abaixo devolvem imediatamente o volante a V20. + # -------------------------------------------------------------- + if ( + bool(self._tracking_row_hybrid_diagonal_enabled) + and bool(getattr(self, "_tracking_row_hybrid_active", False)) + ): + ff_deg_abs = abs(math.degrees(float(delta_ff))) + if ( + e_head_deg <= float(self._tracking_row_hybrid_heading_abort_graus) + and ff_deg_abs <= float(self._tracking_row_hybrid_curve_abort_graus) + ): + return ( + TipoMovimentoDirecional.MovimentoDiagonal, + self._angulo_diagonal_tracking(e_lat, v), + ) + self._tracking_row_hybrid_active = False + self._tracking_row_hybrid_dist_m = 0.0 + + # -------------------------------------------------------------- + # V20) CaminhandoRua: tracking contínuo e acoplado. + # + # Não existe mais a situação "heading pequeno => crab puro". Mesmo em + # reta, um cross-track gera uma curvatura de convergência. Conforme o + # rover cria heading para voltar ao centro, a parcela de heading se + # opõe naturalmente à lateral e amortece a chegada. + # + # Em curva, o feed-forward continua girando junto com a centerline + # enquanto a mesma lei corrige o deslocamento lateral. Assim não + # congelamos yaw justamente quando mais precisamos acompanhar a curva. + # -------------------------------------------------------------- + if bool(self._tracking_row_continuous_enabled): + self._tracking_curve_diag_recovery_active = False + + delta_heading = self._ganho_heading_reta * erro_heading_rad + delta_lateral = self._angulo_lateral_row_tracking(e_lat, v) + + delta_base = delta_ff + delta_heading + delta_lateral + lim = math.radians(self.angulo_max_graus) + delta_base = float(np.clip(delta_base, -lim, lim)) + + if ( + x is not None + and y is not None + and theta is not None + and idx_base is not None + ): + return self._escolher_eixo_yaw_tracking( + delta_base, + float(x), float(y), float(theta), + v, + int(idx_base), + status_carro, + tipo_prev, + contexto=contexto, + ) + + return ( + TipoMovimentoDirecional.RodasDianteiras, + delta_base, + ) + + # -------------------------------------------------------------- + # V15 legado abaixo: mantido para rollback/A-B com + # tracking_row_continuous_enabled=False. + # -------------------------------------------------------------- ff_deg_abs = abs(math.degrees(delta_ff)) curva_ativa = bool( @@ -2792,11 +3107,29 @@ class ControladorMPC: if inversao_significativa: filtrado = 0.0 else: - tau = float(self._filtro_angulo_reta_tau_s) + # V20: Dianteira/Traseira em CaminhandoRua precisam responder mais + # cedo que o filtro legado, mas continuam amortecidas. MovimentoArco + # permanece sem filtro, como antes. + usar_filtro_row = bool( + self._tracking_row_continuous_enabled + and _status_in(status_carro, [StatusCarroMapa.CaminhandoRua]) + and tipo in [ + TipoMovimentoDirecional.RodasDianteiras, + TipoMovimentoDirecional.RodasTraseiras, + ] + ) + + if usar_filtro_row: + tau = float(self._tracking_row_filter_tau_s) + taxa_graus_s = float(self._tracking_row_rate_graus_s) + else: + tau = float(self._filtro_angulo_reta_tau_s) + taxa_graus_s = float(self._taxa_angulo_reta_graus_s) + alpha = 1.0 if tau <= 1e-6 else dt / (tau + dt) filtrado = anterior + alpha * (desejado - anterior) - max_delta = np.radians(self._taxa_angulo_reta_graus_s) * dt + max_delta = np.radians(taxa_graus_s) * dt filtrado = anterior + float(np.clip( filtrado - anterior, -max_delta, @@ -3267,11 +3600,11 @@ class ControladorMPC: if idx_local > 0: pontos_visitados[:idx_local] = True - if havia_avanco_especulativo and not self._ultima_divergencia_visitados: - _log( - f"[MPC] Progresso especulativo descartado; " - f"C# autoritativo no idx_local={idx_local}." - ) + #if havia_avanco_especulativo and not self._ultima_divergencia_visitados: + # _log( + # f"[MPC] Progresso especulativo descartado; " + # f"C# autoritativo no idx_local={idx_local}." + # ) self._ultima_divergencia_visitados = havia_avanco_especulativo @@ -3906,6 +4239,8 @@ class ControladorMPC: self._dubins_pp_ativo_ciclo_anterior = True self._entry_capture_active = False self._entry_capture_dist_m = 0.0 + self._tracking_row_hybrid_active = False + self._tracking_row_hybrid_dist_m = 0.0 return self._comando_seguidor_dubins( contexto=contexto, @@ -3990,6 +4325,16 @@ class ControladorMPC: f"dist={self._entry_capture_dist_m:.2f}m" ) + # ------------------------------------------------------------ + # V21 - atualiza o latch hibrido uma unica vez com a pose REAL. + # A flag desligada preserva integralmente a V20 pura. + # ------------------------------------------------------------ + self._atualizar_tracking_row_hybrid_real( + x=x, y=y, theta=theta, velocidade=velocidade, + idx_base=idx_alvo_correcao, status_carro=status_carro, + contexto=contexto, + ) + # -------------------- Planejamento: horizonte ESPACIAL fixo -------------------- S_ALVO = float(self.horizonte) V_FLOOR = float(getattr(self, "_vmin_planejamento", 0.30)) diff --git a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/shared/contexto_global_redis.py b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/shared/contexto_global_redis.py index 4f1da5091..d0744e543 100644 --- a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/shared/contexto_global_redis.py +++ b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/shared/contexto_global_redis.py @@ -236,53 +236,70 @@ class ContextoGlobalRedis: @classmethod def atualizar_ctx_dict(cls, chave: CtxKey, **novos_campos): """ - Atualiza campos de um JSON salvo no Redis. + Atualiza campos de um JSON salvo no Redis sem perder writes concorrentes. - Suporta caminho aninhado: - atualizar_ctx_dict(CtxKey.DadosControle, "a__b"=1) - vira: - {"a": {"b": 1}} - - Também aceita caminho com ponto: - atualizar_ctx_dict(..., a__b=1) - ou internamente "a.b". + O lock Python é apenas local ao processo. Como ManagerWorker, + HealthWorker e os demais workers são processos separados, um simples + GET -> modifica -> SET pode apagar campos escritos por outro processo + (ou pelo C#). WATCH/MULTI transforma a atualização em CAS otimista. """ - from collections.abc import MutableMapping def _set_nested(d: MutableMapping, path: str, value): keys = str(path).split(".") atual = d - for k in keys[:-1]: if k not in atual or not isinstance(atual[k], dict): atual[k] = {} atual = atual[k] - atual[keys[-1]] = value - try: - chave_str = cls._key(chave) + chave_str = cls._key(chave) + novos_campos = { + str(k).replace("__", "."): v + for k, v in novos_campos.items() + } - novos_campos = { - str(k).replace("__", "."): v - for k, v in novos_campos.items() - } + ultimo_erro = None + for tentativa in range(8): + try: + # Serializa writers do mesmo processo; WATCH protege contra + # writers de outros processos e clientes externos (C#). + with cls._update_lock: + with cls._redis.pipeline() as pipe: + pipe.watch(chave_str) + bruto = pipe.get(chave_str) + dados = cls._json_loads(bruto, default={}) + if not isinstance(dados, dict): + dados = {} - with cls._update_lock: - dados = cls.get(chave_str, {}) + for campo, valor in novos_campos.items(): + _set_nested(dados, campo, valor) - if not isinstance(dados, dict): - dados = {} + payload = cls._json_dumps(dados) + pipe.multi() + pipe.set(chave_str, payload) + pipe.execute() + return True - for campo, valor in novos_campos.items(): - _set_nested(dados, campo, valor) + except redis.WatchError as e: + ultimo_erro = e + time.sleep(min(0.025, 0.002 * (tentativa + 1))) + continue + except (redis.ConnectionError, redis.TimeoutError) as e: + ultimo_erro = e + cls._reconectar() + time.sleep(0.02 * (tentativa + 1)) + continue + except Exception as e: + mostrar_log(f"Erro ao atualizar_ctx_dict: {e} - Chave: {chave}") + return False - return cls.set(chave_str, dados) - - except Exception as e: - mostrar_log(f"Erro ao atualizar_ctx_dict: {e} - Chave: {chave}") - return False + mostrar_log( + f"Erro ao atualizar_ctx_dict após conflitos concorrentes: " + f"{ultimo_erro} - Chave: {chave}" + ) + return False @classmethod def publicar_comando(cls, canal: CmdKey, dados): @@ -647,16 +664,27 @@ class ContextoGlobalRedis: ) revisao_solicitada = cls._int(operacao.get("mapa_revisao_solicitada", 0), 0) - revisao_processada = cls._int(operacao.get("mapa_revisao_processada", 0), 0) - revisao_mpc = cls._int(operacao.get("mapa_revisao_mpc_aplicada", 0), 0) qtd_esperada = cls._int(operacao.get("mapa_quantidade_pontos_esperada", 0), 0) - qtd_processada = cls._int(operacao.get("mapa_quantidade_pontos_processada", 0), 0) - qtd_mpc = cls._int(operacao.get("mapa_quantidade_pontos_mpc_aplicada", 0), 0) qtd_publicada = len(pontos_publicados) status_mapa = str(operacao.get("mapa_status", "") or "") - status_mpc = str(operacao.get("mapa_status_mpc", "") or "") mapa_confirmado_flag = cls._bool(operacao.get("mapa_confirmado", False)) mapa_barreira_critica = cls._bool(operacao.get("mapa_barreira_critica", False)) + command_solicitado = str(operacao.get("mapa_command_id", "") or "") + + # MAPSYNC V2: ACKs possuem dono único. + # HealthWorker nunca mais disputa DadosOperacao com o Manager; + # ManagerWorker nunca mais escreve ACK no JSON central da operação. + health_map = cls._dict(cls.get(CtxKey.DadosHealthWorker, {})) + manager_map = cls._dict(cls.get(CtxKey.DadosManagerWorker, {})) + + health_cmd = str(health_map.get("ultimo_comando_mapa", "") or "") + health_rev = cls._int(health_map.get("mapa_revisao_aplicada", 0), 0) + health_status = str(health_map.get("mapa_status", "") or "") + + manager_cmd = str(manager_map.get("mapa_command_id_aplicado", "") or "") + manager_rev = cls._int(manager_map.get("mapa_revisao_aplicada", 0), 0) + manager_qtd = cls._int(manager_map.get("mapa_quantidade_pontos_aplicada", 0), 0) + manager_status = str(manager_map.get("mapa_status", "") or "") primeiro_esperado = operacao.get("mapa_primeiro_idx_esperado", None) ultimo_esperado = operacao.get("mapa_ultimo_idx_esperado", None) @@ -691,13 +719,17 @@ class ContextoGlobalRedis: not mapa_barreira_critica and mapa_confirmado_flag and status_mapa == "aplicado" - and status_mpc == "aplicado" and revisao_solicitada > 0 - and revisao_solicitada == revisao_processada == revisao_mpc + and health_rev == revisao_solicitada + and manager_rev == revisao_solicitada + and command_solicitado != "" + and health_cmd == command_solicitado + and manager_cmd == command_solicitado + and health_status == "aplicado" + and manager_status == "aplicado" and qtd_esperada > 0 and qtd_publicada == qtd_esperada - and qtd_processada == qtd_esperada - and qtd_mpc == qtd_esperada + and manager_qtd == qtd_esperada and indices_extremos_ok ) ) @@ -705,9 +737,10 @@ class ContextoGlobalRedis: if not mapa_mpc_ok: motivos.append( "Mapa/MPC ainda não confirmado: " - f"rev solicitada/processada/mpc={revisao_solicitada}/{revisao_processada}/{revisao_mpc}; " - f"pontos publicados/esperados/processados/mpc={qtd_publicada}/{qtd_esperada}/{qtd_processada}/{qtd_mpc}; " - f"status={status_mapa}/{status_mpc}" + f"rev solicitada/health/manager={revisao_solicitada}/{health_rev}/{manager_rev}; " + f"cmd solicitado/health/manager={command_solicitado}/{health_cmd}/{manager_cmd}; " + f"pontos publicados/esperados/manager={qtd_publicada}/{qtd_esperada}/{manager_qtd}; " + f"status={status_mapa}/{health_status}/{manager_status}" ) detalhes_regras["mapa_mpc"] = { @@ -716,14 +749,17 @@ class ContextoGlobalRedis: "confirmado": mapa_confirmado_flag, "barreira_critica": mapa_barreira_critica, "status_mapa": status_mapa, - "status_mpc": status_mpc, + "status_health": health_status, + "status_manager": manager_status, "revisao_solicitada": revisao_solicitada, - "revisao_processada": revisao_processada, - "revisao_mpc": revisao_mpc, + "revisao_health": health_rev, + "revisao_manager": manager_rev, + "command_solicitado": command_solicitado, + "command_health": health_cmd, + "command_manager": manager_cmd, "pontos_publicados": qtd_publicada, "pontos_esperados": qtd_esperada, - "pontos_processados": qtd_processada, - "pontos_mpc": qtd_mpc, + "pontos_manager": manager_qtd, "primeiro_idx_publicado": primeiro_publicado, "ultimo_idx_publicado": ultimo_publicado, "primeiro_idx_esperado": primeiro_esperado, @@ -1293,11 +1329,12 @@ class ContextoGlobalRedis: operacao.get("inicializacao_operacao_executada", False) ) if not inicializacao_executada: - cls._iniciar_operacao() - updates_base["tempo_aguardando"] = agora - updates_base["ja_entrou_em_andamento"] = False - updates_base["inicializacao_operacao_executada"] = True - ja_entrou_em_andamento = False + inicio_emitido = cls._iniciar_operacao() + if inicio_emitido: + updates_base["tempo_aguardando"] = agora + updates_base["ja_entrou_em_andamento"] = False + updates_base["inicializacao_operacao_executada"] = True + ja_entrou_em_andamento = False # 4) Bloqueio automático: confirma apenas glitches muito curtos. # As proteções de cada módulo continuam responsáveis por suas próprias @@ -1488,17 +1525,84 @@ class ContextoGlobalRedis: @classmethod def _iniciar_operacao(cls): + """ + Inicia a execução somente depois de o mapa MPC estar confirmado. + + MAPSYNC V3: + - evita o antigo envio de um mapa com cmd/rev=None antes da barreira; + - o mapa já foi carregado pela transação crítica; + - aqui enviamos apenas o comando de INÍCIO para a instância confirmada. + + Retorna True quando o comando de início foi efetivamente emitido. + """ + operacao = cls._dict(cls.get_operacao()) + dados_dir = cls._dict(operacao.get("Dir", {})) + tipo_controle = cls._int( + dados_dir.get( + "tipo_controle", + TiposControladorDirecional.Manual.value, + ), + TiposControladorDirecional.Manual.value, + ) + + if tipo_controle == TiposControladorDirecional.MPC.value: + revisao = cls._int(operacao.get("mapa_revisao_solicitada", 0), 0) + command_id = str(operacao.get("mapa_command_id", "") or "") + qtd_esperada = cls._int( + operacao.get("mapa_quantidade_pontos_esperada", 0), 0 + ) + + health = cls._dict(cls.get(CtxKey.DadosHealthWorker, {})) + manager = cls._dict(cls.get(CtxKey.DadosManagerWorker, {})) + + health_ok = ( + command_id != "" + and str(health.get("ultimo_comando_mapa", "") or "") == command_id + and cls._int(health.get("mapa_revisao_aplicada", 0), 0) == revisao + and str(health.get("mapa_status", "") or "") == "aplicado" + ) + + manager_ok = ( + command_id != "" + and str(manager.get("mapa_command_id_aplicado", "") or "") == command_id + and cls._int(manager.get("mapa_revisao_aplicada", 0), 0) == revisao + and cls._int(manager.get("mapa_quantidade_pontos_aplicada", 0), 0) == qtd_esperada + and str(manager.get("mapa_status", "") or "") == "aplicado" + ) + + mapa_ok = ( + revisao > 0 + and qtd_esperada > 0 + and cls._bool(operacao.get("mapa_confirmado", False)) + and not cls._bool(operacao.get("mapa_barreira_critica", False)) + and str(operacao.get("mapa_status", "") or "") == "aplicado" + and health_ok + and manager_ok + ) + + if not mapa_ok: + return False + cls.publicar_comando( CmdKey.WeedWorkerRx, {"cmd": WeedWorkerCommandType.ReiniciarDeteccoes.value}, ) - cls._atualizar_mpc(do_manager=False, iniciar_operacao=True) + # O mapa já está confirmado no Manager. Não reenviamos params/mapa aqui. + # O comando sem params apenas inicia a instância existente. + iniciado_enviado = cls.publicar_comando( + CmdKey.ManagerWorkerRx, + {"cmd": ManagerWorkerCommandType.IniciarMPC.value}, + ) + + if not iniciado_enviado: + return False cls.atualizar_ctx_dict( CtxKey.DadosOperacao, tempo_aguardando=time.time(), ) + return True @classmethod def _finalizar_operacao(cls): @@ -1749,8 +1853,6 @@ class ContextoGlobalRedis: atualizado = cls.atualizar_ctx_dict( CtxKey.DadosOperacao, pontos_mapa=pontos_info, - mapa_status_mpc="pendente", - mapa_quantidade_pontos_processada=0, ) if not atualizado: @@ -1760,59 +1862,167 @@ class ContextoGlobalRedis: } # ========================================================== - # 5. Enviar EXATAMENTE este mapa ao Manager + # 5/6. Entrega confiável Health -> Manager # ========================================================== - mostrar_log( - f"[MAPSYNC] enviando mapa ao Manager | " - f"cmd={command_id} | rev={map_revision} | " - f"pontos={qtd_esperada}" - ) + # Redis Pub/Sub é efêmero. Se uma publicação cair exatamente + # durante uma janela ruim do listener, ela não fica enfileirada. + # + # Em vez de esperar 4 s e pedir ao C# um NOVO command_id, o Health + # retransmite o MESMO snapshot com o MESMO command_id/revisão. + # O Manager V3 é idempotente: replays da mesma revisão não + # reinicializam o MPC desnecessariamente e apenas reafirmam o ACK. + intervalo_reenvio_s = 0.35 + timeout_total_s = 5.00 + timeout_aplicando_s = 3.00 + poll_s = 0.05 - mpc_enviado = cls._atualizar_mpc( - do_manager=False, - iniciar_operacao=False, - operacao=operacao, - mapa=pontos_info, - map_revision=map_revision, - command_id=command_id, - expected_point_count=qtd_esperada, - ) + inicio_entrega = time.monotonic() + limite_total = inicio_entrega + timeout_total_s + proximo_envio_em = inicio_entrega + numero_envio = 0 + aplicado_mpc = False + aplicando_desde = None - if not mpc_enviado: + def _estado_manager_mapsync(): + estado = cls._dict( + cls.get(CtxKey.DadosManagerWorker, {}) + ) return { - "sucesso": False, - "erro": "Falha ao enviar mapa ao Manager/MPC", + "cmd": str( + estado.get("mapa_command_id_aplicado", "") or "" + ), + "rev": cls._int( + estado.get("mapa_revisao_aplicada", 0), 0 + ), + "qtd": cls._int( + estado.get("mapa_quantidade_pontos_aplicada", 0), 0 + ), + "status": str( + estado.get("mapa_status", "") or "" + ), + "erro": str( + estado.get("mapa_erro", "") or "" + ), } - # ========================================================== - # 6. Esperar Manager confirmar que realmente carregou - # ========================================================== - aplicado_mpc = cls._aguardar_mapa_mpc_aplicado( - command_id=command_id, - map_revision=map_revision, - expected_point_count=qtd_esperada, - timeout_s=2.5, - ) + while time.monotonic() < limite_total: + agora_mono = time.monotonic() + mgr = _estado_manager_mapsync() + + mesma_revisao = ( + mgr["rev"] == cls._int(map_revision, 0) + and mgr["qtd"] == qtd_esperada + ) + mesmo_comando = ( + mgr["cmd"] == str(command_id) + ) + + # ------------------------------------------------------ + # Manager recebeu a revisão e já está construindo o MPC. + # Nesse estado NÃO retransmitimos. Damos uma janela curta + # para finalizar e escrever o ACK, evitando fila de replays. + # Aceitamos "aplicando" da mesma revisão mesmo se o cmd for + # de uma tentativa anterior; ao concluir, o replay corrente + # fará rebind idempotente do command_id. + # ------------------------------------------------------ + if mesma_revisao and mgr["status"] == "aplicando": + if aplicando_desde is None: + aplicando_desde = agora_mono + mostrar_log( + f"[MAPSYNC/TX-MANAGER] Manager aplicando; " + f"retransmissão pausada | cmd={command_id} | " + f"cmd_manager={mgr['cmd']} | rev={map_revision} | " + f"pontos={qtd_esperada}" + ) + + if agora_mono - aplicando_desde > timeout_aplicando_s: + return { + "sucesso": False, + "erro": ( + f"Manager permaneceu em 'aplicando' por mais de " + f"{timeout_aplicando_s:.1f}s | " + f"cmd={command_id} | rev={map_revision} | " + f"pontos={qtd_esperada}" + ), + } + + time.sleep(poll_s) + continue + + aplicando_desde = None + + if mesma_revisao and mgr["status"] == "erro" and mesmo_comando: + return { + "sucesso": False, + "erro": ( + "Manager rejeitou explicitamente o snapshot | " + f"cmd={command_id} | rev={map_revision} | " + f"erro={mgr['erro']}" + ), + } + + if mesma_revisao and mgr["status"] == "aplicado": + if mesmo_comando: + aplicado_mpc = True + mostrar_log( + f"[MAPSYNC/TX-MANAGER] ACK confirmado | " + f"envios={numero_envio} | cmd={command_id} | " + f"rev={map_revision} | pontos={qtd_esperada}" + ) + break + + # Mesma revisão já aplicada sob command_id anterior. + # Enviamos o command atual imediatamente para o Manager + # executar apenas o rebind idempotente do ACK. + proximo_envio_em = min(proximo_envio_em, agora_mono) + + if agora_mono >= proximo_envio_em: + numero_envio += 1 + + mostrar_log( + f"[MAPSYNC/TX-MANAGER] envio={numero_envio} | " + f"cmd={command_id} | rev={map_revision} | " + f"pontos={qtd_esperada}" + ) + + mpc_enviado = cls._atualizar_mpc( + do_manager=False, + iniciar_operacao=False, + operacao=operacao, + mapa=pontos_info, + map_revision=map_revision, + command_id=command_id, + expected_point_count=qtd_esperada, + ) + + if not mpc_enviado: + return { + "sucesso": False, + "erro": "Falha ao enviar mapa ao Manager/MPC", + } + + proximo_envio_em = agora_mono + intervalo_reenvio_s + + time.sleep(poll_s) if not aplicado_mpc: + mgr = _estado_manager_mapsync() return { "sucesso": False, "erro": ( - f"Manager não confirmou aplicação exata do mapa | " + f"Manager não confirmou aplicação exata do mapa após " + f"{numero_envio} envio(s) e {timeout_total_s:.1f}s | " f"cmd={command_id} | rev={map_revision} | " - f"pontos={qtd_esperada}" + f"pontos={qtd_esperada} | " + f"manager_status={mgr['status']} | " + f"manager_cmd={mgr['cmd']} | manager_rev={mgr['rev']}" ), } # ========================================================== - # 7. Só agora considerar revisão processada + # 7. ACK do Health será publicado pelo próprio HealthWorker. + # Não espelhamos estado transacional em DadosOperacao. # ========================================================== - cls.atualizar_ctx_dict( - CtxKey.DadosOperacao, - mapa_revisao_processada=map_revision, - mapa_quantidade_pontos_processada=len(pontos_info), - ) - return { "sucesso": True, "quantidade": len(pontos_info), @@ -1918,36 +2128,20 @@ class ContextoGlobalRedis: command_id, map_revision, expected_point_count=None, - timeout_s=2.5, + timeout_s=4.0, ): + """Espera o ACK de propriedade exclusiva do ManagerWorker.""" limite = time.time() + max(0.1, float(timeout_s)) while time.time() < limite: - operacao = cls.get_operacao() - - cmd_aplicado = str( - operacao.get( - "mapa_command_id_mpc_aplicado", - "" - ) - ) - - revisao_aplicada = operacao.get( - "mapa_revisao_mpc_aplicada", - None - ) + manager = cls._dict(cls.get(CtxKey.DadosManagerWorker, {})) + cmd_aplicado = str(manager.get("mapa_command_id_aplicado", "") or "") + revisao_aplicada = manager.get("mapa_revisao_aplicada", None) qtd_aplicada = cls._int( - operacao.get("mapa_quantidade_pontos_mpc_aplicada", 0), - 0, - ) - - status = str( - operacao.get( - "mapa_status_mpc", - "" - ) + manager.get("mapa_quantidade_pontos_aplicada", 0), 0 ) + status = str(manager.get("mapa_status", "") or "") quantidade_ok = ( expected_point_count is None @@ -1963,10 +2157,7 @@ class ContextoGlobalRedis: ): return True - if ( - cmd_aplicado == str(command_id) - and status == "erro" - ): + if cmd_aplicado == str(command_id) and status == "erro": return False time.sleep(0.05)