centralizacao no corredor e resync de mapa

This commit is contained in:
Diego Freitas 2026-09-19 00:41:03 -03:00
parent 8af660c9b5
commit 5fb8caa88c
5 changed files with 974 additions and 269 deletions

1
.gitignore vendored
View File

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

View File

@ -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<long>(
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<long>(
CtxKey.DadosHealthWorker,
"mapa_revisao_aplicada",
0L
);
int qtdProcessada = RedisService.GetField<int>(
CtxKey.DadosOperacao,
"mapa_quantidade_pontos_processada",
0
string healthStatus = RedisService.GetField(
CtxKey.DadosHealthWorker,
"mapa_status",
""
);
long revisaoMpc = RedisService.GetField<long>(
CtxKey.DadosOperacao,
"mapa_revisao_mpc_aplicada",
string managerCmd = RedisService.GetField(
CtxKey.DadosManagerWorker,
"mapa_command_id_aplicado",
""
);
long managerRev = RedisService.GetField<long>(
CtxKey.DadosManagerWorker,
"mapa_revisao_aplicada",
0L
);
int qtdMpc = RedisService.GetField<int>(
CtxKey.DadosOperacao,
"mapa_quantidade_pontos_mpc_aplicada",
int managerQtd = RedisService.GetField<int>(
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;
}

View File

@ -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,6 +53,10 @@ def main():
mostrar_log(f"❌ Erro no loop do Manager: {e}")
latencia = time.time() - t0
# 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,
@ -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,6 +119,101 @@ def main():
f"esperado={expected_point_count}, recebido={quantidade_recebida}"
)
# 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
and expected_point_count is not None
)
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} | "
@ -115,30 +221,30 @@ def main():
f"pontos={len(mapa or [])}"
)
iniciar_mpc(
resultado_mpc = iniciar_mpc(
parametros,
mapa,
p_ref,
iniciar_operacao
)
# ---------------------------------------------
# ACK somente depois de inicializar()
# retornar sem exception.
# ---------------------------------------------
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.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(),
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(
@ -148,6 +254,10 @@ def main():
f"pontos={len(mapa or [])}"
)
finally:
if eh_mapsync:
_mapsync_em_andamento.clear()
else:
ContextoGlobalRedis()._atualizar_mpc(
do_manager=True,
@ -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()

View File

@ -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(
@ -2791,12 +3106,30 @@ class ControladorMPC:
if inversao_significativa:
filtrado = 0.0
else:
# 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))

View File

@ -236,54 +236,71 @@ 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)
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:
dados = cls.get(chave_str, {})
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 = {}
for campo, valor in novos_campos.items():
_set_nested(dados, campo, valor)
return cls.set(chave_str, dados)
payload = cls._json_dumps(dados)
pipe.multi()
pipe.set(chave_str, payload)
pipe.execute()
return True
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
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):
try:
@ -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,7 +1329,8 @@ class ContextoGlobalRedis:
operacao.get("inicializacao_operacao_executada", False)
)
if not inicializacao_executada:
cls._iniciar_operacao()
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
@ -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,10 +1862,125 @@ class ContextoGlobalRedis:
}
# ==========================================================
# 5. Enviar EXATAMENTE este mapa ao Manager
# 5/6. Entrega confiável Health -> Manager
# ==========================================================
# 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
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
def _estado_manager_mapsync():
estado = cls._dict(
cls.get(CtxKey.DadosManagerWorker, {})
)
return {
"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 ""
),
}
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] enviando mapa ao Manager | "
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}"
)
@ -1784,35 +2001,28 @@ class ContextoGlobalRedis:
"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,
expected_point_count=qtd_esperada,
timeout_s=2.5,
)
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)