487 lines
17 KiB
C#
487 lines
17 KiB
C#
using AgroBase.Models;
|
|
using System;
|
|
using System.Collections.Concurrent;
|
|
using System.Collections.Generic;
|
|
using System.IO.Ports;
|
|
using System.Linq;
|
|
using System.Text;
|
|
using System.Threading.Tasks;
|
|
using static AgroBase.Models.Enums;
|
|
using static WaveshareCANFDService;
|
|
|
|
namespace AgroBase.Services
|
|
{
|
|
internal class CanServiceWaveshare : ICanService
|
|
{
|
|
public WaveshareCANFDService Connector = new WaveshareCANFDService(0);
|
|
public string _portName { get; set; } = "CAN-FD";
|
|
private int _bitRate;
|
|
public bool IsConnected { get; private set; }
|
|
public bool Iniciado
|
|
{
|
|
get
|
|
{
|
|
return SerialService.DispositivosMapeados.Any(x => /*x.Status != StatusModulo.Desconectado &&*/ SerialService.DispositivosCan.Contains(x.Dispositivo));
|
|
}
|
|
}
|
|
private AsyncTaskTimerModel tmrEnviarCAN;
|
|
private AsyncTaskTimerModel tmrProcessarCAN;
|
|
private StringBuilder SerialCANBuffer = new StringBuilder();
|
|
private readonly ConcurrentQueue<(uint id, byte[] data)> FramesRecebidos = new ConcurrentQueue<(uint, byte[])>();
|
|
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> MensagensPendentes = new List<CanMessage>();
|
|
private readonly object _LockMensagens = new object();
|
|
private bool DebugMode = false;
|
|
private DateTime UltimoTx = DateTime.MinValue;
|
|
private DateTime UltimoRx = DateTime.MinValue;
|
|
|
|
private void MostrarLog(string message, byte id)
|
|
{
|
|
if (DebugMode/* && new List<byte>() { 0x01 }.Contains(id)*/)
|
|
{
|
|
Console.WriteLine($"[CAN] {message}");
|
|
}
|
|
}
|
|
|
|
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(SerialPort Porta = null, int BitRate = 250)
|
|
{
|
|
_bitRate = BitRate;
|
|
bool canIniciado = Connector.IsOpen;
|
|
|
|
if (!canIniciado)
|
|
{
|
|
canIniciado = Connector.Open((VCI_CAN_OBJ[] frames) =>
|
|
{
|
|
try
|
|
{
|
|
var now = DateTime.Now;
|
|
for (int i = 0; i < frames.Length; i++)
|
|
{
|
|
var f = frames[i];
|
|
var dlc = f.DataLen;
|
|
if (dlc < 0) dlc = 0; if (dlc > 8) dlc = 8;
|
|
|
|
// clone rápido (já vem deep do conector, mas garantimos)
|
|
var data = dlc > 0 ? (byte[])f.Data.Clone() : Array.Empty<byte>();
|
|
|
|
// nunca bloquear. se quiser cap, faça:
|
|
FramesRecebidos.Enqueue((f.ID, data));
|
|
while (FramesRecebidos.Count > 4096) // cap opcional
|
|
FramesRecebidos.TryDequeue(out _);
|
|
}
|
|
UltimoRx = now;
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
MostrarLog($"[CAN] ERRO no callback consumidor: {ex.Message}", 0);
|
|
}
|
|
}, bitrateKbps: _bitRate, samplePoint: 0.8);
|
|
}
|
|
else
|
|
{
|
|
return true;
|
|
}
|
|
|
|
if (!canIniciado)
|
|
{
|
|
MostrarLog($"Erro ao iniciar o conversor CAN FD no canal {Connector.Canal}.", 0);
|
|
}
|
|
else
|
|
{
|
|
MostrarLog($"Conversor CAN FD iniciado com sucesso no canal {Connector.Canal}.", 0);
|
|
}
|
|
|
|
if (canIniciado || !IsConnected)
|
|
{
|
|
UltimoRx = DateTime.Now;
|
|
UltimoTx = DateTime.Now;
|
|
|
|
tmrEnviarCAN?.Dispose();
|
|
tmrEnviarCAN = new AsyncTaskTimerModel("tmrEnviarCAN", EnviarCAN, 3);
|
|
tmrEnviarCAN.Start();
|
|
|
|
tmrProcessarCAN?.Dispose();
|
|
tmrProcessarCAN = new AsyncTaskTimerModel("tmrProcessarCAN", ProcessarCAN, 5);
|
|
tmrProcessarCAN.Start();
|
|
|
|
IsConnected = true;
|
|
}
|
|
|
|
return canIniciado;
|
|
}
|
|
|
|
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 var 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 void ProcessarFrameCAN((uint, byte[]) frame)
|
|
{
|
|
uint id = frame.Item1;
|
|
byte[] data = frame.Item2;
|
|
|
|
if (data == null || data.Length == 0) return;
|
|
|
|
try
|
|
{
|
|
MostrarLog($"Processando frame: {BitConverter.ToString(data)}", (byte)id);
|
|
|
|
_roteadores.TryGetValue((byte)id, out ICanMessageHandler handler);
|
|
|
|
T_Code dispositivo = handler?.Dispositivo ?? T_Code.Vzo;
|
|
byte fCodeRx = data[0];
|
|
byte parametroRx = data.Length > 1 ? data[1] : (byte)0x00;
|
|
string chave = GerarChaveMensagem(dispositivo, (byte)id, data);
|
|
|
|
CanMessage mensagem = null;
|
|
lock (_LockMensagens)
|
|
{
|
|
mensagem = MensagensPendentes.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);
|
|
}
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
MostrarLog($"Erro ao processar frame! Frame: {frame}, Erro: {ex.Message}", (byte)id);
|
|
}
|
|
}
|
|
|
|
public void AdicionarMensagemNaFila(T_Code Dispositivo, byte idTx, byte idRx, byte funcCodeTx, byte funcCodeRx, byte[] payload = null, bool get = true)
|
|
{
|
|
if (!Connector.IsOpen) 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 = MensagensPendentes.Any(x => (!x.Enviado || (x.ComResposta && x.Enviado && !x.Respondido)) && x.IdTx == idTx && x.IdRx == idRx && x.DataTx.SequenceEqual(dados.ToArray()));
|
|
if (!mensagemNaFila)
|
|
{
|
|
MensagensPendentes.Add(novaMensagem);
|
|
MostrarLog($"Mensagem adicionada na fila. ID: {novaMensagem.IdTx}, Data: {FuncoesGlobais.ConverterComandoBytesParaTexto(novaMensagem.DataTx)}", idTx);
|
|
//if (payload[0] == 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)payload[1]).ToString()} ({((DadosLoRaParse)payload[1]).ToString()})" + string.Join(" ", payload));
|
|
//}
|
|
}
|
|
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 = MensagensPendentes
|
|
.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 readonly ConcurrentDictionary<byte, DateTime> _cooldownPorId = new ConcurrentDictionary<byte, DateTime>();
|
|
private static readonly TimeSpan COOLDOWN_MKS = TimeSpan.FromMilliseconds(6);
|
|
|
|
private bool PodeEnviarAgora(CanMessage m)
|
|
{
|
|
// Exemplo: só MKS precisa de respiro
|
|
if (m.Dispositivo != T_Code.Mks) return true;
|
|
|
|
var now = DateTime.UtcNow;
|
|
var due = _cooldownPorId.GetOrAdd(m.IdTx, DateTime.MinValue);
|
|
if (now < due) return false;
|
|
|
|
_cooldownPorId[m.IdTx] = now + COOLDOWN_MKS;
|
|
return true;
|
|
}
|
|
|
|
private async Task EnviarCAN()
|
|
{
|
|
if (!Connector.IsOpen) 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), 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, CanMessage msg = null)
|
|
{
|
|
if (Connector.IsOpen)
|
|
{
|
|
bool sucesso = Connector.Send(msg.IdTx, msg.DataTx);
|
|
|
|
if (sucesso)
|
|
{
|
|
MostrarLog($"Frame enviado: {comando}", msg.IdTx);
|
|
UltimoTx = DateTime.Now;
|
|
}
|
|
else
|
|
{
|
|
MostrarLog($"Erro ao enviar frame: {comando}", msg.IdTx);
|
|
}
|
|
return sucesso;
|
|
}
|
|
|
|
return false;
|
|
}
|
|
|
|
private void LimparMensagensAntigas()
|
|
{
|
|
var agora = DateTime.Now;
|
|
|
|
lock (_LockMensagens)
|
|
{
|
|
List<CanMessage> mensagens = MensagensPendentes
|
|
.Where(x =>
|
|
(!x.ComResposta && x.Enviado) ||
|
|
(x.ComResposta && x.Respondido) ||
|
|
(x.ComResposta && x.Enviado && !x.Respondido && x.EnviadoEm < agora.AddSeconds(-2)) ||
|
|
(x.Momento < agora.AddSeconds(-30))
|
|
)
|
|
.ToList();
|
|
foreach (var item in mensagens)
|
|
{
|
|
MensagensPendentes.Remove(item);
|
|
}
|
|
}
|
|
}
|
|
|
|
private void VerificaSaudeConexao()
|
|
{
|
|
DateTime Agora = DateTime.Now;
|
|
if (IsConnected && ((Agora - UltimoTx).TotalSeconds > 10 || (Agora - UltimoRx).TotalSeconds > 10))
|
|
{
|
|
UltimoRx = Agora;
|
|
UltimoTx = Agora;
|
|
MostrarLog("Muito tempo sem receber novos dados, reiniciando conexao CAN...", 0);
|
|
try { Connector.Close(); } catch { }
|
|
try
|
|
{
|
|
Inicializar();
|
|
}
|
|
catch { }
|
|
}
|
|
}
|
|
|
|
public void Fechar()
|
|
{
|
|
tmrEnviarCAN?.Dispose();
|
|
|
|
lock (_LockMensagens)
|
|
{
|
|
MensagensPendentes.Clear();
|
|
}
|
|
|
|
SerialCANBuffer.Clear();
|
|
while (FramesRecebidos.Count > 0)
|
|
{
|
|
FramesRecebidos.TryDequeue(out var frame);
|
|
}
|
|
_roteadores.Clear();
|
|
_ultimoIdEnviado = null;
|
|
|
|
Connector.Close();
|
|
|
|
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)
|
|
{
|
|
if (data == null || data.Length == 0) return $"{Id}_empty";
|
|
|
|
// MKS: chaveia por funcCode apenas
|
|
if (Dispositivo == T_Code.Mks)
|
|
return $"{Id}_{data[0]}";
|
|
|
|
// Outros: mantém param se houver
|
|
byte parametro = (data.Length > 1) ? data[1] : (byte)0x00;
|
|
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;
|
|
}
|
|
}
|
|
|
|
}
|
|
}
|