diff --git a/AgroBase/AgroBase/Models/TrajetoriaMapaOperacaoModel.cs b/AgroBase/AgroBase/Models/TrajetoriaMapaOperacaoModel.cs index 41dc01ccb..32061952d 100644 --- a/AgroBase/AgroBase/Models/TrajetoriaMapaOperacaoModel.cs +++ b/AgroBase/AgroBase/Models/TrajetoriaMapaOperacaoModel.cs @@ -8,6 +8,7 @@ using System; using System.Collections.Generic; using System.Globalization; using System.Linq; +using System.Threading.Tasks; using static AgroBase.Models.CorredorTrajetoriaModel; using static AgroBase.Models.Enums; @@ -4654,9 +4655,146 @@ namespace AgroBase.Models LoopAtualizaDadosCore(); - RedisService.AtualizarCampos(CtxKey.DadosOperacao, ("configurado", false)); + /* + * Primeiro publicar TODOS os parâmetros e a nova _TrajetoriaFixa + * no Redis, mas sem disparar o sync comum e SEM liberar a operação. + */ + HealthWorkerService.AtualizarDadosOperacao( + forcar: true, + sincronizarMapa: false, + configurado: false + ); - RedisService.Publish(CmdKey.ManagerWorkerRx, JsonConvert.SerializeObject(new { cmd = ManagerWorkerCommandType.IniciarMPC })); + /* + * A sincronização crítica ocorre fora do fluxo principal. + * + * Enquanto ela não terminar: + * + * configurado = false + * + * portanto o rover não deve navegar. + */ + _ = Task.Run( + async () => + await SincronizarEIniciarRetornoAsync(op) + ); + } + + private async Task SincronizarEIniciarRetornoAsync(OperacaoModel op) + { + try + { + Variaveis.MostrarLog( + $"[RETORNO/MAPSYNC] Iniciando sincronização crítica | " + + $"pontos={_TrajetoriaFixa?.Count ?? 0}" + ); + + bool mapaAplicado = + await HealthWorkerService + .SincronizarMapaCriticoAsync(op) + .ConfigureAwait(false); + + /* + * Durante a espera a operação pode ter sido substituída, + * reiniciada ou encerrada. + */ + if (!ReferenceEquals( + op, + Variaveis.OperacaoEmAndamento)) + { + Variaveis.MostrarLog( + "[RETORNO/MAPSYNC] Operação mudou durante sincronização." + ); + + return; + } + + if (!mapaAplicado) + { + RedisService.AtualizarCampos( + CtxKey.DadosOperacao, + ("configurado", false) + ); + + Variaveis.MostrarLog( + "[RETORNO/MAPSYNC] FALHA: " + + "novo mapa não foi confirmado pelo MPC. " + + "Retorno permanece bloqueado." + ); + + op.Sensoriamento?.InserirLog( + T_Code.Trj, + StatusModulo.Falha, + 0, + "Retorno bloqueado: MPC não confirmou a nova trajetória." + ); + + return; + } + + /* + * Nesse momento temos: + * + * C# -> nova trajetória + * Health -> converteu + * Manager -> inicializou MPC + * Health -> ACK + * C# -> recebeu ACK + * + * Só agora permitimos a operação. + */ + HealthWorkerService.AtualizarDadosOperacao( + forcar: true, + sincronizarMapa: false, + configurado: true + ); + + Variaveis.MostrarLog( + $"[RETORNO/MAPSYNC] Mapa confirmado pelo MPC | " + + $"pontos={_TrajetoriaFixa?.Count ?? 0}" + ); + + /* + * O primeiro IniciarMPC, enviado pelo Health durante o + * map-sync, carregou o mapa com iniciar_operacao=false. + * + * Este segundo comando, SEM params, é o que efetivamente + * libera/inicia a execução usando o mapa já confirmado. + */ + RedisService.Publish( + CmdKey.ManagerWorkerRx, + JsonConvert.SerializeObject( + new + { + cmd = + ManagerWorkerCommandType.IniciarMPC + } + ) + ); + + Variaveis.MostrarLog( + "[RETORNO/MAPSYNC] Retorno liberado." + ); + } + catch (Exception ex) + { + RedisService.AtualizarCampos( + CtxKey.DadosOperacao, + ("configurado", false) + ); + + Variaveis.MostrarLog( + $"[RETORNO/MAPSYNC] Erro crítico: {ex.Message}" + ); + + op?.Sensoriamento?.InserirLog( + T_Code.Trj, + StatusModulo.Falha, + 0, + "Erro sincronizando trajetória de retorno: " + + ex.Message + ); + } } private (bool valido, List pontos) GerarTrajetoriaRetornoAutomatico(GPSModel baseGps) diff --git a/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs b/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs index 1a09aaf85..14b8bac52 100644 --- a/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs +++ b/AgroBase/AgroBase/Services/Operadores/HealthWorkerService.cs @@ -284,7 +284,7 @@ namespace AgroBase.Services.Operadores } } - public static void AtualizarDadosOperacao(bool forcar = false) + public static void AtualizarDadosOperacao(bool forcar = false, bool sincronizarMapa = true, bool configurado = true) { var op = Variaveis.OperacaoEmAndamento; @@ -312,7 +312,7 @@ namespace AgroBase.Services.Operadores RedisService.AtualizarCampos( CtxKey.DadosOperacao, ("id", op.ID ?? ""), - ("configurado", true), + ("configurado", configurado), ("modo", op.Parametros?.Modo ?? ModoOperacao.NaoDefinido), ("tempo_aguardar_inicio_operacao", op.TempoIniciarOperacao), ("tempo_aguardar_retomada", op.TempoAguardarRetomada), @@ -400,13 +400,32 @@ namespace AgroBase.Services.Operadores ) ); - _ = SolicitarAtualizacaoMapaAsync(op).ContinueWith( - task => + if (sincronizarMapa) + { + _ = Task.Run(async () => { - Variaveis.MostrarLog($"[HealthWorkerService] Erro atualizando mapa: {task.Exception?.GetBaseException().Message}"); - }, - TaskContinuationOptions.OnlyOnFaulted - ); + try + { + bool sucesso = await SolicitarAtualizacaoMapaAsync( + op, + aguardarGate: false + ).ConfigureAwait(false); + + if (!sucesso) + { + Variaveis.MostrarLog( + "[MAPSYNC] Sincronização não crítica não foi confirmada." + ); + } + } + catch (Exception ex) + { + Variaveis.MostrarLog( + $"[MAPSYNC] Erro em sincronização não crítica: {ex.Message}" + ); + } + }); + } //Console.WriteLine("Dados de operacao atualizados com sucesso!"); @@ -420,53 +439,179 @@ namespace AgroBase.Services.Operadores private static readonly SemaphoreSlim _mapaSync = new SemaphoreSlim(1, 1); - private static async Task SolicitarAtualizacaoMapaAsync(OperacaoModel op) + public static async Task SincronizarMapaCriticoAsync(OperacaoModel op) { - Variaveis.MostrarLog($"[HealthWorkerService] SolicitarAtualizacaoMapaAsync solicitado"); - if (!await _mapaSync.WaitAsync(0)) + return await SolicitarAtualizacaoMapaAsync( + op, + aguardarGate: true + ).ConfigureAwait(false); + } + + private static async Task SolicitarAtualizacaoMapaAsync(OperacaoModel op, bool aguardarGate = false) + { + if ( + op == null || + !(op.Parametros?.ControleAutomatico ?? false) || + !(op.Trajetoria?._TrajetoriaFixa?.Any() ?? false) + ) + { + Variaveis.MostrarLog( + "[MAPSYNC] Solicitação recusada: operação/mapa inválido." + ); + return false; + } + + bool gateObtido = false; try { + if (aguardarGate) + { + /* + * Sincronização CRÍTICA: + * jamais descartar porque outra sincronização + * está em andamento. + */ + await _mapaSync.WaitAsync().ConfigureAwait(false); + gateObtido = true; + } + else + { + /* + * Sincronização normal: + * pode ser coalescida. + */ + gateObtido = + await _mapaSync.WaitAsync(0).ConfigureAwait(false); + + if (!gateObtido) + { + Variaveis.MostrarLog( + "[MAPSYNC] Sincronização comum ignorada: " + + "outra atualização já está em andamento." + ); + + return false; + } + } + + if (!ReferenceEquals(op, Variaveis.OperacaoEmAndamento)) + { + Variaveis.MostrarLog( + "[MAPSYNC] Operação mudou durante solicitação." + ); + + return false; + } + + /* + * A REVISÃO identifica o MAPA. + * + * Ela permanece a mesma durante todos os retries. + */ + long mapRevision = + DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); + + string operationId = op.ID ?? ""; + + RedisService.AtualizarCampos( + CtxKey.DadosOperacao, + ("mapa_revisao_solicitada", mapRevision), + ("mapa_operacao_id", operationId), + ("mapa_status", "pendente") + ); + for (int tentativa = 1; tentativa <= 3; tentativa++) { - if (!ReferenceEquals(op, Variaveis.OperacaoEmAndamento)) + if (!ReferenceEquals( + op, + Variaveis.OperacaoEmAndamento)) + { return false; + } - bool atualizado = await AtualizarPontosMapa(op); + /* + * commandId identifica a TENTATIVA de transporte. + * + * A revisão continua a mesma. + */ + string commandId = + Guid.NewGuid().ToString("N"); + + RedisService.AtualizarCampos( + CtxKey.DadosOperacao, + ("mapa_command_id", commandId) + ); + + Variaveis.MostrarLog( + $"[MAPSYNC] Tentativa={tentativa} | " + + $"cmd={commandId} | " + + $"rev={mapRevision} | " + + $"pontos={op.Trajetoria._TrajetoriaFixa.Count}" + ); + + bool atualizado = + await AtualizarPontosMapa( + op, + commandId, + operationId, + mapRevision + ).ConfigureAwait(false); if (atualizado) + { + RedisService.AtualizarCampos( + CtxKey.DadosOperacao, + ("mapa_status", "aplicado") + ); + + Variaveis.MostrarLog( + $"[MAPSYNC] SUCESSO | " + + $"rev={mapRevision} | " + + $"tentativa={tentativa}" + ); + return true; + } + + Variaveis.MostrarLog( + $"[MAPSYNC] Falha/timeout | " + + $"tentativa={tentativa} | " + + $"cmd={commandId} | " + + $"rev={mapRevision}" + ); if (tentativa < 3) - await Task.Delay(200); - Variaveis.MostrarLog($"[HealthWorkerService] SolicitarAtualizacaoMapaAsync. Tentativa={tentativa}"); + { + await Task.Delay(250) + .ConfigureAwait(false); + } } + RedisService.AtualizarCampos( + CtxKey.DadosOperacao, + ("mapa_status", "erro") + ); + return false; } finally { - _mapaSync.Release(); + if (gateObtido) + _mapaSync.Release(); } } - private static async Task AtualizarPontosMapa(OperacaoModel op) + private static async Task AtualizarPontosMapa(OperacaoModel op, string commandId, string operationId, long mapRevision) { - if (!(op?.Parametros?.ControleAutomatico ?? false) || !(op?.Trajetoria?._TrajetoriaFixa?.Any() ?? false)) + if ( + !(op?.Parametros?.ControleAutomatico ?? false) || + !(op?.Trajetoria?._TrajetoriaFixa?.Any() ?? false) + ) + { return false; - - var commandId = Guid.NewGuid().ToString("N"); - var operationId = op?.ID ?? ""; - var mapRevision = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(); - - RedisService.AtualizarCampos( - CtxKey.DadosOperacao, - ("mapa_revisao_solicitada", mapRevision), - ("mapa_operacao_id", operationId), - ("mapa_command_id", commandId), - ("mapa_status", "pendente") - ); + } var comando = new { @@ -481,8 +626,63 @@ namespace AgroBase.Services.Operadores JsonConvert.SerializeObject(comando) ); - bool sucesso = await AguardarMapaAplicadoAsync(commandId, mapRevision, TimeSpan.FromSeconds(3)); - return assinantes > 0 && sucesso; + Variaveis.MostrarLog( + $"[MAPSYNC/TX] cmd={commandId} | " + + $"rev={mapRevision} | " + + $"subs={assinantes} | " + + $"pontos={op.Trajetoria._TrajetoriaFixa.Count}" + ); + + /* + * Se ninguém ouviu o Publish, não existe motivo para + * esperar quatro segundos por um ACK impossível. + */ + if (assinantes <= 0) + { + Variaveis.MostrarLog( + $"[MAPSYNC] Nenhum assinante HealthWorker | " + + $"cmd={commandId}" + ); + + return false; + } + + /* + * O Python agora espera o ACK do Manager (~2,5 s). + * Portanto 3 s ficou apertado demais. + */ + bool sucesso = await AguardarMapaAplicadoAsync( + commandId, + mapRevision, + TimeSpan.FromSeconds(4) + ).ConfigureAwait(false); + + if (!sucesso) + { + string status = + RedisService.GetField( + CtxKey.DadosHealthWorker, + "mapa_status", + "" + ); + + string erro = + RedisService.GetField( + CtxKey.DadosHealthWorker, + "mapa_erro", + "" + ); + + Variaveis.MostrarLog( + $"[MAPSYNC] ACK não confirmado | " + + $"cmd={commandId} | " + + $"rev={mapRevision} | " + + $"status={status} | " + + $"erro={erro}" + ); + } + + return sucesso; } private static async Task AguardarMapaAplicadoAsync(string commandId, long revision, TimeSpan timeout) diff --git a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/main.py b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/main.py index 3cbc60b89..6cc85d9f9 100644 --- a/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/main.py +++ b/AgroBase/AgroBase/bin/x64/Debug/Python/Scripts/workers/health_worker/main.py @@ -93,14 +93,22 @@ def main(): operation_id = str(dados.get("operation_id", "")) map_revision = dados.get("map_revision") + mostrar_log( + f"[MAPSYNC/RX] command_id={command_id} | " + f"operacao={operation_id} | revisao={map_revision}" + ) + try: resultado = ContextoGlobalRedis._atualizar_pontos_mapa( operation_id=operation_id, map_revision=map_revision, + command_id=command_id, ) if not resultado.get("sucesso", False): - raise RuntimeError(resultado.get("erro", "Falha desconhecida")) + raise RuntimeError( + resultado.get("erro", "Falha desconhecida") + ) ContextoGlobalRedis.atualizar_ctx_dict( CtxKey.DadosHealthWorker, @@ -131,7 +139,9 @@ def main(): ) mostrar_log( - f"Erro ao aplicar mapa | command_id={command_id} | {e}" + f"Erro ao aplicar mapa | " + f"command_id={command_id} | " + f"revisao={map_revision} | {e}" ) def inicializar(): 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 c6163bd56..8478c0718 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 @@ -54,39 +54,126 @@ def main(): time.sleep(delay_corrigido) def redis_callback(dados): - acao = ManagerWorkerCommandType(dados.get("cmd", 0)) + acao = ManagerWorkerCommandType( + dados.get("cmd", 0) + ) if acao == ManagerWorkerCommandType.IniciarMPC: - try: - _p = dados.get("params") + _p = dados.get("params") + + command_id = None + map_revision = None + + try: + from manager_worker.modulos.mpc import ( + inicializar as iniciar_mpc, + get_iniciado, + ) - from manager_worker.modulos.mpc import inicializar as iniciar_mpc, get_iniciado _mpc_iniciado = get_iniciado() + if _mpc_iniciado and _p is None: - ContextoGlobalRedis()._atualizar_mpc(do_manager=True, iniciar_operacao=True) + ContextoGlobalRedis()._atualizar_mpc( + do_manager=True, + iniciar_operacao=True + ) + elif _p is not None: parametros = _p.get("parametros") mapa = _p.get("mapa") - iniciar_operacao = _p.get("iniciar_operacao", False) + iniciar_operacao = _p.get( + "iniciar_operacao", + False + ) p_ref = _p.get("p_ref", []) - iniciar_mpc(parametros, mapa, p_ref, iniciar_operacao) + + command_id = _p.get( + "command_id" + ) + + map_revision = _p.get( + "map_revision" + ) + + 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 ( + 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_status_mpc="aplicado", + mapa_erro_mpc="", + mapa_mpc_aplicado_em=time.time(), + ) + + mostrar_log( + f"[MAPSYNC/MANAGER] mapa aplicado | " + f"cmd={command_id} | " + f"rev={map_revision} | " + f"pontos={len(mapa or [])}" + ) + else: - ContextoGlobalRedis()._atualizar_mpc(do_manager=True, iniciar_operacao=False) + ContextoGlobalRedis()._atualizar_mpc( + do_manager=True, + iniciar_operacao=False + ) + except Exception as e: - mostrar_log(f"Erro ao Iniciar MPC: {e}") + if ( + command_id is not None + and map_revision is not None + ): + try: + ContextoGlobalRedis.atualizar_ctx_dict( + CtxKey.DadosOperacao, + mapa_command_id_mpc_aplicado=str(command_id), + mapa_revisao_mpc_aplicada=map_revision, + mapa_status_mpc="erro", + mapa_erro_mpc=str(e), + ) + except Exception: + pass + + mostrar_log( + f"Erro ao Iniciar MPC: {e}" + ) + elif acao == ManagerWorkerCommandType.ReiniciarMPC: try: from manager_worker.modulos.mpc import resetar as resetar_mpc _mpc_iniciado = resetar_mpc() except Exception as e: mostrar_log(f"Erro ao Resetar MPC: {e}") + elif acao == ManagerWorkerCommandType.AtualizarDadosControle: - #mostrar_log("Adicionado na fila") fila.adicionar(acao) - #manager.executar(not loop_fixo) - pass + else: - mostrar_log(f"⚠️ Comando desconhecido: {acao.name}") + mostrar_log( + f"⚠️ Comando desconhecido: {acao.name}" + ) 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 779d12510..09cc72d7b 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 @@ -1391,24 +1391,70 @@ class ContextoGlobalRedis: ) @classmethod - def _atualizar_mpc(cls, do_manager=False, iniciar_operacao=False): + def _atualizar_mpc( + cls, + do_manager=False, + iniciar_operacao=False, + operacao=None, + mapa=None, + map_revision=None, + command_id=None, + ): try: - dados_dir = cls._dict(cls.get_operacao().get("Dir", {})) + # Snapshot único. + if operacao is None: + operacao = cls.get_operacao() + + operacao = cls._dict(operacao) + + dados_dir = cls._dict( + operacao.get("Dir", {}) + ) + tipo_controle = cls._int( - dados_dir.get("tipo_controle", TiposControladorDirecional.Manual.value), + dados_dir.get( + "tipo_controle", + TiposControladorDirecional.Manual.value + ), TiposControladorDirecional.Manual.value, ) if tipo_controle != TiposControladorDirecional.MPC.value: return False - dados_mpc = cls._dict(dados_dir.get("mpc", {})) - mapa = cls.get_operacao().get("pontos_mapa", []) - p_ref = cls.get_operacao().get("ponto_mapa_ref", []) + dados_mpc = cls._dict( + dados_dir.get("mpc", {}) + ) + + if mapa is None: + mapa = operacao.get( + "pontos_mapa", + [] + ) + + p_ref = operacao.get( + "ponto_mapa_ref", + [] + ) + + if not mapa: + mostrar_log( + "[MAPSYNC] MPC não atualizado: mapa vazio" + ) + return False if do_manager: - from manager_worker.modulos.mpc import inicializar as iniciar_mpc - iniciar_mpc(dados_mpc, mapa, p_ref, iniciar_operacao) + from manager_worker.modulos.mpc import ( + inicializar as iniciar_mpc + ) + + iniciar_mpc( + dados_mpc, + mapa, + p_ref, + iniciar_operacao + ) + else: cls.publicar_comando( CmdKey.ManagerWorkerRx, @@ -1419,6 +1465,10 @@ class ContextoGlobalRedis: "mapa": mapa, "p_ref": p_ref, "iniciar_operacao": iniciar_operacao, + + # Identidade da transação + "map_revision": map_revision, + "command_id": command_id, }, }, ) @@ -1426,27 +1476,69 @@ class ContextoGlobalRedis: return True except Exception as e: - mostrar_log(f"Erro ao atualizar MPC: {e}") + mostrar_log( + f"Erro ao atualizar MPC: {e}" + ) return False @classmethod - def _atualizar_pontos_mapa(cls, operation_id=None, map_revision=None): + def _atualizar_pontos_mapa( + cls, + operation_id=None, + map_revision=None, + command_id=None, + ): try: + # ========================================================== + # 1. Snapshot único da operação + # ========================================================== + operacao = cls.get_operacao() - operacao_atual = str(operacao.get("mapa_operacao_id", "")) + operacao_atual = str( + operacao.get("mapa_operacao_id", "") + ) if operation_id and operacao_atual != operation_id: return { "sucesso": False, "erro": ( - f"Operação divergente: esperado={operation_id}, " + f"Operação divergente: " + f"esperado={operation_id}, " f"atual={operacao_atual}" ), } - # Conversão atual dos pontos... - pontos_info = cls._converter_pontos_mapa(operacao) + # ========================================================== + # 2. Garantir que o comando pertence à revisão atual + # ========================================================== + + revisao_atual = operacao.get( + "mapa_revisao_solicitada", + None + ) + + if ( + map_revision is not None + and revisao_atual is not None + and int(revisao_atual) != int(map_revision) + ): + return { + "sucesso": False, + "erro": ( + f"Revisão de mapa obsoleta: " + f"esperada={map_revision}, " + f"atual={revisao_atual}" + ), + } + + # ========================================================== + # 3. Converter exatamente este snapshot + # ========================================================== + + pontos_info = cls._converter_pontos_mapa( + operacao + ) if not pontos_info: return { @@ -1454,20 +1546,81 @@ class ContextoGlobalRedis: "erro": "Nenhum ponto válido foi produzido", } + mostrar_log( + f"[MAPSYNC] mapa convertido | " + f"cmd={command_id} | " + f"rev={map_revision} | " + f"pontos={len(pontos_info)}" + ) + + # ========================================================== + # 4. Gravar pontos convertidos + # ========================================================== + atualizado = cls.atualizar_ctx_dict( CtxKey.DadosOperacao, pontos_mapa=pontos_info, - mapa_revisao_processada=map_revision, + mapa_status_mpc="pendente", ) - cls._atualizar_mpc(do_manager=False, iniciar_operacao=False) - if not atualizado: return { "sucesso": False, "erro": "Falha ao gravar pontos_mapa no Redis", } + # ========================================================== + # 5. Enviar EXATAMENTE este mapa ao Manager + # ========================================================== + + mostrar_log( + f"[MAPSYNC] enviando mapa ao Manager | " + f"cmd={command_id} | rev={map_revision}" + ) + + mpc_enviado = cls._atualizar_mpc( + do_manager=False, + iniciar_operacao=False, + operacao=operacao, + mapa=pontos_info, + map_revision=map_revision, + command_id=command_id, + ) + + if not mpc_enviado: + return { + "sucesso": False, + "erro": "Falha ao enviar mapa ao Manager/MPC", + } + + # ========================================================== + # 6. Esperar Manager confirmar que realmente carregou + # ========================================================== + + aplicado_mpc = cls._aguardar_mapa_mpc_aplicado( + command_id=command_id, + map_revision=map_revision, + timeout_s=2.5, + ) + + if not aplicado_mpc: + return { + "sucesso": False, + "erro": ( + f"Manager não confirmou aplicação do mapa | " + f"cmd={command_id} | rev={map_revision}" + ), + } + + # ========================================================== + # 7. Só agora considerar revisão processada + # ========================================================== + + cls.atualizar_ctx_dict( + CtxKey.DadosOperacao, + mapa_revisao_processada=map_revision, + ) + return { "sucesso": True, "quantidade": len(pontos_info), @@ -1567,6 +1720,55 @@ class ContextoGlobalRedis: return pontos_info + @classmethod + def _aguardar_mapa_mpc_aplicado( + cls, + command_id, + map_revision, + timeout_s=2.5, + ): + 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 + ) + + status = str( + operacao.get( + "mapa_status_mpc", + "" + ) + ) + + if ( + cmd_aplicado == str(command_id) + and revisao_aplicada is not None + and int(revisao_aplicada) == int(map_revision) + and status == "aplicado" + ): + return True + + if ( + cmd_aplicado == str(command_id) + and status == "erro" + ): + return False + + time.sleep(0.05) + + return False + # ============================================================ # Mapeamento de parâmetros # ============================================================