ajustado envio dos dados do mapa da trajetorio c# pelo redis para o python com confirmacao

This commit is contained in:
Diego Freitas 2026-09-09 16:58:34 -03:00
parent 0319fd204d
commit f48741f92d
5 changed files with 703 additions and 66 deletions

View File

@ -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<GPSModel> pontos) GerarTrajetoriaRetornoAutomatico(GPSModel baseGps)

View File

@ -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<bool> SolicitarAtualizacaoMapaAsync(OperacaoModel op)
public static async Task<bool> 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<bool> 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<bool> AtualizarPontosMapa(OperacaoModel op)
private static async Task<bool> 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<bool> AguardarMapaAplicadoAsync(string commandId, long revision, TimeSpan timeout)

View File

@ -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():

View File

@ -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}"
)

View File

@ -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
# ============================================================