agrobot_base/AgroBase/AgroBase/Services/CANServiceSerial.cs

759 lines
26 KiB
C#

using AgroBase.Models;
using Emgu.CV.Structure;
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Diagnostics;
using System.IO.Ports;
using System.Linq;
using System.Text;
using System.Threading;
using System.Threading.Tasks;
using static AgroBase.Models.Enums;
namespace AgroBase.Services
{
public class CanServiceSerial : ICanService
{
public SerialPort _PortaCAN;
public string _portName { get; set; }
private int _baudRate;
private int _bitRate;
private string _canVersion;
private object _lock = new object();
public bool IsConnected { get; private set; }
private volatile bool _precisaReabrir;
public bool Iniciado { get; private set; }
private AsyncTaskTimerModel tmrEnviarCAN;
private AsyncTaskTimerModel tmrOuvirCAN;
private AsyncTaskTimerModel tmrProcessarCAN;
private ConcurrentQueue<string> FramesRecebidos = new ConcurrentQueue<string>();
private StringBuilder SerialCANBuffer = new StringBuilder();
private readonly ConcurrentDictionary<byte, ICanMessageHandler> _roteadores = new ConcurrentDictionary<byte, ICanMessageHandler>();
private readonly ConcurrentDictionary<byte, bool> _roteadoresRegistrados = new ConcurrentDictionary<byte, bool>();
private byte? _ultimoIdEnviado = null;
private List<CanMessage> Mensagens = new List<CanMessage>();
private readonly object _LockMensagens = new object();
public double UltimoTx { get; set; } = 0;
private DateTime _ultimoTx = DateTime.MinValue;
public double FreqTx { get; set; } = 0;
public double UltimoRx { get; set; } = 0;
private DateTime _ultimoRx = DateTime.MinValue;
public double FreqRx{ get; set; } = 0;
public int ErrosCriticos { get; set; } = 0;
public double UltimoErroCritico { get; set; } = 0;
private DateTime _ultimoErroCritico = DateTime.MinValue;
private bool DebugMode = false;
private void MostrarLog(string message, byte id)
{
if (DebugMode && new List<byte>() { 0x11 }.Contains(id))
{
Console.WriteLine($"[CAN] {message}");
}
}
public bool PortaIsCan(SerialPort Porta)
{
if (Porta == null && _PortaCAN != null)
{
bool portaAberta = _PortaCAN?.IsOpen ?? false;
if (!portaAberta)
{
try
{
_PortaCAN?.Open();
}
catch
{
Iniciado = false;
}
}
else
{
var portas = SerialPort.GetPortNames();
var portaExiste = !string.IsNullOrEmpty(_portName) && portas.Contains(_portName, StringComparer.OrdinalIgnoreCase);
if (!portaExiste)
{
Iniciado = false;
}
}
}
else if (Porta != null)
{
Iniciado = false;
try
{
Porta.Close();
_portName = Porta.PortName;
_baudRate = 115200;
lock (_lock)
{
_PortaCAN = new SerialPort(_portName, _baudRate, Parity.None, 8, StopBits.One)
{
NewLine = "\r",
ReadTimeout = 500,
WriteTimeout = 500,
DtrEnable = true, // alguns adaptadores pedem
RtsEnable = true
};
_PortaCAN.ErrorReceived += (s, e) =>
{
IsConnected = false;
_precisaReabrir = true;
};
_PortaCAN.PinChanged += (s, e) =>
{
// Se perder DSR/CD ou receber BREAK, considere desconexão
if (e.EventType == SerialPinChange.Break || e.EventType == SerialPinChange.CDChanged || e.EventType == SerialPinChange.DsrChanged)
{
IsConnected = false;
_precisaReabrir = true;
}
};
_PortaCAN.Open();
for (int i = 0; i < 3; i++)
{
EnviarComando("V\r");
Thread.Sleep(50);
byte[] res = SerialService.LerDadosDaPortaSerial(_PortaCAN, 20);
_canVersion = Encoding.ASCII.GetString(res);
if (_canVersion.Contains("canable") || _canVersion.Contains("normaldotcom") || _canVersion.Contains("github"))
{
Iniciado = true;
SerialService.DispositivosMapeados.Add(new DispositivoDetalhesModel()
{
Dispositivo = T_Code.Can,
Endereco = _portName
});
break;
}
}
}
}
catch
{
}
}
else
{
Iniciado = false;
}
if (!Iniciado && !string.IsNullOrEmpty(_portName))
{
Fechar();
}
return Iniciado;
}
public void RegistrarHandler(byte endereco_rx, ICanMessageHandler handler)
{
_roteadoresRegistrados.TryGetValue(endereco_rx, out bool registrado);
if (!registrado)
{
_roteadoresRegistrados.TryAdd(endereco_rx, true);
_roteadores.TryAdd(endereco_rx, handler);
MostrarLog($"Roteador registrado: {endereco_rx.ToString("X")}", endereco_rx);
}
}
public bool Inicializar(int BitRate)
{
try
{
if (_PortaCAN == null)
{
return false;
}
_bitRate = BitRate;
if (!IsConnected)
{
lock (_lock)
{
if (!Iniciado)
{
if (!PortaIsCan(_PortaCAN))
{
IsConnected = false;
return false;
}
}
EnviarComando("C\r"); // fecha CAN (estado conhecido)
Thread.Sleep(50);
EnviarComando("Z0\r"); // desabilita timestamp
Thread.Sleep(50);
EnviarComando(ComandoBitRate(_bitRate)); // 500 kbps
Thread.Sleep(50);
EnviarComando("O\r"); // open
Thread.Sleep(50);
_PortaCAN.DiscardInBuffer();
_PortaCAN.DiscardOutBuffer();
}
_ultimoRx = DateTime.Now;
UltimoRx = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
_ultimoTx = DateTime.Now;
UltimoTx = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
tmrOuvirCAN?.Dispose();
tmrOuvirCAN = new AsyncTaskTimerModel("tmrOuvirCAN", OuvirCAN, 1);
tmrOuvirCAN.Start();
tmrEnviarCAN?.Dispose();
tmrEnviarCAN = new AsyncTaskTimerModel("tmrEnviarCAN", EnviarCAN, 3);
tmrEnviarCAN.Start();
tmrProcessarCAN?.Dispose();
tmrProcessarCAN = new AsyncTaskTimerModel("tmrProcessarCAN", ProcessarCAN, 5);
tmrProcessarCAN.Start();
IsConnected = true;
}
return IsConnected;
}
catch (Exception ex)
{
lock (_lock)
{
_PortaCAN = null;
IsConnected = false;
}
return false;
}
}
private async Task ProcessarCAN()
{
const int MAX_BATCH = 256; // teto de itens por tick
const int TIME_BUDGET_US = 1000; // ~1 ms de processamento
var sw = System.Diagnostics.Stopwatch.StartNew();
int processed = 0;
while (processed < MAX_BATCH && FramesRecebidos.TryDequeue(out string frame))
{
ProcessarFrameCAN(frame);
processed++;
// Para não monopolizar o loop se entrar enxurrada de frames
if (sw.ElapsedTicks * (1_000_000.0 / System.Diagnostics.Stopwatch.Frequency) > TIME_BUDGET_US)
break;
}
LimparMensagensAntigas();
VerificaSaudeConexao();
}
private async Task OuvirCAN()
{
if ((_PortaCAN?.IsOpen ?? false) && (_PortaCAN?.BytesToRead ?? 0) > 0)
{
byte[] tempBuffer = await SerialService.LerDadosDaPortaSerialAsync(_PortaCAN, _PortaCAN.BytesToRead);
foreach (byte b in tempBuffer)
{
char c = (char)b;
if (c == '\r')
{
var frame = SerialCANBuffer.ToString();
SerialCANBuffer.Clear();
if (!string.IsNullOrWhiteSpace(frame))
{
FramesRecebidos.Enqueue(frame);
}
}
else
{
SerialCANBuffer.Append(c);
}
}
if (tempBuffer.Length > 0)
{
var agora = DateTime.Now;
FreqRx = 1.0 / (agora - _ultimoRx).TotalSeconds;
_ultimoRx = agora;
UltimoRx = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
}
}
}
private void ProcessarFrameCAN(string frame)
{
// Ex: t00183301020304
if (!frame.StartsWith("t")) return;
try
{
if (frame.Length < 5) return;
int id = Convert.ToInt32(frame.Substring(1, 3), 16);
int len = int.Parse(frame.Substring(4, 1));
byte[] data = new byte[len];
if (data == null || data.Length == 0) return;
for (int i = 0; i < len; i++)
{
data[i] = Convert.ToByte(frame.Substring(5 + i * 2, 2), 16);
}
MostrarLog($"Processando frame: {frame}", (byte)id);
_roteadores.TryGetValue((byte)id, out ICanMessageHandler handler);
T_Code dispositivo = handler?.Dispositivo ?? T_Code.Vzo;
byte fCodeRx = data.Length > 0 ? data[0] : (byte)0x00;
byte parametroRx = data.Length > 1 ? data[1] : (byte)0x00;
string chave = GerarChaveMensagem(dispositivo, (byte)id, data);
CanMessage mensagem = null;
lock (_LockMensagens)
{
mensagem = Mensagens.FirstOrDefault(x => x.Enviado && !x.Respondido && x.Dispositivo == dispositivo && x.IdRx == id && x.DataTx[1] == data[1]);
}
if (mensagem == null)
{
mensagem = new CanMessage { Dispositivo = dispositivo, IdTx = (byte)id, IdRx = (byte)id, funcCodeRx = fCodeRx, DataRx = data, ChaveInterna = chave };
}
mensagem.DataRx = data;
mensagem.Respondido = true;
mensagem.RespondidoEm = DateTime.Now;
if ((handler?.Dispositivo ?? T_Code.Vzo) != T_Code.Vzo)
{
handler.ProcessarMensagem(mensagem);
}
else
{
MostrarLog($"Roteador não registrado! ID_Num não tratado: 0x{id:X2}, fCodeRx: 0x{fCodeRx:X2}, frame: {frame}", (byte)id);
}
/*byte[] data = new byte[len];
for (int i = 0; i < len; i++)
{
data[i] = Convert.ToByte(frame.Substring(5 + i * 2, 2), 16);
}
byte funcCodeRx = data[0];
var mensagem = new CanMessage();
lock (Mensagens)
{
mensagem = Mensagens.Where(x => x.IdTx == id).OrderBy(x => x.EnviadoEm).FirstOrDefault(x => x.Enviado && !x.Respondido && x.funcCodeRx == funcCodeRx);
}
if (mensagem == null)
{
mensagem = new CanMessage { IdTx = (byte)id, DataRx = data };
}
mensagem.DataRx = data;
mensagem.Respondido = true;
mensagem.RespondidoEm = DateTime.Now;
if (_roteadores.TryGetValue(mensagem.IdTx, out var handler))
{
handler.ProcessarMensagem(mensagem);
}
else
{
Console.WriteLine($"[CAN] ID_Num não tratado: 0x{funcCodeRx:X2}");
}*/
}
catch (Exception ex)
{
Console.WriteLine($"[CAN] Erro ao processar frame: {ex.Message}");
}
}
public void AdicionarMensagemNaFila(T_Code Dispositivo, byte idTx, byte idRx, byte funcCodeTx, byte funcCodeRx, byte[] payload = null, bool get = true)
{
if (_PortaCAN?.IsOpen != true) return;
var dados = new List<byte>() { funcCodeTx };
if (payload != null)
dados.AddRange(payload);
if (Dispositivo == T_Code.Mks)
{
byte crc = CalcularChecksum(idTx, dados.ToArray());
dados.Add(crc);
}
var novaMensagem = new CanMessage
{
Momento = DateTime.Now,
Dispositivo = Dispositivo,
IdTx = idTx,
IdRx = idRx,
funcCodeTx = funcCodeTx,
funcCodeRx = funcCodeRx,
DataTx = dados.ToArray(),
ComResposta = get,
ChaveInterna = GerarChaveMensagem(Dispositivo, idRx, dados.ToArray())
};
bool mensagemNaFila = false;
lock (_LockMensagens)
{
mensagemNaFila = Mensagens.Any(x => (!x.Enviado || (x.ComResposta && x.Enviado && !x.Respondido)) && x.IdTx == idTx && x.IdRx == idRx && x.DataTx.SequenceEqual(dados.ToArray()));
if (!mensagemNaFila)
{
Mensagens.Add(novaMensagem);
MostrarLog($"Mensagem adicionada na fila. ID: {novaMensagem.IdTx}, Data: {FuncoesGlobais.ConverterComandoBytesParaTexto(novaMensagem.DataTx)}", idTx);
//if (dados[1] == Variaveis.ID_Num_sLRA)
//{
// Console.WriteLine($"{DateTime.Now.ToString("dd/MM/yyyy HH:mm:ss.fff")} - MENSAGEM LORA ADICIONADO NA FILA CAN, FUNC_CODE: {((CanMessagePosicaoDados)dados[0]).ToString()} ({((DadosLoRaParse)dados[0]).ToString()}) " + string.Join(" ", dados));
//}
}
else
{
MostrarLog($"Mensagem já pendente e idêntica. ID: {novaMensagem.IdTx}, Data: {FuncoesGlobais.ConverterComandoBytesParaTexto(novaMensagem.DataTx)}", idTx);
}
}
}
private CanMessage ObterProximaMensagem()
{
List<CanMessage> mensagensPendentes;
lock (_LockMensagens)
{
mensagensPendentes = Mensagens
.Where(x => !x.Enviado)
.ToList();
}
if (!mensagensPendentes.Any())
return null;
// Primeiro tenta pegar mensagens com Get == false
var mensagensPrioritarias = mensagensPendentes
.Where(x => !x.ComResposta)
.ToList();
// Se não houver mensagens prioritárias, considera todas
var mensagensParaProcessar = mensagensPrioritarias.Any()
? mensagensPrioritarias
: mensagensPendentes;
// Agrupa por IdTx e ordena os grupos
var idsDisponiveis = mensagensParaProcessar
.Select(x => x.IdTx)
.Distinct()
.OrderBy(x => x)
.ToList();
byte proximoId;
if (!_ultimoIdEnviado.HasValue || !idsDisponiveis.Contains(_ultimoIdEnviado.Value))
{
proximoId = idsDisponiveis.First();
}
else
{
var index = idsDisponiveis.IndexOf(_ultimoIdEnviado.Value);
proximoId = (index + 1 < idsDisponiveis.Count) ? idsDisponiveis[index + 1] : idsDisponiveis.First();
}
// Seleciona a mensagem mais antiga com esse IdTx
var proximaMsg = mensagensParaProcessar
.Where(x => x.IdTx == proximoId)
.OrderBy(x => x.Momento)
.FirstOrDefault();
return proximaMsg;
}
private async Task EnviarCAN()
{
if (_PortaCAN?.IsOpen != true) return;
RefillTokens();
// pequeno burst controlado para reduzir latência:
const int MAX_BURST = 8;
int sent = 0;
while (_tokens >= 1 && sent < MAX_BURST)
{
var msg = ObterProximaMensagem();
if (msg == null) break;
//if (!PodeEnviarAgora(msg)) break; // sem token “perdido”: sai do burst e tenta no próximo tick
// Evite delays arbitrários em caminho “quente”
// Se precisar espaçar MKS, use cooldown por-ID (abaixo)
if (!EnviarComando(CriarFrame(msg)))
{
MostrarLog("Erro ao enviar frame", msg.IdTx);
break; // evita apertar ainda mais se o driver falhar
}
msg.Enviado = true;
msg.EnviadoEm = DateTime.Now;
_ultimoIdEnviado = msg.IdTx;
_tokens -= 1;
sent++;
}
if (sent == 0)
await Task.Delay(1); // cede o scheduler quando não tinha token/mensagem
}
private string CriarFrame(CanMessage msg)
{
string frame = $"t{msg.IdTx:X3}{msg.DataTx.Length}";
foreach (byte b in msg.DataTx)
frame += b.ToString("X2");
frame += "\r";
return frame;
}
private bool EnviarComando(string comando)
{
lock (_lock)
{
if (_PortaCAN?.IsOpen == true)
{
_PortaCAN.Write(comando);
//Console.WriteLine($"[CAN] Frame enviado: {comando}");
var agora = DateTime.Now;
FreqTx = 1.0 / (agora - _ultimoTx).TotalSeconds;
_ultimoTx = agora;
UltimoTx = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
return true;
}
return false;
}
}
private void LimparMensagensAntigas()
{
var agora = DateTime.Now;
lock (_LockMensagens)
{
List<CanMessage> mensagens = Mensagens
.Where(x =>
(!x.ComResposta && x.Enviado) ||
(x.ComResposta && x.Respondido) ||
(x.ComResposta && x.Enviado && !x.Respondido && x.EnviadoEm < agora.AddSeconds(-1)) ||
(x.Momento < agora.AddSeconds(-30))
)
.ToList();
foreach (var item in mensagens)
{
Mensagens.Remove(item);
}
}
}
private void VerificaSaudeConexao()
{
DateTime Agora = DateTime.Now;
// Porta sumiu do SO?
var portas = SerialPort.GetPortNames();
var portaExiste = !string.IsNullOrEmpty(_portName) && portas.Contains(_portName, StringComparer.OrdinalIgnoreCase);
if (!portaExiste || _PortaCAN == null || !_PortaCAN.IsOpen)
{
IsConnected = false;
// peça uma varredura geral (vai tentar outros nomes de COM):
Task.Run(() => SerialService.RealizarVarreduraPortasUSB());
return;
}
if (IsConnected && (((Agora - _ultimoTx).TotalSeconds > 2 && (Agora - _ultimoRx).TotalSeconds > 2) || _precisaReabrir))
{
ErrosCriticos++;
_ultimoErroCritico = Agora;
UltimoErroCritico = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
_precisaReabrir = false;
MostrarLog("Muito tempo sem receber novos dados, reiniciando conexao CAN...", 0);
ReabrirPortaSegura();
}
if ((Agora - _ultimoErroCritico).TotalMinutes > 10)
{
ErrosCriticos = 0;
}
}
public void Fechar()
{
tmrOuvirCAN?.Dispose();
tmrEnviarCAN?.Dispose();
tmrProcessarCAN?.Dispose();
lock (_LockMensagens)
{
Mensagens.Clear();
}
SerialCANBuffer.Clear();
FramesRecebidos = new ConcurrentQueue<string>();
//_roteadores.Clear();
_ultimoIdEnviado = null;
lock (_lock)
{
if ((_PortaCAN?.IsOpen ?? false) == true)
{
EnviarComando("C\r");
_PortaCAN.Close();
}
_PortaCAN = null;
}
var d = SerialService.DispositivosMapeados.FirstOrDefault(x => x.Endereco == _portName);
SerialService.DispositivosMapeados.Remove(d);
_portName = null;
_canVersion = null;
IsConnected = false;
}
public void Dispose()
{
Fechar();
}
private byte CalcularChecksum(byte id, byte[] data)
{
int soma = id;
foreach (var b in data)
soma += b;
return (byte)(soma & 0xFF);
}
private string GerarChaveMensagem(T_Code Dispositivo, byte Id, byte[] data)
{
byte parametro = 0x00;
// Para dispositivos que não sejam MKS, usar DataTx[1] se existir
if (Dispositivo != T_Code.Mks && (data?.Length ?? 0) > 1)
{
parametro = data[1];
}
return $"{Id}_{data[0]}_{parametro}";
}
// alvo: ~200 fps (ajuste conforme bitrate/carga)
private const double TARGET_FPS = 200;
private double _tokens = TARGET_FPS;
private const double _capacity = TARGET_FPS;
private const double _refillPerMs = TARGET_FPS / 1000.0;
private DateTime _lastRefill = DateTime.UtcNow;
private void RefillTokens()
{
var now = DateTime.UtcNow;
var dtMs = (now - _lastRefill).TotalMilliseconds;
if (dtMs > 0)
{
_tokens = Math.Min(_capacity, _tokens + dtMs * _refillPerMs);
_lastRefill = now;
}
}
private string ComandoBitRate(int rate)
{
switch (rate)
{
case 10:
return "S0\r";
case 20:
return "S1\r";
case 50:
return "S2\r";
case 100:
return "S3\r";
case 125:
return "S4\r";
case 250:
return "S5\r";
case 500:
return "S6\r";
case 800:
return "S7\r";
case 1000:
return "S8\r";
default:
return "S6\r";
}
}
private void ReabrirPortaSegura()
{
MostrarLog("Reabrindo porta segura...", 0);
tmrEnviarCAN?.Stop();
tmrOuvirCAN?.Stop();
tmrProcessarCAN?.Stop();
lock (_LockMensagens)
{
Mensagens.Clear();
}
lock (_lock)
{
try { _PortaCAN?.Close(); } catch { }
try
{
_PortaCAN?.Open();
_PortaCAN.DiscardInBuffer();
_PortaCAN.DiscardOutBuffer();
EnviarComando("C\r"); Thread.Sleep(50);
EnviarComando("Z0\r"); Thread.Sleep(50);
EnviarComando(ComandoBitRate(_bitRate)); Thread.Sleep(50);
EnviarComando("O\r"); Thread.Sleep(50);
_ultimoRx = DateTime.Now;
UltimoRx = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
_ultimoTx = DateTime.Now;
UltimoTx = Stopwatch.GetTimestamp() / (double)Stopwatch.Frequency;
}
catch
{
IsConnected = false;
return;
}
}
tmrEnviarCAN?.Start();
tmrOuvirCAN?.Start();
tmrProcessarCAN?.Start();
}
}
}