Inizio modifica x passaggio richieste tramite REIDS channels pub/sub

This commit is contained in:
Samuele Locatelli
2025-07-21 16:49:54 +02:00
parent 4e354a1767
commit 367379542c
3 changed files with 120 additions and 5 deletions
+119 -4
View File
@@ -3,7 +3,11 @@ using Lux.Data;
using Lux.Data.Services;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.Mvc;
using Microsoft.Extensions.Configuration;
using Newtonsoft.Json;
using StackExchange.Redis;
using System.Diagnostics;
using System.Runtime;
namespace Lux.API.Controllers
{
@@ -17,16 +21,119 @@ namespace Lux.API.Controllers
{
_logger = logger;
_config = config;
// setup compoenti REDIS
redisConn = ConnectionMultiplexer.Connect(_config.GetConnectionString("Redis") ?? "localhost");
redisDb = redisConn.GetDatabase();
// json serializer... FIX errore loop circolare https://www.ryadel.com/en/jsonserializationexception-self-referencing-loop-detected-error-fix-entity-framework-asp-net-core/
JSSettings = new JsonSerializerSettings()
{
ReferenceLoopHandling = ReferenceLoopHandling.Ignore
};
// init classe sottoscrizione PubSub CHannel messages REDIS
messageDisp.Subscribe(ChannelName("", true), (channel, message) =>
{
SaveCalcData(channel, message);
});
// verifico se usare engine
enableEgwEng = _config.GetValue<bool>("ServerConf:EgwEngineEnab");
#if false
if (enableEgwEng && EgwProcManager != null)
{
ProcessMan = EgwProcManager;
ProcessMan.m_AnswerReceived += ProcessMan_m_AnswerReceived;
}
}
#endif
ImgService = imgServ;
}
/// <summary>
/// Salva risultato calcolo da broadcast channel REDIS
/// </summary>
/// <param name="channel"></param>
/// <param name="message"></param>
private void SaveCalcData(RedisChannel channel, RedisValue message)
{
string rawData = $"{message}";
if (!string.IsNullOrEmpty(rawData) && rawData.Length > 2)
{
// provo a deserializzare
try
{
var retData = JsonConvert.DeserializeObject<ProcessArgsResult>(rawData);
if (retData != null)
{
// verifico nId di risposta x salvare correttamente
ProcessMan_m_AnswerReceived(retData);
}
}
catch (Exception exc)
{
_logger.LogError($"Errore in fase decodifica messaggio da REDIS Channel{Environment.NewLine}{exc}");
}
}
}
private void sendMessage(string notifyChannel, string message)
{
RedisChannel pubNotifyChannel = new RedisChannel(notifyChannel, RedisChannel.PatternMode.Literal);
messageDisp.Publish(pubNotifyChannel, message);
}
/// <summary>
/// Nome del channel sottoscritto per ritorno calcoli
/// </summary>
/// <param name="servId"></param>
/// <returns></returns>
private string ChannelName(string servId, bool isOut)
{
return isOut ? $"EgwEngineOutput" : $"EgwEngineInput";
// se implementassimo multi-elaborazione calcoli ogni esecutore ha un suo channel
//return $"EgwEngineOutput_{servId}";
}
/// <summary>
/// Oggetto subscriber x pubblicazione/sottoscrizione canali REDIS
/// </summary>
private ISubscriber _currSub { get; set; } = null!;
/// <summary>
/// Message Dispatcher: oggetto comunicazione pub/sub via REDIS channels corrente
/// </summary>
protected ISubscriber messageDisp
{
get
{
ISubscriber answ;
// se già valorizzato uso oggetto private...
if (_currSub != null)
{
answ = _currSub;
}
else
{
// sottoscrizione al dispatcher messaggi
answ = redisConn.GetSubscriber();
_currSub = answ;
}
// restituisco oggetto DB
return answ;
}
}
protected JsonSerializerSettings? JSSettings;
/// <summary>
/// Oggetto per connessione a REDIS
/// </summary>
protected ConnectionMultiplexer redisConn = null!;
/// <summary>
/// Oggetto DB redis da impiegare x chiamate R/W
/// </summary>
protected IDatabase redisDb = null!;
#endregion Public Constructors
#region Public Methods
@@ -47,13 +154,18 @@ namespace Lux.API.Controllers
// ...se ricevo percorso --> leggo jwd/svg cablato
if (!string.IsNullOrEmpty(currJwd))
{
if (enableEgwEng && ProcessMan != null)
if (true || (enableEgwEng && ProcessMan != null))
{
Dictionary<string, string> DictExec = new Dictionary<string, string>();
DictExec.Add("Mode", "1");
DictExec.Add("Jwd", currJwd);
int nId = 1;
ProcessArgs currArgs = new ProcessArgs(nId, DictExec);
sendMessage(ChannelName("", false), currArgs.sProcessArgs);
svgContent = "DONE";
#if false
bool done = ProcessMan.ArgumentsEnqueue(currArgs);
waitResult = true;
@@ -62,8 +174,10 @@ namespace Lux.API.Controllers
{
numWait--;
await Task.Delay(waitDelay);
}
}
#endif
}
#if false
// se ho risultato mostro...
if (!string.IsNullOrEmpty(lastSvg))
{
@@ -74,7 +188,8 @@ namespace Lux.API.Controllers
else
{
svgContent = "EMPTY";
}
}
#endif
}
sw.Stop();
_logger.LogInformation($"svgString | {sw.Elapsed.TotalMilliseconds:N3} ms");