agrobot_base/AgroBase/AgroBase/Services/FilaService.cs

262 lines
8.5 KiB
C#

using AgroBase.Models;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using static AgroBase.Models.Enums;
namespace AgroBase.Services
{
public class FilaService
{
private static List<MensagemFilaModel> MensagensFila;
private static readonly object _lock = new object();
private static int _counter = 1;
private static AsyncTaskTimerModel tmrFila;
private static bool DebugMessages = false;
public static bool Liberado
{
get
{
lock (_lock)
{
return MensagensFila == null || !MensagensFila.Any(x => !x.Enviado && x.AdicionadoEm.AddSeconds(10) > DateTime.Now);
}
}
}
public static void IniciarRotinas()
{
lock (_lock)
{
MensagensFila = new List<MensagemFilaModel>();
}
PararRotinas();
tmrFila = new AsyncTaskTimerModel("tmrFila", tmrFila_Tick, 100);
tmrFila.Start();
}
public static void PararRotinas()
{
tmrFila?.Dispose();
}
public static void AdicionarMensagemFila(string Modulo_ID, T_Code Dispositivo, string Mensagem, bool _prioridade, bool _aguardaCallback)
{
var _Mensagem = new MensagemFilaModel()
{
Id = _counter,
ID_Modulo = Modulo_ID,
Dispositivo = Dispositivo,
Mensagem = Mensagem,
AdicionadoEm = DateTime.Now,
Enviado = false,
Recebido = false,
Prioridade = _prioridade,
AguardaCallback = _aguardaCallback,
Reenvios = 0,
};
RealizarEnvio(_Mensagem);
return;
lock (_lock)
{
if (!MensagensFila.Any(x => !x.Enviado && x.Dispositivo == Dispositivo && x.ID_Modulo == Modulo_ID && x.Mensagem == Mensagem))
{
MensagensFila.Add(_Mensagem);
_counter++;
}
}
}
public static void MarcarMensagemEnviada(int Id)
{
lock (_lock)
{
var Mensagem = MensagensFila.FirstOrDefault(x => x.Id == Id);
if (Mensagem != null && !Mensagem.Enviado)
{
try
{
Mensagem.Enviado = true;
Mensagem.EnviadoEm = DateTime.Now;
if (DebugMessages)
{
Console.WriteLine($"Mensagem {Id} marcada como enviada: {Mensagem.Mensagem}");
}
}
catch (Exception ex)
{
Console.WriteLine($"Erro ao marcar mensagem como enviada: {ex.Message}");
}
}
}
}
public static void MarcarMensagemRecebida(int Id)
{
lock (_lock)
{
var Mensagem = MensagensFila.FirstOrDefault(x => x.Id == Id);
if (Mensagem != null && !Mensagem.Recebido)
{
try
{
Mensagem.Enviado = true;
Mensagem.Recebido = true;
Mensagem.RecebidoEm = DateTime.Now;
if (DebugMessages)
{
Console.WriteLine($"Mensagem {Id} marcada como recebida: {Mensagem.Mensagem}");
}
}
catch (Exception ex)
{
Console.WriteLine($"Erro ao marcar mensagem como recebida: {ex.Message}");
}
}
}
}
private static bool RealizarEnvio(MensagemFilaModel Mensagem)
{
try
{
var Dispositivo = Variaveis.DispositivosConectados.FirstOrDefault(x => x.Dispositivo == Mensagem.Dispositivo);
if (Dispositivo == null) return false;
/*Dispositivo.EnviarDadosMQTT(
//SerialService.BeginLine +
Mensagem.Id +
SerialService.SplitMessage +
Mensagem.ID_Modulo +
SerialService.SplitMessage +
Mensagem.Mensagem +
SerialService.EndLine
);*/
MarcarMensagemRecebida(Mensagem.Id);
return true;
}
catch (Exception ex)
{
Console.WriteLine($"Erro ao enviar mensagem para o dispositivo {Mensagem.Dispositivo}: {ex.Message}");
return false;
}
}
public static bool ExistemMensagensPendentesEnvio(T_Code Dispositivo)
{
lock (_lock)
{
return MensagensFila.Any(x => x.Dispositivo == Dispositivo && !x.Enviado);
}
}
public static bool ExistemMensagensPendentesResposta(T_Code Dispositivo)
{
lock (_lock)
{
return MensagensFila.Any(x => x.Dispositivo == Dispositivo && x.Enviado && !x.Recebido);
}
}
private static bool FilaLiberada(T_Code Dispositivo)
{
lock (_lock)
{
var Fila = MensagensFila.Where(x => x.Dispositivo == Dispositivo).OrderBy(x => x.AdicionadoEm);
var MensagemPendente = Fila.FirstOrDefault(x => x.Enviado && !x.Recebido);
bool Liberado = MensagemPendente == null || (!MensagemPendente.AguardaCallback || MensagemPendente.EnviadoEm.Value.AddMilliseconds(1500) <= DateTime.Now);
return Liberado;
}
}
public static async Task tmrFila_Tick()
{
await ExecutarTickAsync();
}
private static async Task ExecutarTickAsync()
{
RemoverDuplicatas();
lock (_lock)
{
var Filas = MensagensFila.Where(x => !x.Enviado).GroupBy(x => x.Dispositivo);
foreach (var Fila in Filas)
{
if (FilaLiberada(Fila.Key))
{
MensagemFilaModel Mensagem = null;
lock (_lock)
{
Mensagem = MensagensFila.Where(x => x.Dispositivo == Fila.Key)
.OrderByDescending(x => x.Prioridade)
.ThenBy(x => x.AdicionadoEm)
.FirstOrDefault(x => !x.Enviado);
}
if (Mensagem != null)
{
if (RealizarEnvio(Mensagem))
{
MarcarMensagemEnviada(Mensagem.Id);
}
}
}
}
}
VerificarMensagensSemResposta();
}
private static void RemoverDuplicatas()
{
lock (_lock)
{
var mensagensAgrupadas = MensagensFila
.Where(x => !x.Enviado)
.GroupBy(x => new { x.Dispositivo, x.ID_Modulo, x.Mensagem })
.Select(g => g.OrderBy(x => x.AdicionadoEm).First())
.ToList();
// Remover todas as mensagens não enviadas
MensagensFila.RemoveAll(x => !x.Enviado);
// Adicionar de volta apenas as mensagens mais antigas de cada grupo
MensagensFila.AddRange(mensagensAgrupadas);
MensagensFila.RemoveAll(x => x.Enviado && x.Recebido);
}
}
private static void VerificarMensagensSemResposta()
{
lock (_lock)
{
var mensagensSemResposta = MensagensFila
.Where(x => x.Enviado && !x.Recebido && x.EnviadoEm.Value.AddMilliseconds(1500) <= DateTime.Now && x.Reenvios < 5)
.ToList();
foreach (var Mensagem in mensagensSemResposta)
{
Mensagem.Enviado = false;
Mensagem.AdicionadoEm = DateTime.Now;
Mensagem.Reenvios++;
}
}
}
}
}